skip to content

A high-throughput consumer group processes millions of events per hour and checks a shared relational 'processed_messages' table before applying each one. What architectural problems does this single shared table cause at that scale, and what alternative designs would you consider?

level: principalimportance: nice to knowfreq 35%

answer

  1. single shared table = throughput bottleneck
  2. unbounded growth without TTL/pruning
  3. Redis SETNX + TTL as fast dedup store
  4. shard dedup store by same partition key as messages
  5. Kafka exactly-once as narrower alternative

basics

~20 s

One shared 'have we seen this before' table gets hammered by every single message across a huge system, becomes a bottleneck, and keeps growing forever. At large scale, teams split it up, for example one store per partition, or use a fast key-value store with automatic expiry instead of one giant table everyone reads and writes.

solid answer

~60 s

A single shared relational processed_messages table checked by every consumer instance in a high-throughput group becomes both a throughput bottleneck, every message adds a read plus a write against one table/index, serialized by that table's own locking, and a hot-row/hot-index contention point, and it grows unbounded unless something continuously prunes it. At scale, teams typically move to a low-latency key-value store, for example Redis with TTL-based expiry, or a partitioned store, that can absorb far higher read/write throughput than a general relational table, shard the dedup keyspace along the same partition key the messages are naturally partitioned by so a dedup check never needs to look outside its own partition, and set a TTL on entries equal to the maximum plausible redelivery window rather than keeping them forever. Some systems avoid a separate store altogether by pushing the guarantee down into the storage layer the side effect writes to, idempotent upserts, or Kafka's transactional exactly-once semantics within its own ecosystem, trading a general-purpose dedup service for a narrower, cheaper guarantee scoped to what that particular sink actually needs.

go deeper

for a junior

Not expected to design this; a reasonable answer just notices that one shared table could get slow if there's a lot of traffic.

for a middle

Should identify at least the throughput-bottleneck problem and one basic fix like using a faster key-value store.

for a senior

Should discuss TTL-based expiry, sharding by partition key, and the durability trade-off of a faster store versus a relational table.

for a principal

Should weigh multiple alternative architectures, sharded fast store versus pushing the guarantee into the sink versus narrower platform-native exactly-once semantics, articulate the residual risk in each, and describe defense-in-depth, dedup store plus idempotent sink, as the pragmatic answer at real scale.

## Why one shared table caps the whole consumer group A single shared relational `processed_messages` table, checked by every consumer instance before applying a message, works fine at modest volume but becomes a **structural bottleneck** once throughput reaches millions of events per hour. Every message, regardless of which partition, shard, or logical stream it belongs to, funnels through the same table and typically the same unique index on `message_id`, so the table's write path becomes a shared serialization point: - concurrent inserts from many consumer instances contend for index-page locks, - and the read-before-write check adds a network round trip to a single database for every single message, which caps the whole consumer group's throughput at whatever that one table can sustain, however many parallel consumer instances you add. It's the classic mistake of introducing a distributed system's worth of parallel consumers and then routing all of them through one non-partitioned, non-scaled piece of shared state. ## Growth that never stops The second structural problem is unbounded growth: a table that records every processed message and is never pruned grows linearly with total message volume forever, which at millions-per-hour translates to billions of rows within months. As it grows: - index lookup latency degrades as the index grows past what fits comfortably in memory/cache, - storage costs climb, - and routine operations like backups or schema migrations on that table become progressively more painful, all for a table whose actual useful lifespan per entry is only as long as the maximum plausible redelivery window, typically minutes to low hours, not months or years. ## Alternative one: a low-latency store with native expiry The first alternative design is to move the dedup store to a purpose-built low-latency key-value system with native TTL support, **Redis** being the most common choice, where a set-if-not-exists operation on the message key serves the same claim role the unique-constraint INSERT played before, but at far higher throughput than a general relational database, and where entries can be given a TTL matching the redelivery window, say twenty-four hours, so old entries expire automatically instead of needing an explicit cleanup job. The trade-off is durability and consistency semantics: Redis, without careful configuration, can lose recently written keys on a crash/failover, which reopens a small window where a duplicate could slip through right after an outage — an acceptable trade for many systems, but a real one to weigh against the relational table's stronger durability guarantees. ## Alternative two: shard along the messages' own partition key The second alternative is to shard the dedup keyspace along the same partition key the messages are already naturally partitioned by. For a Kafka-based system that typically means: - a dedup store, or table, per partition, - or a store keyed such that lookups for messages from one partition never need to touch data from another. Because a partition is already the unit of ordering and single-consumer-at-a-time processing in Kafka's model, colocating the dedup check with the partition eliminates cross-partition contention entirely and lets the dedup store scale linearly with the number of partitions, the same way the rest of the pipeline already does. ## Alternative three: push the guarantee into the sink The third, often more elegant alternative for a subset of use cases is to avoid a general-purpose dedup store altogether and push the exactly-once guarantee down into the specific sink the consumer writes to. - If the downstream write is a naturally idempotent upsert, no dedup store is needed at all. - If the consumer is entirely inside the Kafka ecosystem, consuming from and producing back to Kafka, or consuming and writing via Kafka Connect, Kafka's own exactly-once semantics, idempotent producer plus read-process-write transactions, can supply the guarantee natively without any application-level dedup table, at the cost of being locked into that narrower transactional boundary and only covering writes that go back through Kafka. ## What teams actually run at that scale **A concrete real-world scenario:** a large e-commerce platform's order-event pipeline processing on the order of ten million events per hour would typically not survive on a single Postgres `processed_messages` table checked synchronously by every consumer instance. Teams at that scale commonly: 1. move to a sharded Redis cluster keyed by `(partition, message_id)` with a TTL of a few hours, sized to the platform's own worst-case redelivery lag, 2. explicitly accepting that a Redis failover in the rare case could theoretically let a duplicate through, 3. and mitigating that residual risk by making the terminal database write itself an idempotent upsert as a **second line of defense**, rather than relying on the dedup store as the sole guarantee.

  • Why is Redis's weaker durability an acceptable trade-off for a dedup store when it wouldn't be for the actual business data?
    The dedup store's only job is to prevent an extra, rare duplicate application of an operation whose downstream effect is usually also guarded some other way, an idempotent upsert, a reconciliation job, or manual correction, so an occasional missed dedup is a low-probability, low-severity, and often recoverable event. The business data itself, the payment record, the order row, has no such fallback, which is why it stays on a fully durable store even while the dedup layer in front of it trades some durability for throughput.
  • If you shard the dedup store by partition, what happens if a message somehow gets reprocessed in a different partition than the original, for example after a topic repartition?
    The per-partition dedup check would miss it, because the new partition's store has no record of the message ID from the old partition; this is a real limitation of partition-local sharding that teams need to either accept, since repartitioning is rare and usually paired with a broader migration/reconciliation step, or handle by keeping a lookup that spans partitions specifically during and after a repartitioning event.
  • When is Kafka's built-in exactly-once semantics NOT sufficient on its own, even for a Kafka-native pipeline?
    It only guarantees exactly-once for writes that stay within the Kafka transactional boundary, reading from and producing back to Kafka topics, or via a transaction-aware Kafka Connect sink. The moment the consumer's side effect is a call to something outside that boundary, like an external payment API, an email send, or a non-transactional database write, Kafka's exactly-once semantics don't cover that external effect at all, and the consumer still needs its own idempotency mechanism for it.

It's like routing every visitor to a thousand-room convention center through one single sign-in desk with one pen and one paper sheet; no matter how many rooms and staff you add, the whole building's throughput is capped by how fast that one desk can write names down. Scaling means giving each wing its own sign-in desk, or replacing the paper sheet with something built for high-speed lookups and letting old entries auto-expire.

saying these in an interview costs you the question

  • Assumes a single relational dedup table scales to any throughput without change
  • No mention of TTL/expiry or growth as a concern at scale
  • Proposes sharding by something unrelated to the messages' own partition key, ignoring cross-shard lookup cost
  • Treats Kafka exactly-once semantics as a universal exactly-once guarantee covering all side effects
  • Doesn't recognize that moving to a faster, less durable store is a deliberate trade-off, not a free upgrade

context