# How Kafka Streams Are Integrated and Utilized for State Persistence in Twitter's Recommendation Algorithm Services

> Discover how Twitter's recommendation algorithm integrates Kafka Streams for state persistence. Learn about stateful processing, RocksDB, and changelog topics in this technical deep dive.

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

---

**Twitter's recommendation algorithm uses Finatra-based Kafka Streams to build stateful processing pipelines for Unified User Actions (UUA), with the `EnrichmentPlannerService` and `EnricherService` handling topology configuration, offset persistence, and repartition state through built-in RocksDB stores and changelog topics.**

The `twitter/the-algorithm` repository implements a sophisticated event-driven architecture where Kafka Streams provide the backbone for real-time data enrichment. The system leverages Finatra's `SecureKafkaStreamsConfig` wrapper to manage consumer groups, producer acknowledgments, and local state stores that persist processing state across service restarts and rebalances.

## Core Services and Topology Architecture

The recommendation algorithm embeds Kafka Streams topologies within two primary services that process Unified User Actions sequentially.

### EnrichmentPlannerService as the Entry Point

The `EnrichmentPlannerService` serves as the first stage of the UUA enrichment pipeline, defined in [`unified_user_actions/service/src/main/scala/com/twitter/unified_user_actions/service/EnrichmentPlannerService.scala`](https://github.com/twitter/the-algorithm/blob/main/unified_user_actions/service/src/main/scala/com/twitter/unified_user_actions/service/EnrichmentPlannerService.scala). This service extends `SecureKafkaStreamsConfig` and implements the `configureKafkaStreams` method to declare the streaming topology.

The topology reads from the raw UUA topic using `UnKeyedSerde` for keys and `ScalaSerdes.Thrift[UnifiedUserAction]` for values. It transforms each record into an `EnrichmentEnvelop` with a repartition key (`EnrichmentKey`), applies Decider-based sampling for feature gating, and writes to the output topic. The `flatMapValues` operation and subsequent repartitioning create internal stateful operations that Kafka Streams persists to local RocksDB stores and changelog topics.

### EnricherService for Asynchronous Hydration

The `EnricherService`, located in [`unified_user_actions/service/src/main/scala/com/twitter/unified_user_actions/service/EnricherService.scala`](https://github.com/twitter/the-algorithm/blob/main/unified_user_actions/service/src/main/scala/com/twitter/unified_user_actions/service/EnricherService.scala), constitutes the second stage. It continues the stream processing by consuming the repartitioned records and performing asynchronous hydration through `flatMapAsync`.

This service configures a commit interval of 5 seconds with 10,000 worker threads to handle the enrichment driver execution. The `flatMapAsync` method materializes stateful processing guarantees, ensuring that in-flight enrichment plans are checkpointed when tasks are reassigned. The service routes the final enriched payloads to the destination topics using `Produced` serializers for `EnrichmentKey` and `EnrichmentEnvelop`.

## Kafka Streams Configuration and State Persistence

State persistence in the recommendation algorithm relies on Finatra's configuration abstractions and Kafka Streams' built-in state store mechanisms.

### SecureKafkaStreamsConfig Implementation

Both services override `streamsProperties` to inject custom `KafkaStreamsConfig` instances. The configuration sets the consumer group ID using `KafkaGroupId(ApplicationId)`, client IDs for both consumer and producer, and critical producer acknowledgments via `AckMode.ALL`. The configuration also specifies `CompressionType.LZ4` for efficient network utilization.

The `SecureKafkaStreamsConfig` wrapper handles security credentials and connects to the underlying `KafkaStreamsConfig` builder pattern. When the Finatra service starts, it instantiates a `KafkaStreams` object using this configuration, which automatically manages local state store directories and changelog topics for fault tolerance.

### State Persistence Mechanics

The recommendation algorithm leverages three layers of state persistence within Kafka Streams:

**Consumer offsets** are managed through the internal `__consumer_offsets` changelog topic. The explicit `groupId` configuration (`uua-enrichment-planner` or `uua-enricher`) scopes these offsets to the specific service instance, enabling accurate recovery after restarts.

**Repartition state** occurs when the topology transforms `UnKeyed` records to keyed `EnrichmentKey` records. Kafka Streams creates internal repartition topics named `<appId>-<operator>-repartition` and maintains local RocksDB stores to buffer these records. The state is periodically checkpointed to the changelog topics.

**Processing guarantees** are enforced through `AtLeastOnceProcessor` wrappers defined in `KafkaProcessorRekeyUuaModule`. The `AckMode.ALL` configuration ensures records are only considered processed after upstream offsets are committed, while the commit intervals in `EnricherService` control the trade-off between latency and durability.

## Integration with the Broader Recommendation Stack

The Kafka Streams services integrate tightly with Twitter's recommendation infrastructure through the Unified User Actions pipeline.

The UUA stream serves as the primary ingestion point for all downstream recommendation systems, including home-timeline and video recommendation pipelines. The `EnrichmentPlannerService` adds deterministic repartition keys that enable downstream services to perform efficient joins without additional lookups.

The Decider-driven sampling mechanism allows fine-grained rollout of new enrichment logic per key, enabling A/B testing and gradual feature deployment without service redeployment. This stateful processing ensures that sampling decisions persist across consumer group rebalances.

## Code Examples

### Minimal Finatra-Kafka-Streams Service Skeleton

```scala
import com.twitter.finatra.kafkastreams.config.{KafkaStreamsConfig, SecureKafkaStreamsConfig}
import com.twitter.finatra.kafkastreams.dsl.FinatraDslToCluster
import org.apache.kafka.streams.StreamsBuilder
import org.apache.kafka.streams.scala.kstream.{Consumed, Produced}
import com.twitter.finatra.kafka.serde.{ScalaSerdes, UnKeyedSerde}
import com.twitter.unified_user_actions.thriftscala.{UnifiedUserAction, EnrichmentKey, EnrichmentEnvelop}

object MyUuaServiceMain extends MyUuaService {
  val ApplicationId = "my-uua-service"
  val InputTopic   = "unified_user_actions"
  val OutputTopic  = "unified_user_actions_keyed"
}
class MyUuaService extends FinatraDslToCluster with SecureKafkaStreamsConfig {
  import MyUuaServiceMain._

  override protected def configureKafkaStreams(builder: StreamsBuilder): Unit = {
    builder.asScala
      .stream(InputTopic)(Consumed.`with`(UnKeyedSerde,
                                          ScalaSerdes.Thrift[UnifiedUserAction]))
      .mapValues(uua => EnrichmentEnvelop(/* … build … */))
      .to(OutputTopic)(Produced.`with`(ScalaSerdes.Thrift[EnrichmentKey],
                                        ScalaSerdes.Thrift[EnrichmentEnvelop]))
  }

  override def streamsProperties(cfg: KafkaStreamsConfig): KafkaStreamsConfig = {
    super.streamsProperties(cfg)
      .consumer.groupId(KafkaGroupId(ApplicationId))
      .producer.ackMode(AckMode.ALL)
  }
}

```

### Adding a Custom State Store

```scala
import org.apache.kafka.streams.state.{Stores, KeyValueStore}
import org.apache.kafka.common.utils.Bytes
import org.apache.kafka.streams.scala.kstream.{Consumed, Produced, KTable}

// Inside configureKafkaStreams
val storeName = "tweet-id-to-user-store"
val storeSupplier = Stores.keyValueStoreBuilder(
  Stores.persistentKeyValueStore(storeName),
  ScalaSerdes.Long,
  ScalaSerdes.Thrift[UserId]
).withCachingEnabled()

builder.addStateStore(storeSupplier)

builder.asScala
  .stream(InputTopic)(Consumed.`with`(UnKeyedSerde,
                                      ScalaSerdes.Thrift[UnifiedUserAction]))
  .transformValues(() => new MyTransformer(storeName), storeName) // access store
  .to(OutputTopic)(Produced.`with`(...))

```

> *The code above illustrates how a service could add an explicit RocksDB-backed `KeyValueStore` that would be automatically persisted to a changelog topic by Kafka Streams.*

## Summary

- **Finatra Integration**: The recommendation algorithm uses `SecureKafkaStreamsConfig` and `FinatraDslToCluster` to embed Kafka Streams topologies within `EnrichmentPlannerService` and `EnricherService`.

- **Topology Design**: The `configureKafkaStreams` method defines processing pipelines that read raw UUA events, transform them into `EnrichmentEnvelop` objects with deterministic keys, and repartition for downstream consumption.

- **State Persistence**: Consumer offsets are managed via `__consumer_offsets` with scoped `groupId` values, while repartition state uses internal RocksDB stores and changelog topics. `AckMode.ALL` ensures at-least-once processing guarantees.

- **Async Processing**: The `EnricherService` utilizes `flatMapAsync` with configurable commit intervals to handle high-throughput hydration without blocking the stream topology.

- **Configuration Management**: Custom `streamsProperties` override consumer group IDs, client IDs, compression types (LZ4), and acknowledgment modes to optimize for the recommendation workload.

## Frequently Asked Questions

### How does Twitter's recommendation algorithm use Kafka Streams for stateful processing?

The algorithm embeds Kafka Streams within Finatra services (`EnrichmentPlannerService` and `EnricherService`) to process Unified User Actions. The topology uses `flatMapValues` and `flatMapAsync` transformations that maintain state through Kafka Streams' internal RocksDB stores and changelog topics, enabling fault-tolerant repartitioning and offset management.

### What is the role of SecureKafkaStreamsConfig in the recommendation services?

`SecureKafkaStreamsConfig` is a Finatra wrapper around `KafkaStreamsConfig` that provides secure credential injection and configuration defaults. Services extend this trait to override `streamsProperties`, setting consumer group IDs, producer acknowledgment modes (`AckMode.ALL`), and compression types (LZ4) while inheriting security policies for Kafka authentication.

### How does the algorithm ensure state persistence across service restarts?

State persistence relies on three mechanisms: consumer offsets stored in the internal `__consumer_offsets` topic scoped by `groupId`, repartition state maintained in local RocksDB stores backed by changelog topics, and at-least-once processing semantics enforced by `AtLeastOnceProcessor` wrappers. When a service restarts, Kafka Streams restores the state stores from changelog topics and resumes consumption from the last committed offset.

### What is the difference between EnrichmentPlannerService and EnricherService?

`EnrichmentPlannerService` acts as the first stage, reading raw UUA events and transforming them into `EnrichmentEnvelop` objects with deterministic repartition keys. `EnricherService` serves as the second stage, consuming the repartitioned stream and performing asynchronous hydration using `flatMapAsync` with configurable worker pools and commit intervals before routing to final destination topics.