Technologies that implement this pattern: Apache Kafka · Flink & Stream Processing · PostgreSQL · Elasticsearch · Time Series DBs · Cassandra
Why This Matters#
Every company past a certain size runs more pipelines than services. Search indexes, recommendation features, billing aggregates, fraud signals, dashboards, warehouse tables, ML training sets, cache invalidations — each is data that was written somewhere, copied, transformed and written somewhere else. Pipelines are how one write becomes twenty derived views. They are also where data quietly goes wrong: a dropped event, a double-counted click, a schema change that turns revenue into nulls, a late-arriving partition nobody reprocessed.
Most candidates treat pipelines as a batch-vs-streaming question: "use Kafka and Flink" or "a nightly Spark job." Staff engineers treat them as a correctness-over-time question: what is the replayable source of truth, what does each output promise about completeness and freshness, and how do we recompute when — not if — the logic or the input was wrong? A pipeline that can't be re-run from its source is a pipeline whose bugs are permanent.
The second reframe: every pipeline is a contract between teams who don't talk to each other. The producer changes a field; three consumers break a week later. The consumer falls behind; the producer's Kafka retention deletes the data before it is read. The bill depends on an aggregate that a data engineer "fixed" without telling finance. The technology choices are downstream of who owns the schema, who owns lag, who owns the backfill, and who signs off that the numbers are right.
If you can walk an interviewer from "what's the source of truth and what does each output promise" to "how the pipeline handles duplicates, late data and failure" to "how we replay, backfill and evolve schemas without breaking consumers", you are answering at Staff level.
The 60-Second Version#
- Start from the log. Make the source an immutable, replayable log — Kafka with 7–30 day retention plus an object-storage archive, or CDC from the database WAL. If you can replay the input, you can fix any output.
- Choose latency by consumer need, not fashion. Seconds → streaming (Flink, Kafka Streams). Minutes → micro-batch. Hours → batch over object storage. A streaming job costs ~2–5× the engineering and on-call of the equivalent hourly batch job.
- Assume at-least-once and make sinks idempotent. Exactly-once end-to-end = replayable source + checkpointed state + idempotent or transactional sink. Keys like
(entity_id, window_start)turn duplicates into no-ops. - Late data is normal, not an edge case. Mobile and IoT events commonly arrive 1–5% late, some hours late. Use event time with watermarks (allowed lateness 1–10 min for dashboards) and a correction path (batch re-run) for billing-grade outputs.
- Schemas are the API. Use a schema registry with compatibility rules (backward-compatible by default); additive changes only on shared topics. Most pipeline outages are schema outages.
- Lag is the SLO. Every consumer has an owner and a lag target (
consumer_lag_secondsp99 < 30s for operational views, < 1h for analytics). Retention must exceed worst-case lag plus recovery time — typically ≥ 3× the longest plausible outage.
The Problem#
An e-commerce company wants: search results that reflect a price change within 5 seconds; a fraud model that sees a card's last 10 minutes of activity within 1 second; a revenue dashboard by the hour; monthly seller payouts that are exactly right; and an ML team that wants three years of clickstream to train on. The same order event feeds all five. The fraud path can't wait for an hourly batch; the payout path can't trust a streaming counter that silently dropped late events; the search indexer falls 40 minutes behind during a sale and nobody notices until merchants complain; and a producer renames amount_cents to amount on a Tuesday afternoon. One pipeline design can't serve all of them — and five unrelated pipelines reading five different copies of "the order" will disagree about revenue by Thursday.
Case Studies That Use This Pattern#
- Stream Processing — The flagship: windows, watermarks, state, checkpoints and exactly-once in a real design
- Ad Click Aggregator — Billing-grade counts from a firehose: dedupe, late clicks, and batch reconciliation of the streaming result
- Search Indexing — CDC-fed index building, reindexing with zero downtime, and drift detection
- Metrics & Monitoring — High-volume ingestion, downsampling and rollups as pipelines
- Message Queues — Queues vs logs, consumer groups, retries and dead-letter handling
- Idempotency & Exactly-Once — The sink-side half of exactly-once
- Web Crawler — A batch-and-stream hybrid: frontier, fetch, parse, index
- Leaderboard — Streaming aggregation with periodic correction
The Four Intents#
"Build a data pipeline" covers four intents with different freshness, correctness and failure postures.
| Intent | Constraint | Strategy | Failure Mode | Correctness Bar |
|---|---|---|---|---|
| Operational derived views (search index, caches, read models) | Freshness 1–10s; must converge to source | CDC → log → idempotent upserts; periodic reconciliation | Lag during spikes; silent drift from missed events | Eventually equal to source, drift-monitored |
| Real-time decisions (fraud, alerting, pricing, rate signals) | Sub-second to seconds; value decays fast | Stream processing with keyed state, windows, watermarks | Late or out-of-order events; state loss on failover | Best-effort timely; approximate is OK |
| Financial / billing aggregates (payouts, invoices, ad spend) | Exact, auditable, reproducible | Streaming for preview + batch over the immutable log as the source of record | Double counting; late data dropped; untracked logic changes | Exact, reconciled, versioned |
| Analytics and ML (warehouse, features, training sets) | Hours OK; huge volume; schema history | Batch ELT over object storage (Parquet/Iceberg), partitioned by date | Small files; schema drift; backfill cost | Complete per partition, reproducible |
🎯 Staff Move: "I'll treat seller payouts as the billing intent: the streaming counter is a preview, but the number we pay is computed by a batch job over the immutable event archive, versioned, and reconciled against the ledger. The fraud path is the real-time intent and accepts approximation. They share the same source topic, not the same pipeline."
The Core Tradeoff#
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Nightly/hourly batch (Spark/SQL over object storage) | Simple, cheap, reproducible, easy to backfill | Freshness in hours; big failures discovered late; long reruns | Consumers who need fresher data; on-call for morning failures |
| Micro-batch (1–15 min) | Most of batch's simplicity at minute freshness | Small-file problem; scheduling overhead; still not seconds | Storage/compaction jobs; platform tuning |
| Stream processing (Flink, Kafka Streams) | Seconds of latency; continuous; event-time correctness | State management, checkpoints, upgrades; 24/7 on-call; harder to backfill | Platform + owning team's on-call; engineering time |
| Lambda (batch + stream in parallel) | Fast preview + correct batch result | Two codebases computing the same thing; they drift | Team maintaining both; users confused by two numbers |
| Kappa (stream only, replay to recompute) | One codebase; reprocess by replaying the log | Needs long retention; replay at speed can be expensive; some joins are hard | Storage for retention; compute during replays |
| CDC / outbox | Captures every DB change in order without dual writes | Couples consumers to source schema; snapshots and re-syncs are operationally fiddly | Source DB team (schema changes ripple); platform |
| ELT into warehouse (load raw, transform in SQL) | Raw kept; transformations versioned in SQL; analysts self-serve | Warehouse cost grows with careless queries; governance needed | Finance (warehouse bill); data platform |
Staff Default Position#
One immutable, replayable log per business fact; every output is a derived view with a declared freshness, correctness bar, owner and recompute path.
Producers publish facts (events, CDC) to Kafka with registered, backward-compatible schemas, and the log is archived to object storage for long-term replay. Operational views consume via idempotent upserts keyed by entity and version. Real-time decisions run in a stream processor with checkpointed state and event-time windows. Anything that moves money is computed (or re-computed) in batch from the archive, versioned, and reconciled against an independent source. Start with batch unless a consumer genuinely needs seconds; promote to streaming when that need is real and owned. Every consumer has a lag SLO, an owner, a dead-letter path and a documented backfill procedure.
When to Deviate#
- Small data, single database — If all the data fits in one Postgres and the outputs are reports, a materialized view or a scheduled SQL job beats any pipeline. Don't introduce Kafka for 50 events/sec.
- Cross-system transactions that must be atomic — A pipeline is eventually consistent by nature. If the derived write must commit with the source (e.g., balance and ledger entry), do it in the same transaction, not downstream.
- Exploratory analytics — Ad hoc questions over raw data belong in the warehouse/lakehouse with schema-on-read. Productionize a pipeline only once the question repeats.
- Extreme freshness with simple logic — If the "pipeline" is "update a counter when X happens" and needs < 100ms, do it in the write path or a cache, not a stream job.
The L5 → L6 → L7 Contrast#
| Behavior | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| First move | "Kafka + Flink" or "nightly Spark job" | "What's the source of truth, can it be replayed, and what does each consumer need for freshness and correctness?" | "How many copies of 'order' exist across the org, who owns the canonical one, and what contract do producers sign?" |
| Correctness | "Kafka gives exactly-once" | Replayable source + checkpoints + idempotent sinks; event time + watermarks; batch reconciliation for money | Org-wide data contracts: schema compatibility, ownership, quality SLOs enforced in CI and at the registry |
| Failure | Retry the job | DLQ with owner, lag SLOs, retention ≥ 3× worst outage, backfill runbook | Designs for the org: lineage to find blast radius, standard backfill tooling, quality incidents treated like outages |
| Evolution | Change the schema and redeploy | Backward-compatible changes; versioned topics for breaking ones; dual-run and compare | Deprecation policy and timelines across 30 consumers; producer accountability for downstream breakage |
| Cost | Not discussed | "Hourly batch is a tenth the cost of streaming for this use" | Prices streaming vs batch fleets, retention tiers and warehouse spend; kills pipelines nobody reads |
| Ownership | The data team owns pipelines | Producers own schemas; consumers own their lag and views; platform owns Kafka/Flink | Domain-owned data products with published contracts; central platform for paved tooling only |
Why "Correctness" separates levels
Kafka's transactions give exactly-once within Kafka. The moment the output is a database, a search index or an API call, exactly-once depends on the sink. The L5 answer isn't wrong about the feature; it's wrong about the boundary. The Staff answer makes duplicates harmless at the sink and keeps a replayable source so mistakes are fixable. The Principal answer recognizes that most data incidents are not duplicates at all — they're contract breaks between teams — and puts the contract into tooling.
Why "Failure" separates levels
"Retry the job" assumes the input is still there and the job is idempotent. With 3-day retention and a consumer that was broken over a long weekend, the input is gone. Staff sizes retention against realistic outages and writes the backfill runbook before it's needed. Principal makes lineage a platform capability so the question "which dashboards and payouts consumed the bad data?" takes minutes, not a week.
Why "Ownership" separates levels
Central data teams that own every pipeline become a ticket queue, and they own outcomes they can't control — the producer changes the schema, the data team gets paged. Staff pushes schema ownership to producers. Principal reorganizes around data products: the domain that produces "orders" publishes a versioned, documented, quality-monitored dataset, and the platform provides the paved road.
The Five Fault Lines#
| # | Fault Line | The Tension |
|---|---|---|
| 1 | Freshness vs Correctness | Answer now with what has arrived, or wait for completeness |
| 2 | One Engine vs Two Paths | Stream-only simplicity vs batch's reproducibility for money |
| 3 | Push Semantics into the Pipeline vs into the Sink | Exactly-once machinery vs idempotent writes |
| 4 | Producer Freedom vs Consumer Safety | Ship schema changes fast vs never break downstream |
| 5 | Retention vs Cost | Replay anything vs pay to keep everything |
Fault Line 1: Freshness vs Correctness#
A 1-minute window closes at 12:01:00 wall clock, but 2% of its events are still on phones in tunnels. Close it now and the count is 2% low forever; wait 10 minutes and the dashboard is 10 minutes stale. Watermarks make this explicit: the watermark is the pipeline's assertion "I've probably seen everything before time T", and allowed lateness says how long to keep windows open for stragglers. Who pays: emitting early — anyone acting on undercounted numbers (an alert that never fires, an advertiser billed short); waiting — every consumer of freshness. Staff default: emit early results on the watermark (bounded out-of-orderness 30s–2min), keep windows updatable for 10–60 minutes with retractions/upserts downstream, and send anything later to a side output that the batch correction picks up. Deviate when: the consumer is a decision that can't be undone (fraud block) — act on what you have and accept the approximation explicitly.
The sequence is the argument for two speeds: the stream gives a number in two minutes, updates it as stragglers arrive, and the batch job over the archive produces the final, versioned answer.
Fault Line 2: One Engine vs Two Paths#
Running batch and streaming implementations of the same logic (the classic lambda shape) gives a fast preview and a correct final answer — and two codebases that drift until the dashboard and the invoice disagree by 0.7% and nobody can explain why. Stream-only (kappa) avoids drift by recomputing via replay, but replaying 90 days of a 2 GB/s topic is a serious compute event, and some logic (large historical joins, model training) is far easier in batch. Who pays: two paths — the owning team (double maintenance) and users (two numbers); stream-only — retention cost and the team running large replays. Staff default: one logic definition, two execution modes where the engine allows it (the same SQL or Flink job run in streaming and batch mode over the archive); otherwise, the batch path is the source of record and the stream is explicitly labeled preview. Deviate when: the output is not money and approximate is fine — stream-only, no reconciliation.
Fault Line 3: Exactly-Once in the Pipeline vs in the Sink#
Exactly-once stream processing (Flink checkpoints with two-phase-commit sinks, Kafka transactions) gives strong guarantees inside the framework at the cost of latency (output visible only on checkpoint commit, e.g., every 10–60s) and operational sensitivity (transaction timeouts, aborted transactions on restarts). Idempotent sinks — upsert keyed by (entity, window, version), or insert with a unique event ID — make at-least-once delivery effectively exactly-once with far less machinery. Who pays: pipeline-level EOS — latency and on-call complexity; sink idempotency — schema design (every output needs a natural key) and storage for dedupe keys. Staff default: at-least-once processing + idempotent, keyed sinks; reserve transactional sinks for Kafka-to-Kafka stages where it's native. Deviate when: the sink can't be idempotent (an external API that sends emails or moves money) — then put an idempotency key in the call and keep a dedupe table in front of it.
Fault Line 4: Producer Freedom vs Consumer Safety#
The team that owns the order service wants to rename a field and ship today. Seven consumers — search, fraud, payouts, the warehouse, two ML features, a partner export — parse that field. Who pays: producer freedom — every consumer, usually discovered days later as nulls in a dashboard or a crashed job; consumer safety — producer velocity (compatibility checks, versioned topics, deprecation windows). Staff default: schema registry with backward compatibility enforced on publish (new readers can read old data; consumers upgrade at their own pace); additive changes only; a breaking change means a new topic version (orders.v2) dual-published for a deprecation window (e.g., 90 days) with consumers migrated and tracked. Deviate when: a topic has a single consumer owned by the same team — relax to "coordinated deploy", but register the schema anyway.
Fault Line 5: Retention vs Cost#
Replayability requires the input to still exist. Kafka on broker disks at 3× replication costs roughly 10–20× per GB what object storage does. Who pays: short retention — whoever must fix a bug discovered after the data aged out (often: nobody can, and the output stays wrong); long retention — the platform's storage bill. Staff default: tiered retention — 3–7 days on brokers (covers outages and quick replays), plus an archive of every topic to object storage in Parquet/Avro partitioned by hour, kept 1–3 years or per legal requirement; replays beyond a week read from the archive. Deviate when: data is regulated (PII, deletion obligations) — retention is set by legal, and the archive needs a deletion mechanism (key-level tombstones or crypto-shredding).
Common Interview Mistakes#
| What Candidates Say | What Interviewers Hear | What Staff Engineers Say |
|---|---|---|
| "Kafka guarantees exactly-once" | "Doesn't know where the guarantee ends" | "Within Kafka. My sink is Postgres, so I upsert on (campaign_id, minute) and replays are no-ops." |
| "Use processing time for windows" | "Will miscount whenever anything is delayed" | "Event time with a 1-minute watermark; late events update the window for 30 minutes; anything later goes to the batch correction." |
| "Stream everything for real-time" | "Paying 24/7 streaming cost for hourly consumers" | "Fraud needs seconds, so it streams. The finance dashboard is read once a morning — an hourly batch is a tenth the cost." |
| "Failed messages go to a DLQ" | "DLQ as a black hole" | "DLQ with an owner, an alert on depth > 0, a replay tool, and a 7-day SLA to drain it." |
| "We'll reprocess if there's a bug" | "Hasn't checked retention" | "Brokers keep 7 days; every topic is archived hourly to object storage, so a backfill can start from any point in the last two years." |
| "Just add the new field" | "Hasn't been paged by a schema break" | "Registered schema, backward-compatible check in CI, field optional with a default. Renames are a new version with a 90-day dual-publish." |
Quick Reference#
Staff Sentence Templates#
"The source of truth is [the orders WAL via CDC / the clicks topic], retained [N days] on brokers and archived to object storage for [M years]. Every output here is a derived view I can recompute from that."
"[Consumer] needs [seconds / minutes / hours] and [approximate / exact] results, so it gets [streaming / micro-batch / batch]. [Consumer B] moves money, so the stream is a preview and the batch over the archive is the number we pay."
"Processing is at-least-once. The sink is idempotent on [key], so replays and retries are no-ops. That's how I get exactly-once results without depending on [framework] transactions end to end."
"Lag is the SLO: [consumer] targets p99 under [X]; retention is [Y], which is [Y/X]× our worst outage. [Team] owns the lag alert and the backfill runbook."
Implementation Deep Dive#
1. Transactional Outbox → CDC → Kafka (PostgreSQL + Debezium)#
The source must be a faithful, ordered log of facts. Dual writes from application code ("write DB, then publish to Kafka") lose events on crash between the two. The outbox puts the event in the same transaction as the state change.
BEGIN;
UPDATE orders SET status = 'paid', version = version + 1 WHERE order_id = $1;
INSERT INTO outbox (event_id, aggregate_type, aggregate_id, event_type, payload, created_at)
VALUES (gen_random_uuid(), 'order', $1, 'OrderPaid', $json, now());
COMMIT;
# Debezium reads the WAL (logical replication slot), routes outbox rows to
# topic 'orders.events', key = aggregate_id -> per-order ordering preserved.
# Outbox rows deleted by a janitor after 24h (the WAL already carried them).
# Watch the replication slot: an unconsumed slot retains WAL on the primary
alert when pg_replication_slots.confirmed_flush_lsn lags > 20 GB # disk risk
Why it matters: this is the difference between "usually consistent" and "provably complete." The cost is coupling to the source database's replication: a stalled CDC connector causes WAL to accumulate on the primary — a pipeline failure that can fill a production database's disk. That alert belongs to the DB platform, not the data team.
2. Event-Time Windowed Aggregation with Idempotent Sink (Flink)#
events = kafka_source("ad.clicks", start="committed offsets")
.assign_timestamps(e -> e.event_time)
.with_watermarks(bounded_out_of_orderness = 60s)
deduped = events
.key_by(e -> e.click_id)
.filter(first_seen_within(ttl = 24h)) # keyed state dedupe
counts = deduped
.key_by(e -> e.campaign_id)
.window(tumbling(1 minute))
.allowed_lateness(30 minutes) # window re-fires on late events
.side_output_late("clicks.too_late") # beyond 30 min -> batch correction
.aggregate(count, sum(cost_micros))
counts.sink(postgres):
INSERT INTO campaign_minute (campaign_id, minute, clicks, cost_micros, updated_at)
VALUES ($1, $2, $3, $4, now())
ON CONFLICT (campaign_id, minute)
DO UPDATE SET clicks = EXCLUDED.clicks, cost_micros = EXCLUDED.cost_micros,
updated_at = now(); # re-fired window overwrites, not adds
checkpointing: every 30s, exactly-once state, RocksDB backend, incremental
Numbers: a 60s bounded out-of-orderness covers ~98–99% of mobile events in typical apps; the 30-minute allowed lateness picks up most of the rest; the side output is usually < 0.5%. Note the sink writes absolute values (SET clicks = EXCLUDED.clicks) rather than increments — overwriting is idempotent, adding is not. See Flink & Stream Processing for checkpoint and state-backend details.
🎯 Staff Insight: The single most common exactly-once bug isn't in the framework. It's a sink that does
clicks = clicks + delta, which double-counts every replay after a restore. Write absolute window values keyed by window, and replays become harmless.
3. Batch Recompute as Source of Record — Partitioned Archive#
# Archive: Kafka Connect S3 sink, Parquet, partitioned by event hour
s3://events/ad.clicks/dt=2026-09-30/hr=14/part-0001.parquet # ~256 MB files
# Nightly (and on-demand) job: idempotent overwrite of a date partition
job billing_daily(date D, logic_version V):
clicks = read("s3://events/ad.clicks/dt=" + D) # includes very late arrivals
.union(read_late_arrivals(D, up_to = D + 3 days))
.dedupe_by(click_id)
.filter(valid_by_fraud_rules(V))
result = clicks.group_by(campaign_id).agg(count, sum(cost_micros))
write_overwrite("warehouse.billing_daily", partition = D,
columns = result + {logic_version: V, computed_at: now()})
reconcile(D):
diff = compare(warehouse.billing_daily[D], sum_by_day(campaign_minute[D]))
emit metric("billing.stream_vs_batch_diff_pct", diff) # alert if > 0.5%
Why it matters: write_overwrite per partition makes the job re-runnable — a backfill is "run it again for these dates with logic version V+1." The logic_version column makes every number traceable to the code that produced it, which is what finance and auditors will ask for. The stream-vs-batch diff is your early warning that the streaming path is drifting.
4. Dead-Letter Queue with Ownership and Replay#
on message m:
try:
process(m)
except RetryableError as e: # timeouts, 5xx from sink
retry with backoff 100ms -> 10s, max 5 attempts, then treat as non-retryable
except NonRetryableError as e: # schema violation, bad data
kafka.produce("orders.events.dlq", key=m.key, value=m.value,
headers={error: e.code, consumer: "search-indexer",
source_offset: m.offset, first_failed_at: now()})
metrics.incr("dlq.produced", tags=["consumer:search-indexer", "code:" + e.code])
commit(m.offset) # never block the partition on one poison message
# Replay tool: after a fix is deployed
dlq-replay --topic orders.events.dlq --consumer search-indexer --since 2026-09-30T10:00Z
Numbers: alert on dlq.depth > 0 for money-adjacent consumers and on dlq.produced_rate > 0.1% of throughput for others. Without the DLQ, one malformed message stalls a partition indefinitely (head-of-line blocking) and lag climbs for every event behind it; with a DLQ and no owner, it becomes a silent data-loss bucket.
Technique Comparison
| Technique | Freshness | Correctness | Reprocessing | Operational Burden |
|---|---|---|---|---|
| Batch over archive | Hours | High, reproducible | Re-run partition | Low |
| Micro-batch | 1–15 min | High | Re-run interval | Medium (small files) |
| Stream + idempotent sink | Seconds | High if keyed correctly | Replay from offset | Medium–High |
| Stream + transactional sink | Seconds + checkpoint interval | Highest inside framework | Replay from checkpoint | High |
| Lambda (both) | Seconds + hours | Batch wins | Batch re-run | Very high (two codebases) |
| CDC consumer | Sub-second–seconds | Complete if slot healthy | Snapshot + stream | Medium (slot management) |
Architecture Diagram#
How to narrate it: everything flows from two kinds of source — database changes via CDC and events via a validating collector — into one log with registered schemas. Every consumer on the right is a derived view with its own lag SLO and owner. The archive turns the 7-day log into a multi-year one, which is what makes batch correction, ML training and backfills possible. The reconciler and lineage are the control plane: they tell you when the derived views disagree and which ones are affected when something goes wrong.
The schema lifecycle is part of the architecture. Breaking changes don't get "coordinated" in Slack; they become a new versioned topic with a dual-publish window and a tracked migration.
Failure Scenarios#
1. Poison Message — Search Index 6 Hours Stale During a Sale#
A product-catalog CDC topic receives a row with a malformed JSON attribute (written by a legacy admin tool). The search indexer's deserializer throws, the consumer retries the same message forever, and the partition stops advancing.
t=0 Malformed row committed. Partition 7 of catalog.changes affected.
t=+1min Indexer crashes on deserialize, restarts, re-reads same offset. Loop.
t=+5min Lag on partition 7 grows; other 31 partitions healthy. Aggregate lag alert
uses average across partitions: below threshold. No page.
t=+2h Merchants report price changes not reflected in search for 1/32 of products.
t=+5h Sale begins. Out-of-stock items appear in search; conversion drops.
t=+6h On-call finds stuck partition; manually skips offset. Index catches up in 20 min.
Detection: consumer_lag_seconds max per partition (not average) > 5 min; consumer.restarts > 3 in 10 min; consumer.offset_unchanged_seconds.
Blast radius: 1/32 of the catalog stale in search for 6 hours, during a sale.
Mitigation: skip-and-DLQ the offset; replay after fix.
Prevention: non-retryable errors go to DLQ automatically with an owner and alert; schema validation at the source (the admin tool) through the registry; per-partition lag alerts.
Owner: search team owns the indexer and its DLQ; catalog team owns source data quality.
🎯 Staff Insight: The alert existed — it averaged lag across partitions and hid one stuck partition behind 31 healthy ones. Pipeline lag must be alerted on the worst partition, because that's what a customer sees.
2. Non-Idempotent Sink — Advertiser Billed 2× for 40 Minutes#
A spend-aggregation job writes per-minute deltas with UPDATE spend SET total = total + $delta. During a broker upgrade the job restores from a checkpoint taken 40 minutes earlier (the latest checkpoint had failed silently due to a misconfigured state backend).
t=0 Checkpoints start failing (state backend permission error). Job keeps running.
t=+40min Broker rolling upgrade triggers task failure; job restores from last good
checkpoint, 40 minutes old, and replays 40 minutes of events.
t=+41min Replayed deltas added again: 40 minutes of spend double-counted.
t=+45min Budget-pacing service sees spend at 2x; pauses 1,800 campaigns early.
t=+3h Advertiser support escalations. Finance notices invoice preview anomaly.
t=+1d Batch recompute from archive corrects totals; credits issued.
Detection: checkpoint.failed_count > 0 / checkpoint.age_seconds > 3× interval (the missing alert); billing.stream_vs_batch_diff_pct; sudden spend.rate step after restart.
Blast radius: every advertiser active in the window — revenue, trust, and pacing decisions made on bad numbers.
Mitigation: recompute affected partitions from the archive; correct pacing state.
Prevention: sink writes absolute window values keyed by (campaign, minute); alert on checkpoint age; billing numbers come from batch over the archive, never directly from the stream.
Owner: ads data team owns the job and sink; streaming platform owns checkpoint health alerts.
3. Retention Cliff — Bug Found After the Data Expired#
A feature-engineering job for a recommendation model had a timezone bug that shifted all "hour of day" features by 5 hours for users in one region. It's discovered 12 days later. Kafka retention is 7 days and the topic was never archived.
t=0 Bug deployed. Features written to feature store with wrong hour for region X.
t=+12d Model quality regression traced to the feature.
t=+12d Fix deployed. Backfill attempted: Kafka has only the last 7 days.
t=+12d 5 days of raw input are gone. Features for that period cannot be recomputed.
t=+14d Team retrains excluding the bad window; 5 days of training data lost.
Detection: feature distribution monitoring (feature.hour_of_day.distribution_shift per region) would have caught it in hours.
Blast radius: model quality for one region; permanent loss of 5 days of training signal.
Mitigation: exclude bad window; retrain.
Prevention: every topic archived to object storage by default (opt-out requires approval); retention ≥ the time-to-detect for quality bugs, which is often weeks, not days; data quality checks on distributions, not just null counts.
Owner: ML platform owns feature quality checks; data platform owns archive-by-default.
Operational Reality Matrix#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Consumer lag | consumer_lag_seconds max per partition > SLO | Staleness of that view | Scale consumers, prioritize, shed | Consuming team |
| Poison message | offset_unchanged_seconds, restart loops | One partition stalled | Skip to DLQ, replay after fix | Consuming team |
| Schema break | deserialize.errors, registry rejections | All consumers of topic | Roll back producer; versioned topic | Producing team |
| Duplicate counting | stream_vs_batch_diff_pct, step after restart | Aggregates, billing | Recompute from archive | Pipeline owner |
| Late data dropped | late_events.side_output_rate > 1% | Undercounted windows | Batch correction | Pipeline owner |
| CDC slot stalled | replication_slot.retained_wal_bytes | Source DB disk | Restart connector; drop slot as last resort | DB platform + data platform |
| Retention exceeded | Lag approaching retention; archive gaps | Permanent data loss | Emergency retention extension | Data platform |
| Small files | Files per partition > 1K, avg < 32 MB | Query cost and latency | Compaction job | Data platform |
The Principal Lens#
Why L7 Sees This Problem Differently#
A Staff engineer builds a correct pipeline. A Principal engineer counts 400 pipelines, 6 definitions of "active user", 3 copies of the order stream that disagree about revenue by 1.2%, a warehouse bill growing 60% a year, and a central data team that gets paged for schema changes made by teams they've never met. At org scale, pipelines stop being an engineering problem and become a data ownership and contract problem. The questions become: who publishes the canonical version of each business fact; what promises does that dataset make (schema stability, freshness, quality, completeness); how do 30 consumers find out about a change before it breaks them; and which of the 400 pipelines could be deleted tomorrow without anyone noticing.
The Org-Level Fault Line#
Central data team owns pipelines vs domain teams publish data as a product.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Central data team builds and owns every pipeline | Consistent tooling; deep data expertise in one place | Ticket queue; owns breakage it can't prevent; domain knowledge lost in translation | Central team's on-call and morale; consumers waiting weeks |
| Every team builds its own pipelines on its own stack | Speed and autonomy | Duplicate copies, conflicting metrics, no lineage, 5 streaming stacks | Finance (duplicate spend), executives (conflicting numbers) |
| Domain-owned datasets with published contracts on a central paved platform | Producers accountable for schema and quality; platform provides Kafka, registry, archive, lineage, orchestration | Requires producers to accept on-call for data quality; platform must be genuinely easy | Domain teams (some new responsibility); platform headcount (6–10) |
The Principal default is the third row: domains publish canonical datasets with schemas, freshness and quality SLOs; the platform owns the paved road (log, registry, archive, stream/batch runtimes, lineage, backfill tooling); consumers own their derived views and lag.
Cost Model#
Assumptions: Kafka ~$0.10–0.15/GB-month on brokers at 3× replication plus broker compute; object storage ~$0.02/GB-month; managed stream processing ~$0.10–0.25 per vCPU-hour; warehouse compute is usage-based and typically the largest variable line; engineer ~$25K/month.
| Scale | Ingest | Infra | People / On-call | Rough Monthly Total |
|---|---|---|---|---|
| Startup | 1K events/s (~2 GB/day) | Managed queue + scheduled SQL in the main DB or a small warehouse (~$2K) | 0.5 FTE; no dedicated on-call | ~$2K + ~$12K people |
| Growth | 100K events/s (~1 TB/day) | Kafka 9–12 brokers ( | 4 FTE data platform; shared rotation for streaming | ~$60K + ~$100K people |
| Large | 5M events/s (~50 TB/day) | Multi-cluster Kafka ( | 15–25 engineers across platform; dedicated rotations | ~$800K + ~$500K people |
The Principal observation: at growth and large scale the biggest controllable line is warehouse compute, and the biggest waste is pipelines and tables nobody reads. A lineage-driven review that deletes outputs with zero reads in 90 days commonly removes 20–40% of jobs. The second lever is streaming vs batch: moving consumers that only need hourly freshness off 24/7 streaming jobs typically cuts their compute cost 5–10×.
The 3-Year Evolution Path#
One-Way Doors vs Two-Way Doors#
| Decision | Door Type | Reversibility Cost |
|---|---|---|
| Batch vs stream for a single consumer | Two-way | Rewrite one job; source unchanged |
| Window sizes, watermarks, lateness | Two-way | Config + recompute from archive |
| Not archiving a topic | One-way | Data that aged out is gone forever |
| Event key choice on a shared topic (ordering unit) | One-way-ish | Every consumer's ordering assumptions depend on it; requires a new topic |
| Public metric definitions ("active user", "revenue") in exec reporting | One-way-ish | Changing them breaks year-over-year comparisons and trust |
| Data format in the archive (Parquet + open table format vs proprietary) | One-way-ish | Re-encoding petabytes; lock-in to one engine |
| Streaming engine choice | Two-way (strategic) | If jobs are SQL or well-abstracted, migration is quarters, not years |
The Standard I'd Write#
RFC: Data Pipeline and Dataset Contract Standard (v1)
Scope: Every shared topic and every dataset consumed by a team other than its producer.
MUST:
- Each shared dataset has a named owning team, a registered schema with backward-compatibility enforced at publish, and a declared freshness SLO and quality checks (completeness, null rates, distribution bounds).
- Events are produced via outbox/CDC or a validating SDK — never by dual writes in request handlers.
- Every topic is archived to object storage in an open format; broker retention ≥ 3× the longest outage in the last year (minimum 3 days).
- Consumers that write to external sinks MUST be idempotent (keyed upserts or dedupe on event ID) and MUST route non-retryable failures to an owned DLQ with alerting.
- Any output that affects money is computed or reconciled in batch from the archive, versioned by logic, with a published diff metric against the streaming preview.
SHOULD: Prefer batch unless a consumer needs < 15 min freshness; alert on max-partition lag; review zero-read outputs quarterly for deletion.
Exceptions: Data platform review; time-boxed to two quarters.
Success metrics: schema-caused incidents down 80%; 100% of shared topics archived; stream-vs-batch billing diff < 0.1%; 25% of pipelines retired or consolidated in year one.
What I'd Tell the VP#
We copy our business data into dozens of other systems — search, dashboards, billing, fraud, machine learning — and today those copies sometimes disagree, which is why two teams brought you different revenue numbers last quarter. The root cause isn't the technology; it's that nobody owns the official version of each piece of data or promises not to change it without warning. I'm proposing that each domain team publishes its data with a clear contract, on a shared platform that keeps a permanent history so we can always recompute when something was wrong. It's about six platform engineers, and it should pay for itself by retiring roughly a quarter of our pipelines and cutting warehouse spend. What I need from leadership is to make data quality part of each team's responsibilities, not just the data team's.
Principal Interview Signals#
| Signal | What It Sounds Like |
|---|---|
| Contracts over pipelines | "The fix isn't another pipeline; it's the order team publishing a versioned dataset with a compatibility guarantee." |
| Replayability as policy | "Archive-by-default is the cheapest insurance we'll buy. Not archiving is the only truly one-way door here." |
| Money gets batch | "Anything we pay or invoice is recomputed from the archive with a logic version. The stream is a preview." |
| Retirement | "A third of these jobs have no readers. Lineage tells us which; I'd delete them before building anything new." |
| Pricing freshness | "Seconds cost 5–10× hours. I want each consumer to tell me which one they actually need." |
Staff answers that L7 interviewers find insufficient:
- "I'd add a schema registry with backward compatibility." — Right tool; doesn't assign producers accountability or a deprecation process across 30 consumers.
- "I'd make the sink idempotent and add reconciliation." — Fixes this pipeline; doesn't make it the standard for every money-adjacent pipeline.
- "We'll use Flink for real-time." — No view on which consumers actually need it or what the org pays for 24/7 streaming.
In the Wild#
LinkedIn: Kafka and the Log as the Integration Point#
Kafka was created at LinkedIn to move activity and operational data between systems, and Jay Kreps's widely read essay "The Log: What every software engineer should know about real-time data's unifying abstraction" (2013) argued for an append-only log as the central integration point: producers write once, every downstream system — search, graph, warehouse, monitoring — consumes the log at its own pace and can rebuild from it.
Staff insight: This is the "one replayable log per business fact, every output a derived view" default with its origin story. Citing it lets you justify why the log, not any single database, is the source you design around.
Netflix: Keystone and a Platform for Event Routing#
Netflix has publicly described Keystone, its data pipeline platform built on Kafka and Flink, which routes trillions of events per day from producers to sinks like object storage, Elasticsearch and other Kafka clusters, and offers self-service stream processing so teams don't each operate their own infrastructure.
Staff insight: Keystone is the platform half of the Principal answer: a paved road for routing and processing with self-service on top. The interview-relevant point is that the platform owns the plumbing while teams own their jobs and datasets.
Uber: Apache Hudi for Incremental Processing#
Uber created and open-sourced Apache Hudi to bring upserts, deletes and incremental pulls to its Hadoop data lake, so downstream jobs could process only changed records rather than re-scanning full partitions, cutting data freshness in the lake from many hours toward minutes.
Staff insight: Hudi addresses the gap between batch reproducibility and streaming freshness: keep data in an archive-friendly table format, but make updates and incremental reads cheap. It's a strong reference for "micro-batch over an open table format" as the middle ground.
Practice Drill#
Prompt: "We charge advertisers per click. Today a Flink job counts clicks per campaign per minute and writes them to Postgres; invoices are generated from those rows. Last month a deploy caused a replay and we overbilled ~$300K, and advertisers also dispute that our counts are lower than their own analytics. Fix the pipeline."
Staff Answer
Two bugs, one root cause: the stream's output is being treated as the source of record. Overbilling came from a non-idempotent sink (total = total + delta) replaying after a checkpoint restore. Undercounting comes from late clicks being dropped at window close (mobile clicks often arrive minutes or hours late). Fix the stream: dedupe by click_id with 24h keyed state; event-time windows with 60s bounded out-of-orderness and 30-minute allowed lateness; sink writes absolute window values with ON CONFLICT (campaign_id, minute) DO UPDATE SET clicks = EXCLUDED.clicks so replays overwrite rather than add; late events beyond 30 minutes to a side topic. Alert on checkpoint age > 3× interval and checkpoint failures. Make batch the billing source: every click topic is archived hourly to object storage (Parquet); a daily job recomputes per-campaign totals for day D at D+3 days (picking up late arrivals), dedupes, applies fraud filtering with a recorded logic_version, and overwrites the partition. Invoices are generated only from these batch partitions. Reconcile: publish billing.stream_vs_batch_diff_pct daily; alert above 0.5%; the stream remains the real-time pacing signal and dashboard preview. Disputes: expose per-campaign daily totals with logic version and late-arrival counts so advertisers can reconcile. Owner: ads data team owns the pipeline; finance signs off on the D+3 invoice cutoff. Metrics: checkpoint.age_seconds, late_events.side_output_rate, consumer_lag_seconds max per partition, billing.stream_vs_batch_diff_pct.
Why this is L6:
- Identifies the real design flaw — streaming output used as source of record — rather than only patching the sink.
- Makes replays harmless with idempotent absolute writes and handles late data with event time, lateness and batch correction.
- Puts a business sign-off (invoice cutoff at D+3) and a reconciliation metric on the design.
What L7 adds:
- Makes "money is computed in batch from the archive with a logic version" an org standard, and audits other money-adjacent pipelines (payouts, usage-based billing) for the same flaw.
- Prices it: archiving the click topic costs a few hundred dollars a month in object storage versus a $300K incident; the daily batch costs less than one engineer-day per month.
- Proposes contract-level transparency for advertisers (published counting methodology and versioning), turning a support problem into a product feature.
Staff Interview Application#
How to Introduce This Pattern#
"Before I choose batch or streaming, I want to pin down the source of truth and make sure it's replayable, then ask each consumer two questions: how fresh, and how exact. Fraud needs seconds and can be approximate; invoices need exact and can wait a day. They'll share the source log but not the pipeline — and the money path gets recomputed from the archive."
Lead with the replayable source, then per-consumer freshness and correctness, then duplicates and late data, then schema evolution, backfills and ownership.
When NOT to Use This Pattern#
- Small, single-database workloads: A materialized view or scheduled SQL query is a pipeline without the operational surface. Under ~100 events/sec with one consumer, don't add Kafka.
- Synchronous invariants: If the derived data must be consistent with the source at commit time (balance + ledger), keep it in one transaction.
- One-off analysis: Query the warehouse/lake directly. Productionize only when the question repeats weekly.
- Request/response integrations: If service A needs an answer from service B now, call it; don't route through a pipeline and wait.
Follow-Up Questions to Anticipate#
| Interviewer Asks | What They Are Testing | How to Respond |
|---|---|---|
| "How do you get exactly-once?" | Understanding the boundary | "Replayable source, checkpointed state, idempotent sink keyed by window or event ID. Framework transactions only for Kafka-to-Kafka." |
| "What about late events?" | Event-time reasoning | "Watermark with 60s out-of-orderness, 30-minute allowed lateness with upserts, side output for later — corrected in batch." |
| "How do you backfill after a bug?" | Reprocessing design | "Archive in object storage partitioned by hour; jobs overwrite partitions idempotently; run for the affected range with the new logic version." |
| "A producer changes the schema — what happens?" | Contract thinking | "Registry rejects breaking changes on shared topics; breaking changes become a v2 topic with a 90-day dual publish and tracked consumer migration." |
| "Batch or streaming?" | Cost/need pragmatism | "Whatever the consumer needs. Seconds are 5–10× the cost of hourly; I'd stream only fraud and search here." |
| "A consumer is 6 hours behind — now what?" | Operations | "Max-partition lag alert should have fired at 5 min. Check for a poison message, DLQ it, scale consumers; confirm retention covers the gap." |
Evaluation Rubric#
| Dimension | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Framing | Batch vs stream tool choice | Replayable source; per-consumer freshness and correctness | Domain-owned datasets with contracts; org-wide metric definitions |
| Correctness | "Exactly-once in Kafka" | Idempotent sinks, event time, watermarks, batch for money | Standards that make money pipelines reproducible and versioned everywhere |
| Failure | Retry | DLQ with owner, per-partition lag, retention sizing, backfill runbook | Lineage for blast radius, archive-by-default, quality incidents as outages |
| Evolution | Edit and redeploy | Registry, compatibility, versioned topics | Deprecation policy and producer accountability across 30 consumers |
| Cost | Not discussed | Batch where freshness allows | Usage-based retirement; warehouse and streaming spend governed |
Strong Hire Signals
| Signal | What It Sounds Like |
|---|---|
| Replayability first | "If I can't replay the input, I can't fix the output." |
| Sink-side exactly-once | "Overwrite keyed windows; never increment in the sink." |
| Late-data honesty | "About 2% arrives late; here's where it goes." |
| Ownership of lag | "Every consumer has a lag SLO and an owner, alerted on the worst partition." |
Lean No-Hire Signals
| Signal | Why It Misses the Bar |
|---|---|
| Dual writes from application code to DB and Kafka | Loses events on every crash between the writes |
| Processing-time windows for billing | Miscounts whenever anything is delayed |
| No backfill or retention story | Bugs become permanent data loss |
Common False Positives: Deep Flink state-backend knowledge ≠ choosing whether to stream. Drawing lambda architecture ≠ explaining which path is the source of record. "We have a DLQ" ≠ someone owns draining it.
Capacity Planning Quick Reference#
Sizing the Pipeline#
ingest_MBps = events_per_sec × avg_event_bytes / 1e6
kafka_partitions = max(ceil(ingest_MBps / 10), max_consumer_parallelism)
broker_storage = ingest_MBps × 86400 × retention_days × replication / compression
retention_days >= 3 × longest_outage_days (minimum 3)
archive_per_day = ingest_MBps × 86400 / parquet_compression # 5-10x typical
stream_parallelism = ceil(peak_events_per_sec / events_per_sec_per_slot) # 10-100K per slot
catch_up_time = lag_events / (consumer_capacity - ingest_rate) # must have headroom
Key Numbers Worth Memorizing#
| Number | Context |
|---|---|
| 10–50 MB/s | Practical throughput per Kafka partition |
| 3–7 days | Broker retention; archive covers the rest |
| 128 MB–1 GB | Target Parquet file size; < 32 MB means a small-files problem |
| 30–120 s | Typical bounded out-of-orderness for mobile/IoT event time |
| 1–5% | Share of mobile events arriving after a 1-minute window closes |
| 10–60 s | Common checkpoint interval; also the visibility delay for transactional sinks |
| 5–10× | Cost of streaming vs hourly batch for the same logic |
| 2× | Consumer headroom over ingest needed to catch up after an outage in reasonable time |
| 90 days | Typical deprecation window for a breaking schema change on a shared topic |
| 20–40% | Share of pipelines commonly found with zero readers in a lineage review |
Common Pitfalls Checklist#
- Source is a replayable log (CDC/outbox), not app-level dual writes
- Every topic archived to object storage in an open format
- Sinks idempotent: keyed upserts with absolute values, not increments
- Event time with watermarks; late-data path defined and measured
- Money-affecting outputs computed or reconciled in batch with a logic version
- Schema registry enforces compatibility; breaking changes use versioned topics
- Lag alerted on the worst partition; retention ≥ 3× longest outage
- DLQ per consumer with an owner, alert and replay tool
- CDC replication slots monitored for retained WAL on the source database