Streaming Joins, Backpressure & Operations — Interval, Temporal, Skew & Backfill
Joins and flow control are where streaming jobs fall over in production. A stream-stream join must buffer both sides, so it needs a time bound or its state grows forever. A stream-table join must decide which version of the row applies: the latest value, or the version that was valid at the event's time. Backpressure is how a healthy engine slows a fast source instead of crashing. Skew is why one subtask is at 100 percent while the rest idle. Operating the job means watching lag, planning a backfill, and upgrading state safely.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
Why can't you run an unbounded stream-stream join in production?
Answer
Each side is buffered until a match could arrive. Without a time bound or a TTL, state grows without limit. An interval join lets the watermark prove when a buffered row can never match, and drops it.
L2
What is a lookup join versus a temporal join?
Answer
A lookup reads the latest value at processing time. Simple, and non-deterministic on replay. A temporal join uses the version valid at the event's timestamp from a versioned table. Deterministic, and correct for prices, rates, and tiers.
L3
How does Flink's backpressure work?
Answer
Credit-based flow control. Receivers grant credits equal to free buffers. Senders send only when they hold credit. A slow operator stops granting credit, pressure reaches the source, and the source stops reading so lag accumulates in Kafka.
L4
How do you find the bottleneck in a backpressured job?
Answer
Look at per-operator busy and backpressured time. The bottleneck is the first busy operator downstream of the backpressured chain. Then check skew, external calls, serialization, or RocksDB.
L5
How do you handle a celebrity hot key in a count?
Answer
Two-phase aggregation. Salt the key across N subtasks for partial counts, then merge by the original key. Or pre-aggregate locally before the shuffle. Only combinable aggregates can be salted.
L6
How do you backfill a month of data after fixing a bug?
Answer
Run the fixed job from earlier offsets, or from the lake, into a new output table. Use watermark alignment and a rate limit. Validate, then cut consumers over. Use a savepoint upgrade only when state can be carried forward.
L7
What lag metric do you alert on?
Answer
Time lag, the age of the oldest unprocessed record or the watermark delay versus now, tied to the freshness SLO, plus the trend. Not raw record counts. A million records can be two seconds or two hours.
Failure modes
Regular join with no TTL
Both sides are buffered forever. The disk fills. Prefer an interval or temporal join.
Lookup against the production database
No cache and no async I/O. The join becomes a denial of service against your own database.
Bigger network buffers as the fix
You postpone the out-of-memory kill and make checkpoints larger.
Salting a non-combinable key
Per-key order breaks. Partial counts cannot be merged if the function is not a sum, a count, or a sketch.
Misconceptions
Backpressure is the bug.
Backpressure is the job refusing to OOM. The bug is the slow operator, the hot key, or the external call.
A temporal join is just a lookup with a cache.
A lookup uses whatever is current at processing time. Replay reprices history. Temporal uses the version valid at event time.
Record lag is the SLO.
Time lag is the SLO. Record counts lie across topics.
Interviewer traps
Adding parallelism before naming the hot key.
Show the per-subtask load. Salt or split the celebrity key. More workers do not help one key.
Replaying a month at full speed into the live table.
Write a new table, rate-limit, validate, then cut over. A backfill can take down the sink.
Design scenario
Same prompt for every reader.
Requirements
An interval join with watermark cleanup, a temporal join against a versioned rate table, salted two-phase aggregation, credit-based backpressure, and a backfill into a new table.
Traffic / scale
Impressions and clicks on the order of tens of thousands per second. One key dominates the count.
Latency
Matches within the ten-second bound. Backpressure may grow Kafka time-lag. Alert on that lag, not on buffer size.
Consistency
A replay prices orders at the historical rate, not today's rate. Join output is deterministic.
Availability
A slow sink stalls the source instead of killing the task. The log holds the backlog.
Failure assumptions
- One partition is hot.
- A bug fix requires reprocessing thirty days.
Constraints
- Do not lookup the FX table at processing time.
- Do not salt a per-key ordered processor.
Prompt
Clicks must match impressions within ten seconds, orders must be priced with the FX rate that was valid at order time, and one celebrity key is 80 percent of a count. The sink database cannot absorb a full-speed month replay.
API
What does the backfill consumer cut over to, and how do you know it caught up?
Data
Which timestamp bounds the impression buffer, and which version of the rate applies?
Architecture
Where does credit stop, and which operator is busy just downstream?
Which version of the row
Prefer
The version valid at event time
A temporal join prices the order with the rate that was current when the order happened. Replay stays deterministic.
- The table is a changelog with valid-from timestamps.
- The watermark on the table side bounds how much history you keep.
- Interval joins on two streams use the same idea: a bound, then delete.
Alternative
Whatever is latest when we process it
A lookup join is easy and wrong on replay. Last week's orders reprice at today's rate.
- Processing time is not a version.
- A cache makes it faster and still non-deterministic.
- An unbounded stream-stream join never gets to delete.
Buffer, then prove you can delete
The watermark is the minimum of both inputs. That minimum is what makes the buffer finite.
- 1
Both sides arrive
Impressions and clicks, or orders and a rate changelog. Key them the same way. - 2
Keep only a bounded buffer
An interval join stores a side for the width of the interval. A temporal join stores versions, not one latest row. - 3
Emit the match you can defend
A click inside the bound, or the rate whose valid-from is the latest one at or before the order. - 4
Drop what the watermark closed
Once the watermark passes the impression time plus the bound, that impression can never match.
Join types
| Join | Semantics | State kept | Cleanup | Output | Use for |
|---|---|---|---|---|---|
| Regular stream-stream | Any match, any time | Both sides forever | Only a TTL | Updates and retractions | Small, bounded key spaces only |
| Interval | The other side's timestamp is inside a lower and upper bound | Each side for the interval width | Watermark passes the bound | Append | Impression and click, order and payment within N minutes |
| Window | Same key in the same window | One window per key | Window end plus lateness | One append per window | Per-window correlations |
| Lookup | Probe an external table at processing time | None, or a cache | Not applicable | Latest value. Non-deterministic on replay | Enrichment where latest is acceptable |
| Temporal | The table version valid at the event's time | Version history per key | Watermark on the table side | Deterministic | FX rates, prices, the tier at purchase time |
| Stream-table in Kafka Streams | Probe a local materialized table | The table in RocksDB | Compaction | Latest value at processing time | Kafka-native enrichment |
Flow
- 1
1a. Impression stream
- next2. Interval join
- 2
2. Interval join
- next3. Temporal join
- 3
1b. Click stream
- next2. Interval join
- 4
3. Temporal join
- next4. Price at event time
- 5
1c. Rate versions
- next3. Temporal join
- 6
4. Price at event time
Lesson map
Streaming Joins, Backpressure & Operations — Interval, Temporal, Skew & Backfill
Joins and flow control are where streaming jobs fall over in production. A stream-stream join must buffer both sides, so it needs a time bound or its state grows forever. A stream-table join must decide which version of the row applies: the latest value, or the version that was valid at the event's time. Backpressure is how a healthy engine slows a fast source instead of crashing. Skew is why one subtask is at 100 percent while the rest idle. Operating the job means watching lag, planning a backfill, and upgrading state safely.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB both["1. Impressions and clicks by ad"] emit["2. Interval join, then drop stale"] temporal["3. Temporal join at event time"] price["4. Price with the rate then valid"] both -->|1. Impressions and clicks by ad| emit emit -->|2. Interval join, then drop stale| temporal temporal -->|3. Temporal join at event time| price
Press Run. Snippets must be self-contained — no network, files, or native modules.
interval join matches: [('ad1', 1, 5), ('ad3', 8, 16), ('ad4', 30, 33)]
unmatched clicks: ['ad2'] (ad2 clicked at 14 > 3+10) | peak buffered impressions: 3
order@ 50: temporal=110.00 USD vs latest-lookup=108.00 USD
order@150: temporal=112.00 USD vs latest-lookup=108.00 USD
order@250: temporal=108.00 USD vs latest-lookup=108.00 USD
order@120: temporal=112.00 USD vs latest-lookup=108.00 USDad2's click at 14 is later than impression 3 plus 10, so it does not match, and the buffer peaks at 3 impressions. The late order at t=120 still gets the rate that was valid at t=120, 1.12. A lookup join would use whatever is current when the record is processed, so replaying last week would reprice everything at today's rate. That non-determinism is the interview point.
Backpressure
Flink's network stack, since 1.5, sends data only when the receiver grants credits equal to free buffers. A slow operator stops granting credit. Upstream buffers fill. The source stops polling Kafka. The backlog stays in the log as lag, not in the heap.
Flow
- 1
1. Slow sink stops credit
- next2. Upstream buffers fill
- 2
2. Upstream buffers fill
- next3. Source pauses the poll
- 3
3. Source pauses the poll
- next4. Lag stays in the log
- 4
4. Lag stays in the log
- next1. Slow sink stops credit
Press Run. Snippets must be self-contained — no network, files, or native modules.
no flow control : {"sent":50,"done":20,"queueAtEnd":30} <- queue grows 3/tick until OOM
credit-based : {"sent":22,"done":20,"queueAtEnd":2} <- bounded; source throttled to sink rate
plain keyBy load : [840,40,60,60] max/avg=3.36
salted 2-phase : [240,260,240,260] max/avg=1.04 celebrity total=800| Approach | What happens | Verdict |
|---|---|---|
| No flow control | Memory grows until the process dies. A restart loses in-flight data | Never |
| Drop or sample | Bounded and lossy | Only telemetry where loss is acceptable |
| Credit-based backpressure | The source slows to the sink. Lag accumulates in Kafka | The default. Pair it with lag alerts and autoscaling |
| Add parallelism | Raises the throughput ceiling | Fixes sustained overload. Needs key groups so rescale is a range move |
Backpressure is a symptom. In the Flink UI, find the first operator that is busy. The bottleneck is usually the operator just after the last backpressured one. Kafka's own lag and rebalance story is consumer lag and backpressure.
Skew
Plain keyBy put 840 of 1000 records on one subtask. The salted pass spread the load to about 1.04 times the average and still totaled 800 for the celebrity key.
| Technique | How | Trade-off |
|---|---|---|
| Two-phase aggregation | keyBy the key plus a salt, partial aggregate, then keyBy the original key and merge | Only combinable aggregates: sum, count, sketches. Extra hop |
| Local pre-aggregation | Combine inside a subtask before the shuffle | A small latency bump |
| Split hot keys | Detect the top keys and route them on a dedicated path | More logic. Detection has to be live |
| Rebalance stateless stages | Rebalance before an expensive stateless map | Only for non-keyed work |
| Broadcast the small side | Broadcast state for a join with a small table | Memory on every instance |
Sketches are the combinable form of a unique count. The merge rules live on approximate aggregations.
Operating the job
Watch consumer lag in records and in seconds. Time lag is what the SLO cares about. Also watch watermark delay versus wall clock, checkpoint duration, size, and failures, backpressured versus busy ratios, late-record and side-output counts, restarts, state size per operator, and sink commit latency.
| Backfill | How | Pros | Cons |
|---|---|---|---|
| Rewind offsets | A new job version from the earliest offset into a new output, then cut over | One codebase | Needs retention or tiered storage. Watermark alignment bounds state during catch-up |
| Savepoint and upgrade | Stop with a savepoint, deploy, resume | Keeps state. No full replay | Only forward fixes. State schema changes must be compatible |
| Bounded batch of the same job | Flink batch mode, or Spark batch over lake history | Efficient for a large history | Two run modes to test |
| Hybrid source | Read files from the lake, then switch to Kafka | Fast bootstrap of state | More moving parts |
A CDC replay that is really a connector backfill, with schema and tombstones, is CDC failure modes. Landing the historical files is object storage.
Pitfalls
From the Python output, say which FX version the order at t=120 uses, and which version a lookup join uses. Then say what a full replay of last week would do to every order if you had chosen lookup.
Interview Q&A
Why can't you run an unbounded stream-stream join in production?
Answer
Each side must be buffered until a match could arrive. Without a time bound or a TTL, state grows without limit. An interval join lets the watermark prove when a buffered row can never match, and drops it.
Lookup join versus temporal join?
Answer
A lookup reads the latest value at processing time. Simple, and non-deterministic on replay. A temporal join uses the version valid at the event's timestamp from a versioned table or changelog. Deterministic, and correct for prices, rates, and tiers.
How does Flink's backpressure work?
Answer
Credit-based flow control. Receivers grant credits equal to free buffers. Senders send only when they hold credit. A slow operator stops granting credit. Pressure reaches the source, which stops reading, and lag accumulates in Kafka.
How do you find the bottleneck in a backpressured job?
Answer
Look at per-operator busy and backpressured time. The bottleneck is the first busy operator downstream of the backpressured chain. Then check skew, external calls, serialization, or RocksDB access.
How do you handle a celebrity hot key in a count?
Answer
Two-phase aggregation. Salt the key across N subtasks for partial counts, then merge per original key. Or pre-aggregate locally before the shuffle.
How do you backfill a month after fixing a bug?
Answer
Run the fixed job from earlier offsets, or from the lake, into a new output table, with watermark alignment and a rate limit. Validate against the old output, then cut consumers over. Use a savepoint upgrade only when state can be carried forward.
What lag metric do you alert on?
Answer
Time lag, the age of the oldest unprocessed record or the event-time watermark delay versus now, tied to the freshness SLO, plus the trend. Not raw record counts.