intermediate 3 min answer Multiple choice

A team must join an order stream keyed by order id to a payment stream that carries order id but is keyed by payment id, with 12 partitions on one topic and 48 on the other, matching within seven days and with a standing requirement to reprocess the last 30 days after a logic fix. Which runtime fits and what is the deciding property?

flinkkafka streamscopartitioningstatereprocessing
Pick one
Show the full answer Hide the answer

The deciding property

Not throughput, and not the partition mismatch. It is seven days of keyed event-time state plus a requirement to reprocess 30 days through the same code. Those two together decide it: the runtime must hold hundreds of gigabytes of state off-heap with incremental checkpoints, and it must be able to run the same job over historical input with event-time semantics rather than wall-clock ones.

The partition mismatch is real but it is a solved problem in every option: 12 against 48 means the payment stream must be re-keyed on order id and redistributed before the join.

Why this choice

A keyed interval join with RocksDB state keeps the seven-day buffers on local disk with incremental checkpoints, so checkpoint duration tracks the change in state rather than its size. Re-keying is an ordinary shuffle rather than a new topic. Reprocessing runs the same job from a historical offset with watermarks derived from the data, which is what makes the 30-day requirement routine instead of a project.

Why the other options fail

  • A consumer loop with state in Redis. It works until it has to be correct. You now own windowing, out-of-order handling, expiry of seven-day entries, and the atomicity of offset commit against state write — and a crash between those two writes corrupts the join silently. This option is right when the "join" is a lookup against small reference data and the volume is modest, which is a different problem.
  • Kafka Streams with a repartition topic. Genuinely viable, and the right answer if the state were hours rather than days and the team already runs Kafka Streams. What it costs here: the repartition topic is a full second copy of the payment stream in the broker, and window stores holding seven days must be restored from changelog topics on every rebalance, so deploy and recovery time scale with state size. At seven days that is tens of minutes of unavailability per deploy unless standby replicas are configured and kept warm.
  • The payment stream as a GlobalKTable. A GlobalKTable is fully replicated to every instance and holds the latest value per key, so it is for bounded reference data. A payment stream is neither bounded nor latest-value-per-key, and every instance would hold a copy of all of it.

What would flip the decision

If this changes Choose Because
The match window drops to 30 minutes Kafka Streams with repartition State is small; restore is seconds and you avoid a second runtime
Payments arrive within seconds of orders Lookup join plus nightly reconciliation No join state at all; the rare late payment is reconciled in batch
Reprocessing is never required Either framework The replay requirement is what most favours Flink here
The team runs no stream processor today Kafka Streams It is a library in your service rather than a cluster to operate

When not to build the join at all

If 99% of payments land within a minute and the rest are exceptions, do not build a seven-day join at all. Join the fast path with a short window, route unmatched orders to a table, and reconcile in a daily batch job. That design has no long-lived state, no restore time, and its failure mode is a delayed row in a report rather than a corrupted join.