Hiring BarSupport

Data Pipeline Patterns — Cross-Cutting Pattern

Pattern37 min read5 diagrams

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_seconds p99 < 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.

IntentConstraintStrategyFailure ModeCorrectness Bar
Operational derived views (search index, caches, read models)Freshness 1–10s; must converge to sourceCDC → log → idempotent upserts; periodic reconciliationLag during spikes; silent drift from missed eventsEventually equal to source, drift-monitored
Real-time decisions (fraud, alerting, pricing, rate signals)Sub-second to seconds; value decays fastStream processing with keyed state, windows, watermarksLate or out-of-order events; state loss on failoverBest-effort timely; approximate is OK
Financial / billing aggregates (payouts, invoices, ad spend)Exact, auditable, reproducibleStreaming for preview + batch over the immutable log as the source of recordDouble counting; late data dropped; untracked logic changesExact, reconciled, versioned
Analytics and ML (warehouse, features, training sets)Hours OK; huge volume; schema historyBatch ELT over object storage (Parquet/Iceberg), partitioned by dateSmall files; schema drift; backfill costComplete 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#

StrategyWhat WorksWhat BreaksWho Pays
Nightly/hourly batch (Spark/SQL over object storage)Simple, cheap, reproducible, easy to backfillFreshness in hours; big failures discovered late; long rerunsConsumers who need fresher data; on-call for morning failures
Micro-batch (1–15 min)Most of batch's simplicity at minute freshnessSmall-file problem; scheduling overhead; still not secondsStorage/compaction jobs; platform tuning
Stream processing (Flink, Kafka Streams)Seconds of latency; continuous; event-time correctnessState management, checkpoints, upgrades; 24/7 on-call; harder to backfillPlatform + owning team's on-call; engineering time
Lambda (batch + stream in parallel)Fast preview + correct batch resultTwo codebases computing the same thing; they driftTeam maintaining both; users confused by two numbers
Kappa (stream only, replay to recompute)One codebase; reprocess by replaying the logNeeds long retention; replay at speed can be expensive; some joins are hardStorage for retention; compute during replays
CDC / outboxCaptures every DB change in order without dual writesCouples consumers to source schema; snapshots and re-syncs are operationally fiddlySource DB team (schema changes ripple); platform
ELT into warehouse (load raw, transform in SQL)Raw kept; transformations versioned in SQL; analysts self-serveWarehouse cost grows with careless queries; governance neededFinance (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#

BehaviorSenior (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 moneyOrg-wide data contracts: schema compatibility, ownership, quality SLOs enforced in CI and at the registry
FailureRetry the jobDLQ with owner, lag SLOs, retention ≥ 3× worst outage, backfill runbookDesigns for the org: lineage to find blast radius, standard backfill tooling, quality incidents treated like outages
EvolutionChange the schema and redeployBackward-compatible changes; versioned topics for breaking ones; dual-run and compareDeprecation policy and timelines across 30 consumers; producer accountability for downstream breakage
CostNot 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
OwnershipThe data team owns pipelinesProducers own schemas; consumers own their lag and views; platform owns Kafka/FlinkDomain-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 LineThe Tension
1Freshness vs CorrectnessAnswer now with what has arrived, or wait for completeness
2One Engine vs Two PathsStream-only simplicity vs batch's reproducibility for money
3Push Semantics into the Pipeline vs into the SinkExactly-once machinery vs idempotent writes
4Producer Freedom vs Consumer SafetyShip schema changes fast vs never break downstream
5Retention vs CostReplay 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.

Diagram: Fault Line 1: Freshness vs Correctness

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 SayWhat Interviewers HearWhat 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#

Diagram: 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.

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

TechniqueFreshnessCorrectnessReprocessingOperational Burden
Batch over archiveHoursHigh, reproducibleRe-run partitionLow
Micro-batch1–15 minHighRe-run intervalMedium (small files)
Stream + idempotent sinkSecondsHigh if keyed correctlyReplay from offsetMedium–High
Stream + transactional sinkSeconds + checkpoint intervalHighest inside frameworkReplay from checkpointHigh
Lambda (both)Seconds + hoursBatch winsBatch re-runVery high (two codebases)
CDC consumerSub-second–secondsComplete if slot healthySnapshot + streamMedium (slot management)

Architecture Diagram#

Diagram: 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.

Diagram: Architecture Diagram

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#

FailureDetection SignalBlast RadiusMitigationOwner
Consumer lagconsumer_lag_seconds max per partition > SLOStaleness of that viewScale consumers, prioritize, shedConsuming team
Poison messageoffset_unchanged_seconds, restart loopsOne partition stalledSkip to DLQ, replay after fixConsuming team
Schema breakdeserialize.errors, registry rejectionsAll consumers of topicRoll back producer; versioned topicProducing team
Duplicate countingstream_vs_batch_diff_pct, step after restartAggregates, billingRecompute from archivePipeline owner
Late data droppedlate_events.side_output_rate > 1%Undercounted windowsBatch correctionPipeline owner
CDC slot stalledreplication_slot.retained_wal_bytesSource DB diskRestart connector; drop slot as last resortDB platform + data platform
Retention exceededLag approaching retention; archive gapsPermanent data lossEmergency retention extensionData platform
Small filesFiles per partition > 1K, avg < 32 MBQuery cost and latencyCompaction jobData 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.

OptionWhat WorksWhat BreaksWho Pays
Central data team builds and owns every pipelineConsistent tooling; deep data expertise in one placeTicket queue; owns breakage it can't prevent; domain knowledge lost in translationCentral team's on-call and morale; consumers waiting weeks
Every team builds its own pipelines on its own stackSpeed and autonomyDuplicate copies, conflicting metrics, no lineage, 5 streaming stacksFinance (duplicate spend), executives (conflicting numbers)
Domain-owned datasets with published contracts on a central paved platformProducers accountable for schema and quality; platform provides Kafka, registry, archive, lineage, orchestrationRequires producers to accept on-call for data quality; platform must be genuinely easyDomain 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.

ScaleIngestInfraPeople / On-callRough Monthly Total
Startup1K 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
Growth100K events/s (~1 TB/day)Kafka 9–12 brokers ($15K), stream jobs ($10K), archive ($2K), warehouse ($30K)4 FTE data platform; shared rotation for streaming~$60K + ~$100K people
Large5M events/s (~50 TB/day)Multi-cluster Kafka ($200K), streaming ($150K), archive ($60K), warehouse/lakehouse ($400K+)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#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionDoor TypeReversibility Cost
Batch vs stream for a single consumerTwo-wayRewrite one job; source unchanged
Window sizes, watermarks, latenessTwo-wayConfig + recompute from archive
Not archiving a topicOne-wayData that aged out is gone forever
Event key choice on a shared topic (ordering unit)One-way-ishEvery consumer's ordering assumptions depend on it; requires a new topic
Public metric definitions ("active user", "revenue") in exec reportingOne-way-ishChanging them breaks year-over-year comparisons and trust
Data format in the archive (Parquet + open table format vs proprietary)One-way-ishRe-encoding petabytes; lock-in to one engine
Streaming engine choiceTwo-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:

  1. 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).
  2. Events are produced via outbox/CDC or a validating SDK — never by dual writes in request handlers.
  3. Every topic is archived to object storage in an open format; broker retention ≥ 3× the longest outage in the last year (minimum 3 days).
  4. 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.
  5. 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#

SignalWhat 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 AsksWhat They Are TestingHow 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#

DimensionSenior (L5)Staff (L6)Principal (L7)
FramingBatch vs stream tool choiceReplayable source; per-consumer freshness and correctnessDomain-owned datasets with contracts; org-wide metric definitions
Correctness"Exactly-once in Kafka"Idempotent sinks, event time, watermarks, batch for moneyStandards that make money pipelines reproducible and versioned everywhere
FailureRetryDLQ with owner, per-partition lag, retention sizing, backfill runbookLineage for blast radius, archive-by-default, quality incidents as outages
EvolutionEdit and redeployRegistry, compatibility, versioned topicsDeprecation policy and producer accountability across 30 consumers
CostNot discussedBatch where freshness allowsUsage-based retirement; warehouse and streaming spend governed

Strong Hire Signals

SignalWhat 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

SignalWhy It Misses the Bar
Dual writes from application code to DB and KafkaLoses events on every crash between the writes
Processing-time windows for billingMiscounts whenever anything is delayed
No backfill or retention storyBugs 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#

NumberContext
10–50 MB/sPractical throughput per Kafka partition
3–7 daysBroker retention; archive covers the rest
128 MB–1 GBTarget Parquet file size; < 32 MB means a small-files problem
30–120 sTypical bounded out-of-orderness for mobile/IoT event time
1–5%Share of mobile events arriving after a 1-minute window closes
10–60 sCommon 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 daysTypical 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
  1. Loading the index…