How UnifiedGraphWriter Handles Concurrent Writes and Data Conflicts in Twitter's Algorithm
UnifiedGraphWriter eliminates data conflicts by isolating each writer thread to a distinct graph segment and using a lock-free ConcurrentLinkedQueue with semaphore-based back-pressure to safely ingest high-volume Kafka streams.
The UnifiedGraphWriter trait in Twitter's open-source recommendation system (twitter/the-algorithm) powers the real-time construction of bipartite recommendation graphs from Kafka event streams. By strictly isolating write operations to specific graph segments and leveraging Java's concurrent utilities, it achieves thread-safe ingestion without explicit locks on the graph structure.
Thread Isolation and the Writer Pool
The writer employs a specialized thread model that separates active ingestion from historical back-filling. This design ensures that no two threads ever mutate the same segment concurrently.
Live Writer Thread
A single long-running live writer thread continuously consumes new RecosHoseMessage batches and appends edges to the current live graph segment. This thread runs indefinitely until the service receives a shutdown signal, handling the real-time stream of recommendation events.
Catch-up Writer Threads
During service startup, catch-up writers back-fill older graph segments. The number of these short-lived threads is configurable via the catchupWriterNum parameter (defaulting to catchupWriterNum-1 instances). Each catch-up writer is assigned a distinct, immutable segment, ensuring their write sets never overlap.
Executor Service Initialization
All threads are submitted to a cached thread pool created in UnifiedGraphWriter.scala. The pool dynamically manages thread lifecycle while keeping the live and catch-up writers isolated:
val threadPool: ExecutorService = Executors.newCachedThreadPool()
...
threadPool.submit(liveWriter) // live writer runs forever
catchupWriters.map(threadPool.submit(_)) // each catches up a single segment
Source: UnifiedGraphWriter.scala, lines 138‑160
Lock-Free Queue and Back-Pressure Control
Kafka processor threads decouple message consumption from graph mutation by feeding a shared queue. To prevent uncontrolled memory growth under high load, the system uses a bounded semaphore.
The queue is instantiated as a ConcurrentLinkedQueue wrapped with a Semaphore permit limit of 1024:
val queue: java.util.Queue[Array[RecosHoseMessage]] =
new ConcurrentLinkedQueue[Array[RecosHoseMessage]]()
val queuelimit: Semaphore = new Semaphore(1024)
Source: UnifiedGraphWriter.scala, lines 77‑80
When producers enqueue new batches, they must first acquire a permit from queuelimit. If the queue reaches 1024 batches, producers block until writers drain the queue, creating natural back-pressure against Kafka consumer lag.
Edge Buffering and Batch Aggregation
Writer threads do not write individual edges directly to the graph. Instead, they use the BufferedEdgeWriter runnable to pull batches and the BufferedEdgeCollector to aggregate edges before insertion.
The BufferedEdgeWriter polls the shared queue and forwards each message to its assigned EdgeCollector:
while (running) {
val currentBatch = queue.poll
if (currentBatch != null) {
// … iterate over the batch
edgeCollector.addEdge(currentBatch(i))
} else {
Thread.sleep(100L)
}
}
Source: BufferedEdgeWriter.scala, lines 28‑45
The BufferedEdgeCollector accumulates edges in a local array. Once the buffer reaches the configured bufferSize, it acquires a permit from queuelimit and enqueues the full batch back onto the shared queue for the actual graph writer:
if (index >= bufferSize) {
queuelimit.acquireUninterruptibly()
queue.add(oldBuffer)
}
Source: EdgeCollector.scala, lines 29‑38
Segment-Based Write Isolation
The definitive mechanism preventing data conflicts is segment isolation. Each writer thread operates on exactly one graph segment, and segments are immutable once a catch-up writer claims them.
- Live Writer: Calls
addEdgeToGraph(liveGraph, message), writing to the mutable live segment. - Catch-up Writers: Each receives a distinct segment via
liveGraph.getLiveSegmentfollowed byliveGraph.rollForwardSegment(), then callsaddEdgeToSegment(segment, message).
Because rollForwardSegment() atomically advances the live pointer and returns the previous segment, no two threads ever hold a reference to the same writable segment:
val liveEdgeCollector = new EdgeCollector {
override def addEdge(message: RecosHoseMessage): Unit = addEdgeToGraph(liveGraph, message)
}
...
val catchupEdgeCollector = new EdgeCollector {
override def addEdge(message: RecosHoseMessage): Unit = addEdgeToSegment(segment, message)
}
Source: UnifiedGraphWriter.scala, lines 66‑70 and lines 85‑93
Graceful Shutdown Coordination
Clean termination relies on an AtomicBoolean flag named isRunning. All writer threads evaluate isRunning.get() in their main loops. The shutdown() method flips this flag, closes Kafka consumers, and awaits thread pool termination:
isRunning.set(false) // signals writers to stop
threadPool.shutdown()
Source: UnifiedGraphWriter.scala, lines 60‑71
Once the flag is cleared, writers exit their while (running) loops, drain remaining buffered edges, and terminate without corrupting partially written segments.
Concrete Implementation Example
To implement a custom graph writer, extend UnifiedGraphWriter and define how RecosHoseMessage objects map to graph edges:
import com.twitter.recos.hose.common._
import com.twitter.graphjet.bipartite._
import com.twitter.recos.internal.thriftscala.RecosHoseMessage
class UserTweetGraphWriter extends UnifiedGraphWriter[
LeftIndexedBipartiteGraphSegment,
MultiSegmentPowerLawBipartiteGraph[LeftIndexedBipartiteGraphSegment]
] {
// Configuration parameters
override val shardId = "utg-shard-01"
override val env = "prod"
override val hosename = "recos-hose"
override val bufferSize = 5000
override val consumerNum = 4
override val catchupWriterNum = 3
override val clientId = "user-tweet-writer"
// Live graph mutation
override def addEdgeToGraph(
graph: MultiSegmentPowerLawBipartiteGraph[LeftIndexedBipartiteGraphSegment],
msg: RecosHoseMessage
): Unit = {
graph.addEdge(msg.sourceId, msg.destinationId, msg.weight)
}
// Historical segment mutation
override def addEdgeToSegment(
segment: LeftIndexedBipartiteGraphSegment,
msg: RecosHoseMessage
): Unit = {
segment.addEdge(msg.sourceId, msg.destinationId, msg.weight)
}
}
// Runtime wiring
val liveGraph = MultiSegmentPowerLawBipartiteGraph.create(/* config */)
val writer = new UserTweetGraphWriter()
writer.initHose(liveGraph) // launches Kafka processors and writer threads
// ... application runs ...
writer.shutdown() // triggers graceful termination
Summary
- Segment Isolation: Each writer thread mutates only its assigned graph segment, preventing cross-thread write conflicts.
- Lock-Free Queue: A
ConcurrentLinkedQueuepaired with a 1024-permitSemaphoreprovides thread-safe handoff and memory-bound back-pressure. - Dual Writer Model: One live writer handles real-time ingestion while configurable catch-up writers back-fill historical segments.
- Atomic Coordination: An
AtomicBooleanflag ensures writers detect shutdown signals and terminate cleanly without data loss.
Frequently Asked Questions
How does UnifiedGraphWriter prevent two threads from writing to the same graph segment?
Each catch-up writer receives a unique segment reference obtained via liveGraph.rollForwardSegment() before it starts processing. The live writer exclusively accesses the mutable live segment via addEdgeToGraph, while catch-up writers exclusively access their frozen segments via addEdgeToSegment. Since segments are never shared between threads, write conflicts are architecturally impossible.
What prevents memory overflow if Kafka produces faster than the graph can consume?
The Semaphore named queuelimit caps the shared queue at 1024 batches. Kafka processor threads must call acquireUninterruptibly() before enqueuing new data. When the queue is full, producers block until writers drain batches, creating back-pressure that throttles Kafka consumption to match graph ingestion capacity.
Can the number of concurrent catch-up writers be configured?
Yes. The catchupWriterNum override parameter controls how many catch-up threads spawn during initialization. Setting this to N creates N-1 catch-up writers plus one live writer. This allows tuning based on available CPU cores and the number of historical segments requiring back-fill.
How does the system handle service shutdown without losing in-flight edges?
The shutdown() method sets the AtomicBoolean isRunning to false, which signals all BufferedEdgeWriter loops to exit their while (running) blocks. Threads then complete their current batch processing, submit any remaining buffered edges to the graph, and terminate before the ExecutorService shuts down, ensuring no partial writes remain uncommitted.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →