Design a Product Analytics System, stage 7 of 9: decide
Write the batch inserter
Write the loop each ingestion worker runs for one stream partition. Inserts can fail, and workers can crash at any point. The database deduplicates an insert whose dedupToken it has seen recently.
System so far· 8 parts
Select a component to see what it is responsible for and which state it owns.
- 1Customers' apps → Capture API: Batches of events
- 2Customers' dashboards → Query service: Chart request
- 3Query service → Postgres: Teams and saved charts
- 4Query service → Events (column store): Aggregate by column
- 5Capture API → Event stream: Append, then acknowledge
- 6Ingestion workers → Event stream: Read a partition
- 7Ingestion workers → Postgres: Who is this ID?
- 8Ingestion workers → Events (column store): Batch insert
What you need to know
0 of 1 checks done
A stream consumer has two steps per batch: insert the batch, then commit its position ("I've processed up to offset 5,000"). The order decides what a crash does:
- Commit then insert: a crash in between skips the batch. Lost.
- Insert then commit: a crash in between re-reads the batch. Duplicated, unless the database can recognise it.
Check
To let the database recognise a retried batch, what should its deduplication token be?"Exactly once" in practice is at least once plus deduplication. See Delivery guarantees. And batches must be large (by size or time), because a column store turns every insert into a part it later merges.