An hourly per-customer total emits early and corrects later, and the sink upserts on customer id: what breaks, and what identity is right?
answer
- one row per what, exactly
- the next interval overwrites the last
- the grouping key is not enough
- identity includes the interval's start
- replace the whole value, never increment
basics
~20 sEach hour's value overwrites the previous hour's, so the sink ends up holding one row per customer showing only the latest hour. The identity must be the grouping key together with the interval's start, so revisions land on their own row.
solid answer
~50 sThe sink is applying two different things on one identity. Within an hour, the corrections for that hour are supposed to replace each other — that part works. Across hours, the next interval's emission also matches `customer id`, so it replaces the previous hour instead of standing beside it, and history quietly disappears. The identity a revisable sink upserts on must be **the grouping key together with the interval's start** (add the interval length or the definition's name when one job writes several groupings, and note that overlapping intervals are distinguished by their starts). Two further conditions: the emission should carry the group's **whole value**, not the difference since the last one, so a replace needs no knowledge of what was already applied; and the consumer must treat a second arrival for the same identity as a revision, never as a duplicate to discard.
go deeper
Recall that a sink applying corrections has to match each arriving value against something. If that something names only the customer, the next hour's value lands on the same row and the previous hour is gone.
Explain the identity in full: the grouping key together with the interval's start, plus the grouping's length or name when one destination receives several. Explain why a corrected emission carries the whole value instead of a delta.
Show the failure with numbers, note that nothing errors while the table quietly becomes wrong, and raise the harder case where a data-defined group's start can move after it has already been published.
Own the identity as a contract between the job and everyone reading the destination. It determines whether history is answerable at all, and the version that destroys history is the one that looks cheapest on a storage review.
## What the sink is actually being asked to do Under the emit-early-and-correct-later contract, one group hands its value downstream several times: a provisional value, zero or more corrections, and a final one. The receiving system therefore has to answer a question the job cannot answer for it — *is this arriving value a new fact, or a better version of a fact I already hold?* It answers that by matching an **identity**. Get the identity wrong and the contract still "works" in the sense that every write succeeds; it is the contents of the table that are wrong, and nothing errors. ## Why the grouping key alone is the wrong identity Take the worked case: hourly totals per customer, upserted on `customer id`. - 09:00–10:00 for customer 42 fires early at 09:00:30 with 3 orders. The sink writes a row for 42. - It fires again at 09:15 with 47, at 09:45 with 118, and finally, once the group is declared finished, with 131. Each of those correctly replaces the last. **This is the part that works.** - At 10:00:30 the *next* hour's group for customer 42 fires early with 2 orders. It matches the same identity. The sink now shows 2. The table has one row per customer, holding whatever hour is currently in progress. Every settled hour is destroyed by the arrival of the next. A downstream reader sees a plausible number that is not the number they asked for, and there is no error anywhere: the job emitted correctly, the sink applied correctly, the identity was wrong. The mirror-image mistake is upserting on the interval's start alone, which collapses every customer's value for that hour into one row. ## What the identity must carry 1. **The grouping key** — which entity the value is about. 2. **The interval's start** — which interval's value this is. With overlapping intervals (a span of fixed length started again every step shorter than that length, so one record sits in several groups at once), the starts differ, which is exactly what keeps those groups apart in the sink. 3. **The interval's length or the grouping's name**, when one job writes several groupings of the same entity into one destination; otherwise a five-minute group and an hourly group with the same start collide. 4. **Nothing else.** In particular the emission's own sequence number or timestamp must *not* be in the identity: if it were, every correction would land as a new row and the contract would degrade to insert-only. A settlement marker ("this emission is final") is useful, but it is a **column**, not part of the identity. ## Replace, do not add A corrected emission should carry the group's whole value, not the difference since the last one. The reason is not transmission size — a difference is smaller. It is that a replace requires no knowledge of what the consumer already applied, while an increment requires the consumer to know exactly which earlier emissions it counted, and that bookkeeping is where consumers go wrong. A whole-value replace also means applying the same emission twice leaves the same row. (What happens to downstream effects when a job restarts and replays is a separate subject with its own owner; the point here is only that the sink's write is a replace.) ## When the group's own identity can move One shape makes this harder. A group with no fixed length that stays open for one key while records keep arriving and closes once a stated quiet span passes with none — commonly a session window — has a start that is a property of the data. If the group fires early under a start of 09:04, and a record then arrives that bridges the quiet gap to an earlier group, two groups may merge into one whose start is 09:01. The sink now holds a row under an identity that no longer corresponds to any group. There are three honest answers, and which is available depends on the runtime: - Emit only once, on close, for that shape, so no identity is ever published that can later move. - Use a contract where the emission carries an explicit retraction of the superseded value as well as the new one — some runtimes offer this and others emit only the new value, so the sink must be built for the one it actually receives. - Make the sink's identity an opaque group identifier the job assigns and keeps stable across a merge, which pushes the problem into the job where the merge is visible. The shapes themselves are a separate subject; what belongs here is the consequence for the consumer. ## What it costs A correct identity means the destination holds one row per entity per interval, not one per entity. That is a bigger table, and it grows with retained history rather than with the number of customers. It is also the only version of the table that can answer "what was the total for customer 42 between nine and ten" after ten o'clock — which is usually the reason the job was written.
- Why should a correction carry the whole value rather than a delta?Because a replace needs no knowledge of history: whatever the row held is overwritten. An increment obliges the consumer to know precisely which earlier emissions it has already added, and any gap or repeat in that bookkeeping silently corrupts the number with no error raised.
- Should the settlement marker be part of the upsert identity?No. Put it in a column. In the identity it would separate the final emission from the provisional ones it is meant to replace, leaving both rows in the table — exactly the accumulation the revisable contract exists to avoid.
- The destination now holds far more rows. Is that a sign the identity is wrong?No, it is the correct shape. One row per entity per interval grows with retained history, while one row per entity does not grow at all because each interval destroys the last. The small table looked cheap because it was discarding the answer.
saying these in an interview costs you the question
- Upserts a corrected value on the grouping key alone
- Sends deltas and expects the sink to accumulate them
- Puts the emission timestamp into the upsert identity
- Drops a second emission for a group as a duplicate
- Assumes a group's start can never change after emitting
- Reads a smaller destination table as evidence of efficiency