# How UnifiedGraphWriter Handles Concurrent Writes and Data Conflicts in Twitter's Algorithm

> Discover how UnifiedGraphWriter prevents data conflicts with thread isolation and lock-free queues, ensuring safe, high-volume Kafka stream ingestion for Twitter's algorithm. Learn its concurrency strategy.

- Repository: [X (fka Twitter)/the-algorithm](https://github.com/twitter/the-algorithm)
- Tags: internals
- Published: 2026-03-03

---

**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`](https://github.com/twitter/the-algorithm/blob/main/UnifiedGraphWriter.scala). The pool dynamically manages thread lifecycle while keeping the live and catch-up writers isolated:

```scala
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`](https://github.com/twitter/the-algorithm/blob/main/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:

```scala
val queue: java.util.Queue[Array[RecosHoseMessage]] =
  new ConcurrentLinkedQueue[Array[RecosHoseMessage]]()
val queuelimit: Semaphore = new Semaphore(1024)

```

*Source: [`UnifiedGraphWriter.scala`](https://github.com/twitter/the-algorithm/blob/main/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`:

```scala
while (running) {
  val currentBatch = queue.poll
  if (currentBatch != null) {
    // … iterate over the batch
    edgeCollector.addEdge(currentBatch(i))
  } else {
    Thread.sleep(100L)
  }
}

```

*Source: [`BufferedEdgeWriter.scala`](https://github.com/twitter/the-algorithm/blob/main/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:

```scala
if (index >= bufferSize) {
  queuelimit.acquireUninterruptibly()
  queue.add(oldBuffer)
}

```

*Source: [`EdgeCollector.scala`](https://github.com/twitter/the-algorithm/blob/main/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.getLiveSegment` followed by `liveGraph.rollForwardSegment()`, then calls `addEdgeToSegment(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:

```scala
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`](https://github.com/twitter/the-algorithm/blob/main/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:

```scala
isRunning.set(false)   // signals writers to stop
threadPool.shutdown()

```

*Source: [`UnifiedGraphWriter.scala`](https://github.com/twitter/the-algorithm/blob/main/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:

```scala
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 `ConcurrentLinkedQueue` paired with a 1024-permit `Semaphore` provides 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 `AtomicBoolean` flag 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.