Writing

Apache Kafka along with RxJS operators

How to buffer received data while the producer is unavailable?

Summary

On one project I owned a processor sitting between two Kafka topics: consume from the first, validate, publish to the second. The transformation was the easy part. The hard part was waiting for the producer, which could be paused or disconnected from time to time, without losing a single message consumed in the meantime.

RxJS merge operators look interchangeable and are not. A takeUntil followed by a repeatWhen resubscribes when the producer comes back, but everything received during the pause is gone. A switchMap that returns an empty source while the producer is away loses the same data. A combineLatest does barely better, with a trap on top: filter directly on the producer and the consumer stops noticing later disconnections.

The answer is a pair of twin operators. windowToggle says when to emit and when to ignore; bufferToggle says when to buffer. The consumer stream switches between the two states according to the producer’s availability, then both results are merged and flattened so that each value is published one at a time. Data is consumed when the producer is back and kept when it stops.

The full example runs on NestJS, node-rdkafka and RxJS, and the repository that goes with the article replays the scenario. The case is narrow; the lesson is that choosing a merge operator is a decision about what you are willing to lose.

Key ideas

  • The hard part of a Kafka processor is not the transformation but waiting for the producer without losing what the consumer receives meanwhile.
  • takeUntil with repeatWhen, switchMap and combineLatest all resume when the producer comes back, and all drop the data received during the pause.
  • windowToggle decides when to emit, bufferToggle when to buffer; merged then flattened, they let nothing fall through.
  • Filtering directly on the producer is a trap: once it is available, the consumer ignores every later disconnection.

Why I wrote this

In this tutorial I wanted to explain a few merge operators and the problem I had run into shortly before with Kafka and RxJS: the project I was working on had a processing function between an incoming topic and an outgoing one.

Companion repositories