Pipeline Anatomy: Advanced
What you will be able to do
Pipelines as Products
Apply the pipeline-as-product framing: name the consumer, write the contract, and define the deprecation path before the pipeline ships.
What a Pipeline Contract Contains
| Element | What It Specifies | Why It Matters |
|---|---|---|
| Producer | The team that owns the pipeline and is paged when it fails | Without an owner, no one fixes failures; orphaned pipelines rot |
| Consumer | The named downstream that depends on the output | If no consumer can be named, the pipeline is dead code |
| Schema | Column names, types, nullability, primary key | Consumer code depends on the shape; schema changes break consumers |
| Freshness SLA | How current the data is guaranteed to be (e.g., 'within the last hour') | Consumer planning depends on freshness; without a stated bar, every delay is a crisis |
| Quality SLA | Row count bounds, null thresholds, distribution checks | The pipeline is allowed to fail; what is not allowed is silent corruption |
| Backfill policy | How far back data can be re-derived and at what cost | Consumers ask for fixes to historical periods; the policy answers in advance |
| Deprecation policy | How the pipeline ends: notice period, migration path, sunset date | Without an end-of-life policy, every pipeline runs forever |
Why the Contract Has to Be Written Down
The Producer-Consumer Relationship
When Contracts Fail
- ▸Asking 'who owns this?' returns a long pause or a name of someone who left the company
- ▸The freshness expectation is in someone's head, not in a YAML file
- ▸Consumer teams maintain their own copies of the same logic 'just in case'
- ▸Schema changes are announced in standups rather than in PRs
- ▸Deprecating a pipeline requires a months-long archaeology project
The Maturity Model
| Level | Pipeline Treated As | Visible Behavior |
|---|---|---|
| 0 | A script someone wrote | Owner is whoever last touched it; failures are firefighting |
| 1 | A scheduled job | Owner is named; failures route to a Slack channel; no SLA |
| 2 | A service | Freshness SLA exists; alerts route to on-call; consumers expect uptime |
| 3 | A product | Contract is written; quality SLAs exist; deprecation has a process; new consumers sign on explicitly |
- Write a contract for every new pipeline at the time it is built, not after
- Name the consumer; if no consumer can be named, do not build the pipeline
- Treat schema, freshness, and quality SLAs as PR-reviewable artifacts
- Build pipelines in service of 'someone might want this later'
- Allow contracts to live in tribal knowledge; they leave the company when people leave
- Treat deprecation as an afterthought; build the end-of-life path with the start-of-life path
The Cross-Cutting Undercurrents
Identify the six cross-cutting undercurrents and map each one onto the four pipeline roles.
The Six Undercurrents
Where the Undercurrents Touch Each Role
| Undercurrent | At the Source | At the Transform | At the Consumer |
|---|---|---|---|
| Orchestration | When to extract; how to coordinate with the source's availability | DAG dependencies; retries; partial failure recovery | Notify when fresh data is available; trigger downstream |
| Observability | Did the extract pull the expected volume; source schema check | Did the transform produce the expected row count and distributions | Did the consumer-facing table update; did dashboards stay green |
| Security | Who has read access to the source; how is the credential stored | PII redaction during transform; column-level encryption | Access control on consumer tables; row-level security |
| Data management | Cataloging the source; lineage from source to raw | Lineage between transforms; quality tests at each step | Discoverability of consumer tables; ownership metadata |
| DataOps | Source change management; sandbox environments | Version-controlled transform code; PR review; CI tests | Backwards-compatible schema changes; deprecation process |
| Cost | Egress costs from source systems | Compute cost of transforms; warehouse credits | Storage cost of consumer tables; query cost |
The Cost of Ignoring an Undercurrent
- Runs successfully on the demo, fails silently in production
- Failure modes are discovered by consumers, not by the pipeline
- Cost is unknown until the cloud bill arrives
- Schema changes are deployed without testing
- Failures route to on-call within minutes via observability
- Quality issues caught at the transform, not at the dashboard
- Cost is tagged per pipeline; budget overruns are predictable
- Schema changes go through CI; consumers are notified before deploy
The Senior Engineer's Allocation
The four roles every pipeline has: a source produces data, transforms reshape it, storage holds it, and a consumer reads it. Data flows left to right.
When to Split, When to Merge
Decide when to split a pipeline into multiple DAGs and when to merge several into one, based on ownership, cadence, and operational cost.
Costs of a Single Large Pipeline
| Cost | What It Looks Like |
|---|---|
| Blast radius | One task fails, the whole DAG halts; unrelated downstream work is delayed |
| Deployment friction | Any change requires testing the whole DAG; small changes carry the risk of large ones |
| Ownership ambiguity | When sixty tasks span four teams, no one team owns the DAG |
| Scheduling rigidity | All tasks run on the same cadence; faster cadences are awkward to add |
| Long tail | The end-to-end latency is the sum of every task; tail latency dominates |
Costs of Many Small Pipelines
| Cost | What It Looks Like |
|---|---|
| Cross-DAG coordination | DAG B depends on DAG A; sensors or asset triggers add lag and complexity |
| Discoverability | Twenty DAGs are harder to find than one; lineage is fragmented |
| Operational overhead | Each DAG has its own alerts, runbooks, on-call rotation |
| Drift | Two DAGs that should share logic implement it twice |
| Latency from sensors | Polling sensors add minutes of lag at every cross-DAG boundary |
The Right Place to Split
- ▸Two teams contribute to the same DAG and disagree about deployment cadence
- ▸The DAG mixes hourly and daily work; the slow tasks dominate the schedule
- ▸A single failure in one branch halts unrelated downstream branches
- ▸On-call cannot tell which team to page when the DAG fails
- ▸End-to-end latency is dominated by tasks that no consumer waits for
The Right Place to Merge
The Cross-DAG Contract
The Decision in Practice
- Single team owns the entire pipeline
- All tasks run on the same cadence
- Tasks share most inputs and have no independent consumers
- End-to-end latency is acceptable as the sum of all tasks
- Multiple teams contribute to the work
- Different parts must run at different cadences
- Branches have independent consumers and independent failure tolerance
- Tasks have no shared inputs and could deploy independently
Build vs Buy at Each Layer
Apply the build-versus-buy decision per layer, weighing differentiation, total cost of ownership, and strategic risk.
The Calculus
| Factor | Pushes Toward Build | Pushes Toward Buy |
|---|---|---|
| Differentiation | The capability is core to the company's product | The capability is generic infrastructure |
| Volume | Volume is so high that vendor pricing exceeds engineering cost | Volume is moderate; vendor pricing is reasonable |
| Cycle time | Iteration speed matters more than total cost of ownership | Cycle time is acceptable at vendor pace |
| Compliance | Data cannot leave the company's environment | Vendor offers compliant deployment options |
| Engineer time | Engineers are available and want to build | Engineers are scarce and the build cost is opportunity cost |
Layer-by-Layer
The Hidden Cost of Building
The Hidden Cost of Buying
| Hidden Cost | Build | Buy |
|---|---|---|
| Maintenance | Tail of engineering effort that grows with feature surface | Vendor invoices that scale with usage |
| On-call | Every internal tool needs an on-call rotation | Vendor handles infrastructure on-call; usage on-call remains |
| Lock-in | Locked into the company's own implementation; rewrite is the only exit | Locked into the vendor; switching cost depends on the layer |
| Outage | The team owns every outage and every fix | Vendor outages are external; no internal recourse |
| Hiring signal | Engineers want to work on novel problems, not internal Airflow clones | Engineers want to work on the latest tools, which vendors usually offer |
The Decision Framework
- Buy commodity infrastructure (warehouses, orchestrators, ingestion connectors)
- Build the business logic, the differentiated transform, and the company-specific source connector
- Price the tail: maintenance, on-call, lock-in, outage, hiring signal
- Build a warehouse, an orchestrator, or a generic ingestion framework in 2026
- Buy a black-box solution for the company's most differentiated capability
- Compare a year-zero build cost to a year-zero buy cost; the tail dominates
Redesigning a Tangled Graph
Diagnose a tangled production pipeline graph and apply the lesson's framings to redesign it for operability.
The Symptoms
- ▸On-call gets paged 8 to 12 times per night, half for failures with no clear owner
- ▸Three different teams compute weekly active users; numbers disagree by 4 to 7 percent
- ▸Schema changes in Postgres orders break six different DAGs simultaneously
- ▸The Snowflake bill grew 3.5x year over year with no clear cause
- ▸Adding a new source takes four weeks of cross-team coordination
- ▸Two engineers are leaving and the team is losing institutional memory
Diagnosis: What Each Symptom Reveals
| Symptom | Underlying Cause | Concept From This Lesson |
|---|---|---|
| Pages with no clear owner | Pipelines treated as scripts; no contract names the producer | Pipelines as products |
| Three teams compute WAU differently | No shared curated layer; each team rebuilds the logic | The shared middle layer (intermediate tier) |
| Schema changes break six DAGs | Direct dependencies on raw schema; no decoupling at curated layer | Layered architecture; raw zone decoupling |
| Snowflake bill 3.5x | No cost observability; no cost attribution per pipeline | Cost as an undercurrent |
| Four weeks to add a new source | Each source built bespoke; no ingestion platform | Build vs buy at the ingestion layer |
| Losing institutional memory | Knowledge not captured in contracts, lineage, or catalog | Data management as an undercurrent |
Redesign Step 1: Establish the Layered Shape
Redesign Step 2: Write Contracts for the Top 10 Pipelines
Redesign Step 3: Buy What Should Be Bought
Redesign Step 4: Split the Mega-DAGs
Redesign Step 5: Instrument the Undercurrents
The Result, Six Months Later
- 80 DAGs, no clear ownership
- Three definitions of WAU
- On-call paged 8 to 12 times per night
- Snowflake bill growing 3.5x year over year
- Four weeks to onboard a new source
- Schema changes break six DAGs
- Top 10 pipelines have contracts; remainder triaged for deprecation
- One canonical fact_active_users in the curated layer
- On-call paged 1 to 2 times per night, all with named owners
- Cost attributed per pipeline; growth budgeted and explained
- Two days to onboard a new source via managed ingestion
- Schema changes propagate via the raw layer; one DAG affected
The Underlying Lesson
> A senior data engineer joins a fintech that has 412 production DAGs accumulated over four years. Roughly 80 are mission critical; the other 332 have no clear ownership. The CTO asks: 'How do we make this system operable, and how do we make sure we never end up here again?'
Pipelines are products with owners, contracts, and lifecycles, not scripts that move data
- Category
- Pipeline Architecture
- Difficulty
- advanced
- Duration
- 35 minutes
- Challenges
- 0 hands-on challenges
Topics covered: Pipelines as Products, The Cross-Cutting Undercurrents, When to Split, When to Merge, Build vs Buy at Each Layer, Redesigning a Tangled Graph
Lesson Sections
- Pipelines as Products (concepts: paDataQuality)
A script copies data; a pipeline serves consumers. The difference is not size. The difference is the existence of a contract. A contract names the consumer, names the producer, names what is delivered, names how often, and names what happens when the delivery fails. Pipelines without contracts accumulate, drift, and rot. The accumulated rot is the largest hidden cost in the data engineering organizations of mature companies. The discipline of treating pipelines as products is the only known anti
- The Cross-Cutting Undercurrents (concepts: paDagOrchestration)
The four roles (source, transform, storage, consumer) describe what a pipeline does. They do not describe the cross-cutting concerns that touch every role. Joe Reis and Matt Housley call these concerns 'undercurrents' in their data engineering lifecycle framework, and the term is apt: they run beneath the surface of every layer. A pipeline that addresses the four roles but ignores the undercurrents is a pipeline that works on the demo and breaks in production. Senior engineers spend much of thei
- When to Split, When to Merge (concepts: paDagOrchestration)
Two pipeline architectures are equivalent in what they produce and very different in how they operate. One large DAG with sixty tasks runs as a single unit. Six DAGs with ten tasks each run as separate units. The choice is one of the most consequential architectural decisions a senior engineer makes, and it cannot be made once for all time; the right boundary changes as the system grows. The principle is simple to state and hard to apply: split when the cost of coupling exceeds the cost of coord
- Build vs Buy at Each Layer (concepts: paEltVsEtl)
Every layer of a pipeline can be built in-house or bought from a vendor. The choice is rarely all build or all buy; the right answer differs per layer. Ingestion has mature SaaS options (Fivetran, Airbyte) that solve the boring 80% of source extraction at a real per-row cost. Orchestration has open-source options (Airflow, Dagster, Prefect) that have absorbed most of what custom schedulers used to do. Storage and warehousing have been almost entirely commoditized into Snowflake, BigQuery, Databr
- Redesigning a Tangled Graph (concepts: paDagOrchestration)
The synthesis exercise is a real-shaped problem. A mid-size company has accumulated 80 production DAGs over four years. The data team has grown from three engineers to twelve. The new tech lead has been asked to make the system operable. The exercise walks through the diagnosis and the redesign, using every concept from the lesson and the prior tiers. The Symptoms Diagnosis: What Each Symptom Reveals Redesign Step 1: Establish the Layered Shape The first move is to introduce a shared raw layer a