How to Extend and Customize the EdgeCollector in Twitter's Recommendation Algorithm
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:
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 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 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:
// 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:
// 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:
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:
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:
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:
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-algorithmthat defines a singleaddEdge(message: RecosHoseMessage): Unitmethod, serving as the primary extension point for edge processing. - BufferedEdgeCollector provides production-ready batching and back-pressure via configurable
bufferSizeand semaphore-basedqueuelimitparameters. - 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
EdgeCollectorand delegates to an underlying collector, then wiring it intogetLiveWriterorgetCatchupWriter. - 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 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.
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 →