data pipelines
Design
Streaming ingestion
Transactional outbox delivers to CDC
- CDC
- Idempotent sink
- Checkpoint
- Dedup in window
- Apply watermark (allow late events before closing a window)
- DLQ - quarantine records that can't be parsed
- Replay and backfill wrong output from retained input
Batch ingestion
- Handle late arriving updates (watermarks, reprocessing windows, correctness guarantees)
- Handle duplicates and retries - add idempotency
- Handle schema changes
Add backfills
Handle emrge cost
Define what correct means - latest state per entry? full history? both?
- Choose ingestion model
- pure CDC
- snapshot+CDC (better for bootstrapping and recovery. Can miss events and still be safe)
- Pick stable primary key
- Create surrogate key if upstream is messy. Required for dedup/merge logic.
- Use a reliable ordering signal
- LSN, commit timestamp, monotically increasing version id etc (updated_at can be unreliable)
- Implement idempotency
- Decide how to handle deletes
- Use delete flags, and define how downstream tables will interpret them
- Preserve raw events before transform (bronze)
- Use watermarks, but also keep a safety lookback window (late data is a fact of life)
- Implement reprocessing window for late arrivals (last N hours or days, depending on observed lateness)
Plan schema evolution
- Keep backwards compatibality? Or allow breaking changes (using versioned schemas)
Perf stuff
- Avoid merge at scale (avoid retwriting massive files repeatedly)
- Merge only recent partitions
- Cluster by primary key
- Compact regularly
- Partitions and buckets should match query/update patterns
- Time partition helps loads, primary key helps merges
- Avoid merge at scale (avoid retwriting massive files repeatedly)
Operational stuff
- Make backfills a first class workflow
- Not just a "same job run with an older date". Put in backfill mode with rate limits to protect sources, have correctness checks, and data merges without double counting
- Add observability
- Track i/p and o/p rows, dedupe rate, late event rate, merge duration and rewrite size, null spikes on critical columns
- Make backfills a first class workflow
Handle
- Late data
- idempotency
- merge costs
Change merge heavy workloads to append only workloads. It will avoid file rewrites.
Then, create materialized views on the data for read snapshots.
Questions
How to ingest and process over 1 million events per day efficiently?
Roughly: 1m per day = 10events per sec (trivial scale)
Design around durability, replayability, batching, cheap storage
- Process in batches. Not processing every event individually. Preferably in 5-50mb parquet files (instead of millions of tiny objects)
- Take the original raw event and store it into bronze layer. Preserved for future replayability in case of logic changes, bugs etc.
- Convert raw json to parquet early (convert in the silver layer itself)
- Implement idempotency. The producer may deliver twice, the consumer may read twice (eg - consumer reads and crashes before commiting kafka offset)
- Separate ingestion from transformation
Layers
- Bronze - stores the raw payload, source metadata and , ingestion metadata (
_meta_ingestion_time, _meta_source, _meta_batch_id) - Silver - delta table - parse json, cast data types, standardize timestamps (timezone etc), handle nulls, handle duplicates, validate records, flatten, standardize col names, enrich with dimensions where applicable, apply cdc logic if required