Command Palette

Search for a command to run...

Hectal
PHASE 8Advanced ~9 min· topic 1 of 3

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.

0/3 · 0%

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

  1. 01

    Pipeline: PostgreSQL (business write + outbox row in one transaction) → Debezium CDC → Kafka topic products.created keyed by ID → Bloom updater group → BF.MADD into Redis (shared) or each instance's local filter (every instance in its own group, or a broadcast pattern).

  2. 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.

  3. 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.

  4. 04

    Replay: rebuilding means starting a new consumer group (or resetting offsets) from the snapshot position into a new filter version.

  5. 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

pipeline.mermaiddiagram
Rendering diagram…
BloomUpdater.javajava
@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

01

Why is a Bloom filter updater naturally idempotent, and why does that matter with Kafka?

Practice

01

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.