# How to Extend and Customize the EdgeCollector in Twitter's Recommendation Algorithm

> Learn how to extend the EdgeCollector in Twitters recommendation algorithm. Inject custom logic into writer pipelines by implementing the addEdge method. Customize edge processing with ease.

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

---

**Developers can extend the `EdgeCollector` trait to inject custom filtering, transformation, or metrics logic into the live and catch-up writer pipelines by implementing the `addEdge` method and wiring the custom instance into `UnifiedGraphWriter`.**

The `EdgeCollector` abstraction in the `twitter/the-algorithm` repository defines the contract for processing `RecosHoseMessage` edges before they are persisted to the bipartite graph. Because the trait is minimal and unopinionated, it serves as an extension point for domain-specific edge handling without modifying core pipeline infrastructure.

## Understanding the EdgeCollector Trait and Core Contract

The entire contract is defined in [`src/scala/com/twitter/recos/hose/common/EdgeCollector.scala`](https://github.com/twitter/the-algorithm/blob/main/src/scala/com/twitter/recos/hose/common/EdgeCollector.scala):

```scala
trait EdgeCollector {
  def addEdge(message: RecosHoseMessage): Unit
}

```

Any class implementing this trait can be substituted into the writer pipeline. The `BufferedEdgeWriter` (the consumer thread) continuously polls the queue and invokes `edgeCollector.addEdge` for each message, as seen in [`BufferedEdgeWriter.scala`](https://github.com/twitter/the-algorithm/blob/main/BufferedEdgeWriter.scala) lines 13-18.

## Built-in Implementation: BufferedEdgeCollector

`BufferedEdgeCollector` provides a reusable buffering strategy that batches edges to optimize write throughput. It accumulates edges in an in-memory array of size `bufferSize`. When the buffer fills, it pushes the batch onto a bounded concurrent queue (`java.util.Queue[Array[RecosHoseMessage]]`) that is later consumed by `BufferedEdgeWriter`.

Key configuration points available in the constructor:

| Parameter | Purpose |
|-----------|---------|
| `bufferSize` | Number of edges to accumulate before enqueueing a batch. |
| `queue` | Shared `ConcurrentLinkedQueue` that decouples producers from consumers. |
| `queuelimit` | `Semaphore` controlling maximum queue depth to provide back-pressure. |
| `statsReceiver` | Metrics scope for tracking `queueAdd` and `waitEnqueue` latencies. |

## Where EdgeCollectors Are Instantiated in the Pipeline

[`UnifiedGraphWriter.scala`](https://github.com/twitter/the-algorithm/blob/main/UnifiedGraphWriter.scala) orchestrates the live and catch-up writer threads. Each thread receives its own `EdgeCollector` implementation tailored to its specific requirements.

### Live Writer Collector

The live path uses an anonymous class that immediately forwards edges to the current graph segment:

```scala
// UnifiedGraphWriter#getLiveWriter (lines 66-70)
val liveEdgeCollector = new EdgeCollector {
  override def addEdge(message: RecosHoseMessage): Unit = addEdgeToGraph(liveGraph, message)
}

```

### Catch-up Writer Collector

The catch-up path tracks how many edges have been processed for a historic segment and stops when the segment reaches capacity:

```scala
// UnifiedGraphWriter#getCatchupWriter (lines 85-93)
val catchupEdgeCollector = new EdgeCollector {
  var currentNumEdges = 0
  override def addEdge(message: RecosHoseMessage): Unit = {
    currentNumEdges += 1
    addEdgeToSegment(segment, message)
  }
}

```

## How to Extend EdgeCollector for Custom Processing

Because `EdgeCollector` is a plain trait, developers can inject any custom logic by implementing `addEdge` and delegating to an underlying collector. Typical customizations include filtering spam, enriching messages, or emitting side-effect metrics.

### Filtering Edges Based on Business Rules

Create a decorator that drops edges from blacklisted users before forwarding to the graph writer:

```scala
class FilteringEdgeCollector(
  underlying: EdgeCollector,
  stats: StatsReceiver,
  blacklist: Set[Long]
) extends EdgeCollector {

  private val filteredCounter = stats.counter("edges_filtered")
  private val processedCounter = stats.counter("edges_processed")

  override def addEdge(message: RecosHoseMessage): Unit = {
    val userId = message.sourceUserId
    if (blacklist.contains(userId)) {
      filteredCounter.incr()
    } else {
      processedCounter.incr()
      underlying.addEdge(message)
    }
  }
}

```

### Transforming Messages Before Insertion

Enrich edges with derived fields or normalize data formats before persistence:

```scala
class EnrichingCollector(inner: EdgeCollector, stats: StatsReceiver) extends EdgeCollector {
  private val enriched = stats.counter("edges_enriched")
  
  override def addEdge(msg: RecosHoseMessage): Unit = {
    val enrichedMsg = msg.copy(customFlag = true)
    enriched.incr()
    inner.addEdge(enrichedMsg)
  }
}

```

### Wiring Custom Collectors into the Pipeline

Replace the anonymous collector in `getLiveWriter` with your decorated implementation:

```scala
override def getLiveWriter(
  liveGraph: TGraph,
  queue: java.util.Queue[Array[RecosHoseMessage]],
  queuelimit: Semaphore
): BufferedEdgeWriter = {

  val baseCollector = new EdgeCollector {
    override def addEdge(message: RecosHoseMessage): Unit =
      addEdgeToGraph(liveGraph, message)
  }

  val customCollector = new FilteringEdgeCollector(
    baseCollector,
    statsReceiver.scope("liveWriterFiltering"),
    Set(12345L, 67890L)
  )

  BufferedEdgeWriter(
    queue,
    queuelimit,
    customCollector,
    statsReceiver.scope("liveWriter"),
    isRunning.get
  )
}

```

The same pattern applies to catch-up writers: wrap the collector that tracks `currentNumEdges` with your custom logic, ensuring the counter increments remain accurate.

## Tuning Buffer Sizes and Queue Limits Without Code Changes

Both `bufferSize` and the semaphore `queuelimit` are constructor arguments of `BufferedEdgeCollector`. Adjust these via Finatra/Finagle configuration flags or `.conf` files to adapt to traffic patterns:

```scala
val largeBatchCollector = BufferedEdgeCollector(
  bufferSize   = 10000,  // Larger batches for high-throughput
  queue        = new ConcurrentLinkedQueue[Array[RecosHoseMessage]](),
  queuelimit   = new Semaphore(4096),  // Deeper queue for bursts
  statsReceiver
)

```

Increasing `bufferSize` reduces queue contention but increases latency per batch. Increasing the semaphore permits allows more batches to queue during traffic spikes, preventing producer blocking.

## Summary

- **EdgeCollector** is a minimal Scala trait in `twitter/the-algorithm` that defines a single `addEdge(message: RecosHoseMessage): Unit` method, serving as the primary extension point for edge processing.
- **BufferedEdgeCollector** provides production-ready batching and back-pressure via configurable `bufferSize` and semaphore-based `queuelimit` parameters.
- **UnifiedGraphWriter** instantiates distinct collectors for live and catch-up writers; developers can replace the anonymous implementations with custom decorators.
- **Custom logic** (filtering, transformation, metrics) is implemented by creating a class that extends `EdgeCollector` and delegates to an underlying collector, then wiring it into `getLiveWriter` or `getCatchupWriter`.
- **Performance tuning** requires only configuration changes to buffer sizes and queue limits, no code modifications.

## Frequently Asked Questions

### What is the EdgeCollector trait in Twitter's recommendation algorithm?

The `EdgeCollector` trait is a minimal abstraction defined in [`src/scala/com/twitter/recos/hose/common/EdgeCollector.scala`](https://github.com/twitter/the-algorithm/blob/main/src/scala/com/twitter/recos/hose/common/EdgeCollector.scala) that specifies a single method `addEdge(message: RecosHoseMessage): Unit`. It acts as the integration point between the Kafka message stream and the bipartite graph storage, allowing developers to intercept and process edges before they are persisted.

### How does BufferedEdgeCollector handle back-pressure?

`BufferedEdgeCollector` uses a `java.util.concurrent.Semaphore` passed as `queuelimit` to cap the number of buffered batches waiting in the queue. When the semaphore permits are exhausted, the collector blocks on `queuelimit.acquire()`, preventing memory exhaustion during traffic spikes and providing natural back-pressure to upstream producers.

### Can I use multiple custom EdgeCollectors simultaneously?

Yes. The pipeline architecture in `UnifiedGraphWriter` creates separate writer threads for live and catch-up segments, each with its own `EdgeCollector` instance. You can provide distinct implementations for each path—for example, a filtering collector for the live writer and a counting collector for catch-up—or instantiate multiple decorated collectors within the same JVM for different graph segments.

### Where should I implement edge filtering logic?

Implement filtering logic inside a custom `EdgeCollector` implementation that wraps the base collector. Override `addEdge` to inspect the `RecosHoseMessage`, apply your business rules (such as checking against a blacklist), and only delegate to the underlying collector for edges that pass the filter. Wire this custom collector into `UnifiedGraphWriter.getLiveWriter` or `getCatchupWriter` to activate it in the production pipeline.