A BigQuery streaming pipeline is producing duplicate rows. How do you diagnose and eliminate them?
answer
- look at the timestamp gap between the copies
- at-least-once is the documented default
- offsets are the server-side idempotency token
- a stable source-generated key makes cleanup possible
basics
~20 sDuplicates usually mean at-least-once writes: the default stream or legacy insertAll retried an append whose acknowledgement was lost. Fix it upstream with application-created committed streams and per-request offsets, or downstream by deduplicating on a business key.
solid answer
~50 sFirst confirm where they come from. Group by your business key and count — if duplicates are exact copies arriving seconds apart, it is write retry; if they differ in a field, the source is emitting twice. The usual cause is at-least-once ingestion: the Storage Write API's `_default` stream takes no offsets, and legacy `tabledata.insertAll` only deduplicates on `insertId` best-effort within a short window, so a lost acknowledgement plus a client retry lands the rows twice. The upstream fix is an application-created committed stream where each `AppendRows` carries an offset — the service ignores an offset already written, so retries are safe. The downstream fix, which you often want anyway, is to keep the raw stream table append-only and expose a deduplicated view using `ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY ingest_ts DESC) = 1`, or to `MERGE` periodically into a curated table.
code
sql · 11 lines-- Merge the last two hours of raw stream into a keyed curated table
MERGE curated.events AS t
USING (
SELECT * EXCEPT(rn) FROM (
SELECT *, ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY ingest_ts DESC) rn
FROM raw.events
WHERE ingest_ts > TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 2 HOUR)
) WHERE rn = 1
) AS s
ON t.event_id = s.event_id
WHEN NOT MATCHED THEN INSERT ROW;go deeper
Know that streamed writes into BigQuery can arrive more than once, and that removing duplicates needs a stable identifier carried on each row.
Explain why a retried append duplicates on the default stream, and write the ROW_NUMBER-over-key deduplication query correctly, including which column orders the tie-break.
Diagnose from the data — timestamp gaps and payload differences point to retry, replay or a genuine double-emit — then argue the upstream-versus-downstream fix on cost and operational grounds.
Decide where the organisation buys exactly-once and where it buys idempotent modelling instead, and make source-generated event ids a platform requirement so every downstream cleanup is possible at all.
## Establish the shape of the duplication Before changing anything, characterise it: ```sql SELECT event_id, COUNT(*) AS n, MIN(ingest_ts), MAX(ingest_ts) FROM raw.events WHERE ingest_ts > TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 1 DAY) GROUP BY event_id HAVING n > 1 ORDER BY n DESC; ``` Three signatures, three different causes: - **Byte-identical rows, timestamps seconds apart** — client retried an append it could not confirm. This is at-least-once ingestion behaving as documented. - **Identical rows, timestamps hours or days apart** — a batch or backfill was replayed, or two pipelines write the same source. - **Same key, different payload** — the producer genuinely emitted the event twice, or it is a legitimate update and your "duplicate" is actually a new version. That is a modelling question, not an ingestion bug. ## Why streaming duplicates in the first place Both streaming paths default to at-least-once: - The Storage Write API's **default stream** (`_default`) accepts no offsets. It is shared by many writers, which is what makes it fast and simple, and is precisely why it cannot tell a retry from a new write. - The legacy **`tabledata.insertAll`** method offers `insertId`-based deduplication that Google documents as **best effort** and only within a limited window. Designing for correctness on top of it is a mistake; it is also the path being superseded by the Storage Write API. Any network path can lose an acknowledgement. A client that does not retry loses data; a client that retries without an idempotency token duplicates. There is no third option unless the server can recognise the retry — which is what offsets do. ## The upstream fix: committed streams with offsets Create a stream per writer with `CreateWriteStream(type=COMMITTED)` and send each `AppendRows` request with the offset at which its rows begin. BigQuery tracks the stream's current offset; an append whose offset has already been written is rejected rather than applied, so the client can retry freely. Handling that rejection as success is the whole exactly-once protocol. What it costs you: one stream per writer with a lifecycle (create, append, finalize), offset state that must survive a client crash, and lower throughput per stream than the shared default stream. Exactly-once is a design commitment, not a flag. If your writes are batch-shaped, **pending** streams give a different form of the same protection: nothing is visible until `BatchCommitWriteStreams`, so a failed run publishes nothing and the rerun is clean. ## The downstream fix: deduplicate on read or on merge Even with exactly-once writes, producers replay and backfills overlap, so mature pipelines keep a deduplication layer: ```sql CREATE OR REPLACE VIEW curated.events AS SELECT * EXCEPT(rn) FROM ( SELECT *, ROW_NUMBER() OVER ( PARTITION BY event_id ORDER BY ingest_ts DESC) AS rn FROM raw.events ) WHERE rn = 1; ``` The view is always correct but re-scans and re-sorts on every read, which is expensive on a large table. The usual production shape is therefore a periodic `MERGE` from the raw append-only table into a curated table keyed by `event_id`, with the view kept only for the most recent, not-yet-merged window. Restrict both by partition — the raw table should be partitioned on ingest time so the dedup window scans one partition rather than history. ## Prerequisites you must insist on Deduplication needs a **stable business key** produced at the source — an event id, or a deterministic hash of the immutable fields. A key generated at ingestion time is useless: the retry gets a different one. This is frequently the real finding of the investigation, and pushing the key upstream is the durable fix. ## Choosing between the two fixes Exactly-once writes remove duplicates at the cost of pipeline complexity and per-stream throughput. Downstream deduplication tolerates duplicates from every source — retries, replays, double-running jobs — at the cost of query-time or merge-time compute. Most teams end up with both: the default stream for throughput, and a merge into a keyed curated table because they need protection from replays regardless of the write path.
- Why is insertId deduplication not a correctness guarantee?Google documents it as best effort within a limited time window. It relies on BigQuery remembering recently seen ids, so a retry delayed beyond that window, or a failover, can let the duplicate through. It is a convenience on the legacy insertAll path, not the basis for a correctness argument; the Storage Write API's offsets are.
- Deduplicating with a ROW_NUMBER view over a 40 TB table got expensive. What next?Stop paying for it on every read. Keep the raw table partitioned on ingest time and MERGE the recent partitions into a curated table keyed by the business key on a schedule; the view then covers only the unmerged tail. That bounds the sort to a small window instead of re-scanning history for every consumer.
- The duplicates have the same key but different payloads. Is that still an ingestion bug?Usually not. That signature means the producer emitted the event twice with different content, or you are looking at legitimate updates to the same entity. Decide the semantics first: if later versions supersede earlier ones, keep the latest by ingest time; if both are real events, the key is not unique and the model needs a distinct event identifier.
saying these in an interview costs you the question
- Expecting the default stream to deduplicate on its own
- Trusting insertId as an exactly-once guarantee
- Generating the deduplication key at ingest time rather than at the source
- Running a ROW_NUMBER dedup view over full table history on every read
- Concluding a bug in the source when the write path is simply at-least-once