Topic 8.1
Reliable Updates with Outbox, CDC, and Kafka
In one line
Feed filters from the database's change stream: an outbox row or CDC event per insert goes to Kafka; a filter-updater consumer group applies idempotent adds to Redis or to local filters in every instance. You get retries, replay for rebuilds, multiple consumers and monitoring through consumer lag.
Think of it like this
A newsroom wire service. Every new story goes out on the wire once; each branch office (filter) receives it and updates its own index. If a branch was offline, it catches up from the wire's archive.
Key ideas
- 01
Pipeline: PostgreSQL (business write + outbox row in one transaction) → Debezium CDC → Kafka topic
products.createdkeyed by ID → Bloom updater group →BF.MADDinto Redis (shared) or each instance's local filter (every instance in its own group, or a broadcast pattern). - 02
Idempotency is free: adding an item twice sets the same bits, so at-least-once delivery is fine. Deletes (for deletable filters) need ordering per key: key events by ID so create and delete for one ID stay ordered.
- 03
Retries and DLT: transient Redis errors retry with backoff; malformed events go to a DLT with alerts. A dropped event is a potential false negative, so never skip silently.
- 04
Replay: rebuilding means starting a new consumer group (or resetting offsets) from the snapshot position into a new filter version.
- 05
Lag is correctness-relevant: the updater's lag is the window where new items may be missing. Expose the watermark (latest applied ID or timestamp) so readers can bypass the filter for newer items.
Code & diagrams
@KafkaListener(topics = "products.created", groupId = "bloom-updater-v7", batch = "true")
public void onCreated(List<ConsumerRecord<String, ProductCreated>> batch, Acknowledgment ack) {
String[] args = new String[batch.size() + 1];
args[0] = "bf:products:v7";
for (int i = 0; i < batch.size(); i++) args[i + 1] = batch.get(i).key();
redis.execute(conn -> conn.execute("BF.MADD", toBytes(args))); // idempotent
long maxId = batch.stream().mapToLong(r -> Long.parseLong(r.key())).max().orElse(0);
redis.opsForValue().set("bf:products:v7:watermark", String.valueOf(maxId));
ack.acknowledge(); // commit after apply
}Interview problem
The problem
Billion-item membership pipeline: 100K updates/sec, 10K queries/sec
Design reliable filter updates for 1 billion records with 100K inserts/sec and 10K membership queries/sec: PostgreSQL → outbox → Kafka → updater → Redis/local filters. Discuss partitioning, ordering, lag, retries, idempotency, replay, rebuild, failure and monitoring.
You're given
- 1B records
- 100K inserts/sec
- 10K queries/sec
- Target FPR 0.1%
When it breaks
Updater consumer crashes for an hour
What you see
New items created in that hour are missing from the filter; without a watermark fallback, their lookups return false 404s.
Fix & prevent
Watermark-based bypass for recent items, lag alerts on the updater group, and automatic restarts.
Explain it without notes
Why is a Bloom filter updater naturally idempotent, and why does that matter with Kafka?
Practice
Compute Redis calls per second for 100K inserts/sec with BF.MADD batches of 500, and the lag if the updater stops for 10 minutes.
Trade-offs
- ↔
Event-driven updates add Kafka and CDC infrastructure but give reliability, replay and multiple consumers.
Done when you can
I can design a reliable, idempotent, monitorable filter update pipeline.