skip to content

Lineage and Pipeline Observability

Knowing which datasets a table was built from and whether today's run produced sane output. Interviewers raise it with impact-analysis scenarios — "this upstream column is changing, what breaks?" — which you can only answer with lineage rather than guesswork.

on this pageshow

questions

6

A pipeline task finished successfully but wrote a tenth of the usual rows — which metrics catch that?

level: middleimportance: must knowfreq 62%

answer

  1. exit code says the code ran
  2. compare against recent history, not zero
  3. rows in versus rows out
  4. catch the run that never happened
  5. a metric nobody watches is not a check

basics

~20 s

Exit status only proves the code ran. Catching this needs data-level metrics recorded per run: output row counts compared with recent history, input-to-output ratios, null and distinct rates on key columns, and dataset freshness — with thresholds that fail the run.

solid answer

~50 s

A green task means the process exited zero, which says nothing about the data. Pipeline observability adds a second layer of signals recorded **per run**: rows read and rows written, the ratio between them, bytes or partitions written, duration, and distribution facts on important columns — null rate, distinct count, min/max of the date column. Those only become alerts when compared to a baseline. A fixed "at least 1,000 rows" threshold breaks on quiet weekends; a comparison against the same weekday over the last few weeks, or a relative band such as ±40% of the trailing median, survives seasonality. The consumer-facing metric is **freshness**: how old the newest record in the dataset is, or how long since it was successfully produced. It catches the case where a run succeeds and writes nothing at all — the failure a row-count-per-run check can miss entirely. Best of all, assert it *in the pipeline* so the run fails rather than publishing quietly.

code

text · 5 lines
text
run_id=2f9c  job=build_orders_daily  state=SUCCESS  duration=41s
  input  raw.orders            rows_read=1_204_388
  output mart.orders_daily     rows_written=9_133   bytes=2.1MB
  ratio  out/in=0.0076         trailing_median_ratio=0.41
  max_event_ts=2026-08-20T23:59:12Z   dataset_age=17h

go deeper

for a junior

Know that success status and data correctness are different things, and be able to name row count and freshness as the two most basic data-level signals to record per run.

for a middle

Explain how a raw count becomes an alert: baselining against comparable recent runs, input/output ratio bands, and why day-of-week seasonality breaks naive thresholds.

for a senior

Show the diagnosis path — compare rows in against rows out, then walk upstream — and argue for placing the assertion inside the pipeline rather than on a dashboard.

for a principal

Own where the tiering line sits: which datasets get strict, run-blocking checks and which get advisory ones, given the alert-fatigue and engineering cost of covering everything.

## Why a green run proves almost nothing A task's terminal state answers one question: did the process finish without raising. Every interesting data failure passes that test. A source system exports a partial file; a filter predicate now excludes most rows because an upstream enum value changed; a join loses rows because a key column started arriving with different casing; an incremental load reads the wrong window and finds almost nothing. In every case the code runs perfectly and the pipeline goes green. This is the gap that **pipeline observability** exists to close, and it is why "is the pipeline healthy?" and "is the data correct?" are separate questions with separate instrumentation. ## The metrics worth recording on every run Record these as structured run metadata, not as log lines someone greps later: - **Rows in and rows out** per input and output dataset, and the **ratio** between them. A join or filter has a characteristic yield; a change in yield is often the earliest signal. - **Bytes or files written**, and the **partitions touched** — useful when a run writes to the wrong partition entirely. - **Duration**, and where possible duration per stage. A run that finishes in a tenth of the usual time processed a tenth of the usual work. - **Column-level facts** on a handful of important columns: null rate, distinct count, min and max of the event-time column, sum of a key measure. - **Freshness**: the maximum event timestamp present in the output, and the wall-clock age of the last successful production. The last two are what let you answer "is today's data sane?" rather than "did the job run?". ## Turning a number into a signal Raw counts are not alerts. You need a comparison: - **Against history** — the trailing median or mean of the same metric over recent comparable runs. Day-of-week matters: Sunday volume is legitimately a fraction of Tuesday volume, so compare like with like or the check will cry wolf every weekend and be muted. - **Against a business floor** — an absolute lower bound where one genuinely exists, for example a subscription table that must never shrink. - **Against the input** — an output/input ratio band is often more stable than either number alone, because it survives genuine volume swings. For the scenario in the question, the ratio check and the trailing-median check both fire immediately: one-tenth of the usual output is far outside any sane band. A more insidious version — 8% of rows silently dropped — is caught only if the band is tight enough, which is why teams tier checks: hard thresholds on datasets that feed billing or regulators, looser ones elsewhere. ```sql -- volume check against the trailing median for the same weekday with history as ( select run_date, rows_written from pipeline_run_metrics where dataset = 'orders_daily' and extract(dow from run_date) = extract(dow from current_date) and run_date >= current_date - interval '56 days' ) select (select rows_written from today) < 0.6 * (select percentile_cont(0.5) within group (order by rows_written) from history) as volume_alert; ``` ## Freshness is the metric consumers actually feel A row-count check on a run that *did not happen* never fires, because there is no run to check. Freshness is measured on the dataset, not on the run, so it catches skipped schedules, permanently stuck sensors and silently disabled pipelines. Two flavours are worth separating: **data freshness** (how old the newest record is) and **pipeline freshness** (how long since the dataset was last successfully written). A dataset written on time but containing only stale records fails the first and passes the second — usually pointing at a broken source rather than a broken pipeline. ## Where the check belongs A metric on a dashboard that nobody watches is not observability. Put the assertion inside the pipeline, as a task that runs after the write and before anything downstream consumes the result, so that a violated threshold fails the run and leaves downstream work unstarted. Emit the measured values as run metadata regardless of pass or fail, so you accumulate the history that makes tomorrow's baseline meaningful. ## Diagnosing this particular incident With per-run metrics in place, the walk is short. Compare rows in against rows out for the failing run: if input was normal and output collapsed, the loss is inside the transformation — check the filter, the join keys, the window used by the incremental predicate. If input also collapsed, walk upstream via lineage to the producing run and repeat. If input looks normal but the *distinct count of the join key* dropped, you have a key-format change. Every one of those steps needs numbers you only have if you recorded them at the time — reconstructing them after the fact is usually impossible once the table has been overwritten.

  • Why is a fixed minimum row-count threshold a weak volume check?
    It has to be set low enough to survive the quietest legitimate day, which makes it useless on busy days — a 90% drop on a Tuesday still clears a floor calibrated for Sunday. Comparing to a trailing baseline for the same weekday, or to an input/output ratio band, keeps sensitivity without false alarms.
  • What does a freshness check catch that a per-run data check cannot?
    The run that never happened. Row-count and distribution checks execute as part of a run, so a schedule that stopped firing, a stuck sensor or a disabled pipeline produces no check and no alert. Freshness is measured against the dataset's own age, so silence itself becomes the signal.
  • Where should these checks run so a bad result does not reach consumers?
    After the write and before publication, as a task that downstream work depends on. If the assertion fails the task fails, downstream tasks do not start, and the previous good version of the published dataset remains in place. Recording the measured values on every run, pass or fail, builds the baseline the checks compare against.

saying these in an interview costs you the question

  • Treats a green task status as proof the data is correct
  • Proposes a fixed row-count floor with no seasonality handling
  • Puts checks only on a dashboard nobody looks at
  • Never mentions freshness, so a skipped run goes unnoticed
  • Measures only duration and success rate, never data-level metrics

context

open as a page

How do you stop a failed data-quality check from publishing bad rows to downstream pipeline consumers?

level: seniorimportance: must knowfreq 55%

basics

~20 s

Write to a staging location first, run the checks against it, and publish only if they pass — the write-audit-publish pattern. Make the check a task that downstream work depends on, so a failure leaves the previous good version in place and stops dependent tasks from running.

open as a page

What is data lineage in a data pipeline, and what questions does it let a team answer?

level: juniorimportance: should knowfreq 58%

basics

~20 s

Data lineage is the recorded graph of which datasets each pipeline job read and which it wrote. It answers where a table's numbers came from, what breaks if an upstream column changes, and which reports a failed job affected.

open as a page

How does column-level lineage differ from table-level lineage for impact analysis in a data platform?

level: middleimportance: should knowfreq 52%

basics

~20 s

Table-level lineage says table B was built from table A; column-level says B.total came from A.amount and A.qty. Table-level over-reports impact — every downstream consumer looks affected — while column-level narrows a schema change to the handful of fields that actually break.

open as a page

How would you build end-to-end lineage across a platform where several different tools transform data?

level: principalimportance: should knowfreq 33%

basics

~20 s

Pick one interchange format and one metadata store, then feed it from three sources: runtime events emitted by each tool, SQL or query-log parsing where no integration exists, and manual edges for opaque hops. The hard part is agreeing dataset naming so the pieces connect.

open as a page

In the OpenLineage standard, what does a run event carry and what are facets used for?

level: middleimportance: nice to knowfreq 28%

basics

~20 s

An OpenLineage run event ties three entities together — a job, one run of it, and the input and output datasets — plus an event type such as START, COMPLETE or FAIL and a timestamp. Facets are optional typed blocks that attach extra metadata to any of them.

open as a page