Scaling GraphJet-Based Services to Handle Millions of Users and Interactions
Twitter's GraphJet library achieves horizontal scalability for recommendation services through a multi-segment power-law bipartite graph architecture, Kafka-backed parallel ingestion pipelines, and deterministic sharding, allowing production deployments to handle tens of GB/s of interaction events while maintaining bounded memory usage.
The twitter/the-algorithm repository open-sources the recommendation stack that powers Twitter's real-time suggestion systems. At its core, GraphJet (com.twitter.graphjet) serves as the in-memory bipartite graph engine that stores user-to-content interactions. Scaling GraphJet-based services to millions of users requires careful tuning of segment sizing, power-law distributions, and write parallelism to prevent garbage collection pauses and memory pressure. The following sections detail the architectural patterns and configuration parameters that enable this scale.
Core Architecture: Multi-Segment Power-Law Graphs
The MultiSegmentPowerLawBipartiteGraph implementation in src/scala/com/twitter/recos/graph_common/MultiSegmentPowerLawBipartiteGraphBuilder.scala partitions the graph into a series of time-ordered segments. This design isolates write traffic and limits memory pressure through three key properties:
- Immutable Historical Segments: Older segments become read-only, allowing the JVM to optimize garbage collection for the single mutable segment.
- Sliding Window Retention: Once
maxNumSegmentsis reached, the oldest segment is dropped, bounding total memory usage regardless of total historical data volume. - Power-Law Edge Pools: Each segment allocates edge storage using power-law distributions, ensuring O(1) random access while efficiently handling both low-activity users and viral content with millions of interactions.
Tuning Graph Builder Configuration
Production scaling depends on proper sizing of the GraphBuilderConfig defined in src/scala/com/twitter/recos/user_video_graph/UserVideoGraphConfig.scala. The following parameters control memory layout and ingestion capacity:
Segment Sizing and Counts
maxNumSegments: Determines the look-back window (e.g., 8 segments ≈ 2 days of data at typical ingestion rates). Increasing this linearly increases memory usage.maxNumEdgesPerSegment: Sets the edge capacity per segment (e.g.,1 << 28≈ 268 million edges ≈ 1 GiB of raw storage). Size this to accommodate peak ingestion bursts plus a safety factor.
Power-Law Exponents
The leftPowerLawExponent and rightPowerLawExponent parameters control edge pool allocation strategies:
leftPowerLawExponent = 16.0: A steep exponent for the user side, optimizing for the typical case where most users have few interactions.rightPowerLawExponent = 4.0: A shallower exponent for the content side, reserving capacity for viral tweets that accumulate millions of edges.
Adjust these exponents if monitoring shows "edge pool overflow" metrics, indicating that the distribution assumptions no longer match the traffic pattern.
Efficient Edge Encoding with Bitmasks
To maintain a dense graph representation while supporting multiple interaction types, GraphJet uses the UserVideoEdgeTypeMask in src/scala/com/twitter/recos/user_video_graph/UserVideoEdgeTypeMask.scala. This implementation packs the interaction type (e.g., video playback, retweet) into the high-order bits of the right-node ID.
- Density: Uses a single
intper edge regardless of interaction type diversity. - Capacity: Reserves 4 bits for edge types, leaving 27+ bits for node IDs.
- Performance: Enables fast filtering by edge type during graph traversals without secondary indexes.
Horizontal Scaling via Kafka-Backed Writers
The ingestion pipeline uses UserVideoGraphWriter in src/scala/com/twitter/recos/user_video_graph/UserVideoGraphWriter.scala, which extends UnifiedGraphWriter from src/scala/com/twitter/recos/hose/common/UnifiedGraphWriter.scala. This architecture decouples ingestion from graph updates and enables horizontal scaling through three mechanisms:
Live Ingestion
Live writers consume RecosHoseMessages from Kafka and call addEdgeToGraph, inserting edges into the mutable segment. The consumerNum parameter (typically 4) controls the parallelism per shard.
Catch-Up Writers
After service restarts, catchupWriterNum threads (defaulting to maxNumSegments - 1) replay historic events into older segments via addEdgeToSegment. This ensures the graph rebuilds without warm-up latency spikes, processing approximately 25 MiB/s per consumer.
Sharding Strategy
Each writer runs inside a logical shard (shardId argument). The service launches one writer per shard, each managing its own MultiSegmentPowerLawBipartiteGraph instance. Total cluster capacity equals #shards × (maxNumSegments × segmentSize), allowing linear horizontal scaling by adding machines and Kafka partitions.
Memory Management and JVM Tuning
Production deployments require specific JVM configurations to handle the large in-memory graphs:
- Garbage Collection: Enable G1GC with
-XX:MaxGCPauseMillis=200to minimize pause times during high-throughput ingestion. - Heap Sizing: Allocate heap based on
expectedNumLeftNodes * 4 bytes + expectedNumRightNodes * 4 bytesper segment, ensuring the working set fits comfortably below 70% of available heap to avoid evacuation failures.
Observability and Monitoring
GraphJet exposes detailed statistics via the StatsReceiver interface passed through the builder configuration. Critical metrics to monitor include:
- Segment metrics:
segmentCount,edgeCount, andevictedSegments - Ingestion rates: Bytes processed per second per consumer thread
- Error rates: Edge pool overflow events and Kafka lag
Set alerts on evictedSegments spikes, which indicate that the retention window is too short for the current ingestion volume or that maxNumSegments requires adjustment.
Code Examples
Building a Graph with Custom Configuration
import com.twitter.recos.graph_common.MultiSegmentPowerLawBipartiteGraphBuilder
import com.twitter.recos.graph_common.MultiSegmentPowerLawBipartiteGraphBuilder.GraphBuilderConfig
import com.twitter.recos.user_video_graph.RecosConfig
import com.twitter.graphjet.stats.NullStatsReceiver
val cfg: GraphBuilderConfig = RecosConfig.graphBuilderConfig
val graph = MultiSegmentPowerLawBipartiteGraphBuilder(cfg, NullStatsReceiver)
Writing an Edge (Live Writer)
def addEdge(graph: MultiSegmentPowerLawBipartiteGraph,
leftId: Long,
rightId: Long,
action: Byte,
card: Option[Byte]): Unit = {
val metaEdge = card match {
case Some(c) if isPhotoCard(c) => TweetIDMask.photo(rightId)
case Some(c) if isPlayerCard(c) => TweetIDMask.player(rightId)
case Some(c) if isSummaryCard(c) => TweetIDMask.summary(rightId)
case Some(c) if isPromotionCard(c)=> TweetIDMask.promotion(rightId)
case _ => rightId
}
val edgeType = UserVideoEdgeTypeMask.actionTypeToEdgeType(action)
graph.addEdge(leftId, metaEdge, edgeType)
}
Initializing a Writer for a Shard
val writer = UserVideoGraphWriter(
shardId = "shard-01",
env = "prod",
hosename = "recos-host",
bufferSize = 10000,
kafkaConsumerBuilder = myConsumerBuilder,
clientId = "user-video-writer",
statsReceiver = myStatsReceiver
)
writer.start() // spawns consumerNum live + catchupWriterNum catch‑up threads
Summary
Scaling GraphJet-based services to millions of users requires a multi-layered approach combining data structure optimization, horizontal sharding, and careful JVM tuning:
- Segmented Architecture: The
MultiSegmentPowerLawBipartiteGraphuses immutable historical segments and a single mutable segment to bound memory usage and isolate write traffic. - Power-Law Tuning: Configuring
leftPowerLawExponentandrightPowerLawExponentensures efficient storage for both typical users and viral content. - Horizontal Scaling: Kafka-backed writers with configurable sharding (
shardId) and parallel consumer threads (consumerNum) enable linear capacity expansion. - Operational Resilience: Catch-up writers (
catchupWriterNum) ensure fast recovery after restarts without warm-up latency spikes.
Frequently Asked Questions
How does GraphJet prevent unbounded memory growth when handling millions of daily interactions?
GraphJet implements a sliding window retention policy through its multi-segment architecture. The MultiSegmentPowerLawBipartiteGraph maintains a fixed number of segments (maxNumSegments), and when the limit is reached, the oldest immutable segment is dropped. This bounds memory usage to #shards × (maxNumSegments × segmentSize) regardless of total historical data volume.
What is the purpose of power-law exponents in GraphJet configuration?
The power-law exponents (leftPowerLawExponent and rightPowerLawExponent) control how edge storage pools are allocated within each segment. A steep exponent (e.g., 16.0) optimizes for the long-tail distribution typical of user interactions, while a shallower exponent (e.g., 4.0) reserves capacity for high-degree nodes representing viral content. Tuning these prevents "edge pool overflow" errors during traffic spikes.
How does GraphJet handle service restarts without losing the real-time graph state?
GraphJet uses catch-up writers that replay historical Kafka events into older segments after a restart. The catchupWriterNum parameter (typically set to maxNumSegments - 1) spawns dedicated threads that consume from earlier offsets, calling addEdgeToSegment to rebuild immutable segments. This eliminates warm-up latency spikes and restores the full sliding window without requiring persistent graph snapshots.
When should I increase the number of shards versus tuning segment sizes?
Increase shard count when per-shard memory usage exceeds 70% of available heap or when CPU saturation occurs during edge insertion. Increase segment size (maxNumEdgesPerSegment) or segment count (maxNumSegments) when the retention window is too short for your use case or when evictedSegments metrics spike frequently. Sharding provides horizontal scalability across machines, while segment tuning optimizes single-node resource utilization.
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 →