Data Pipeline Architecture Projects

10 Real Systems to Build and Break

The best data pipeline architecture projects are running systems that stay correct when a run repeats, a record changes after it was loaded, a row is deleted or a batch comes in broken. These 10 projects build that from real public repositories and courses, in 4 tiers from a scheduled batch load with backfills to change data capture into an Iceberg lakehouse, and 9 of them run on a laptop for $0. Each dossier says what you build, what it teaches, how to extend it and what it signals in an interview.

Last updated: Proudly published by: Jeff Wahl22 min read

What makes a good data pipeline architecture project

A data pipeline architecture project is a system that keeps moving data correctly when things go wrong: when a run repeats, a record changes after you loaded it, a row is deleted at the source or a batch comes in broken. The 10 projects on this page each build 1 such system from a public repository or course you can open today, grouped in 4 tiers that build on each other. Batch and orchestration come first and operations last, with streaming and change data capture, then the lakehouse, in between.

A script that downloads a file and loads a table is not one of them. It shows the code works once, and the architecture only shows on the second run. Each project here is chosen for the failure it makes you handle: a scheduled load teaches reruns and backfills, an API whose records change teaches merges, a database that deletes rows teaches change data capture, and a table that people query while you write to it teaches the audit gate.

Every dossier gives the same facts. It says what you build and what that teaches, then how to extend it into a project of your own and what signal it sends in an interview. Each also rates its difficulty, estimates the time and cost it takes, names its tools and links to the project's own page and its data. 9 of the 10 run on a laptop with Docker for $0, and the BigQuery project fits in Google's free tier.

The 10 pipeline architecture projects at a glance

#ProjectDifficultyTimeCore stackCost
01Scheduled taxi loads with backfills in KestraBeginner+1 to 2 weekendsKestra · PostgreSQL · SQL$0
02Partitioned taxi assets with backfills and a sensor in DagsterBeginner+1 to 2 weekendsDagster · DuckDB · Python$0
03Incremental GitHub issue loads that absorb edits, with dltBeginner1 weekenddlt · Python · DuckDB$0
04Windowed taxi revenue with Redpanda, PyFlink and PostgresIntermediate+2 weekendsFlink · Redpanda · Python$0
05Change data capture from Postgres with DebeziumIntermediate1 to 2 weekendsDebezium · Kafka · PostgreSQL$0
06A local Iceberg lakehouse on taxi dataIntermediate1 to 2 weekendsSpark · Iceberg · MinIO$0
07Replicating Postgres into Iceberg with Debezium ServerAdvanced2 to 3 weekendsDebezium · Iceberg · PostgreSQL$0
08A write-audit-publish gate on Iceberg branchesIntermediate1 weekendSpark · Iceberg · SQL$0
09Lineage for Airflow with OpenLineage and Marquez, plus freshness targetsIntermediate1 weekendAirflow · OpenLineage · Marquez$0
10Partitioned and clustered taxi tables with a BigQuery cost ceilingIntermediate1 weekendBigQuery · GCP · SQL$0 (free tier)

Star counts as each repository's page showed them on 27 September 2026.

How to choose a pipeline architecture project

Choose by the failure you cannot yet explain, not by the tool you want on your resume. If you cannot say what happens when a load runs twice, the batch tier comes first, because every other project assumes a safe rerun. Once reruns are safe but you have never handled a record that changes or disappears after you loaded it, go to the dlt merge or the Debezium tutorial. With both familiar, the lakehouse and operations tiers are where senior interviews spend their time.

Then check the stack against the job ads you are answering. Kestra and Dagster teach orchestration ideas that carry to Airflow, but a posting that names Airflow is better served by the OpenLineage quickstart, which runs on it. A posting that names Kafka or Flink points at the streaming tier. Projects 6 to 8 suit one that names Iceberg or Databricks, or just says lakehouse.

These projects cover the pipeline side of the craft only. The overall list of data engineering projects covers modelling and analytics engineering as well as platform work, and it shares a few repositories with this one.

Which pipeline project to build first

If your situation is
Pick
Why
You have never scheduled a pipeline
Projects 1 → 3 → 5
A backfilled batch load, then records that change, then records that are deleted: the 3 cases every later design assumes you can handle.
You have 1 weekend
Project 3
The dlt tutorial runs in plain Python against a public API and teaches the merge on a key that makes an incremental load safe.
You already run Airflow or Dagster at work
Projects 5 → 7 → 9
Skip the orchestration basics and capture a database's changes into Iceberg, then add lineage and freshness targets.
The roles you want name Kafka or Flink
Projects 4 → 5 → 7
Event time and watermarks first, then change events from a database log, then those changes applied to a lakehouse.
You are aiming at senior or platform roles
Projects 6 → 8 → 9 → 10
Senior loops ask about table maintenance and a quality gate, then lineage and cost, and all of it builds on 1 local stack.

Batch and orchestration: make every run safe to repeat

Tier 1 · Projects 1-3

These 3 projects each own a slice of data per run, a month of taxi trips or the issues changed since the last load, and each teaches 1 way to make rewriting that slice harmless: a MERGE on a row key, a partition deleted before it is reloaded, a merge on a primary key. Projects 4 to 10 assume you can reload any slice without duplicating it.

Project 01 · Batch orchestration

Scheduled taxi loads with backfills in Kestra

Beginner+1 to 2 weekendsLocal & free

Module 2 of the Data Engineering Zoomcamp: a Kestra flow that loads 1 month of NYC yellow or green taxi trips per run into Postgres. Each run derives its file from the trigger date, loads it into a staging table it truncates first, stamps every row with an md5 unique_row_id over 7 columns, and runs a MERGE that inserts only the rows the target lacks. Cron triggers on the 1st of each month run it, and the same flow backfills the months before the schedule existed. The module runs it locally first, then again into Cloud Storage and BigQuery.

Each run owns a slice of time, not a moment. A flow that reads its month from the trigger date loads the same file on every retry and every backfill, while one that reads the clock loads whatever month it happens to be. The staging table and the key-matched MERGE make the load safe to repeat: rerun January and the target gains 0 rows. The flow also sets a concurrency limit of 1, and the guide to idempotent data pipelines covers the write patterns behind both choices.

To make it yours, break it on purpose. Rerun a finished month and compare row counts. Kill a run between the staging load and the MERGE. Then raise the concurrency limit and backfill a year to see what 2 runs sharing 1 staging table do. Write down each result; the notes are the part of the project that is yours.

Interview signal

"How would you backfill 6 months?" and "what happens if a task runs twice?" are the first follow-ups in a pipeline design round. You can answer both from a flow you have rerun and backfilled, with the row counts before and after. Interviewers know the Zoomcamp, so the failure you caused and the change you made are what they will ask about.

KestraPostgreSQLSQLDockerBigQuery
Project 02 · Partitions and sensors

Partitioned taxi assets with backfills and a sensor in Dagster

Beginner+1 to 2 weekendsLocal & free

Dagster Essentials, the free course of Dagster University: a pipeline of software-defined assets over NYC yellow taxi trips in DuckDB. The trips asset is partitioned by month and reloads a partition by deleting that month's rows before inserting the new file, so any month can be rebuilt on its own. The course adds a schedule, a backfill over the partitions and a sensor that watches a data/requests folder and starts a report job for each new or changed request file.

Partitions turn a backfill into a list of independent, repeatable runs: rebuilding March touches only March, and a failed partition is retried without the other 11. The sensor teaches the other half of batch scheduling, running on arrival instead of the calendar. Its run key joins the file name to the file's modification time, and its cursor remembers the newest time it has seen, so the same file never starts 2 runs and an edited file starts exactly 1 more.

Extend it by pointing the sensor at the monthly taxi file itself: poll for the next month's file, check that it has rows and that its pickup times fall inside the month, and only then materialise the partition. Downstream assets then never read a month that is missing or half written.

Interview signal

Late or missing upstream data is a standard follow-up in pipeline design. This project answers it with a mechanism you ran: a sensor with a cursor, a run key that makes it idempotent, and partitions that rebuild 1 month without touching the rest.

DagsterDuckDBPythonSQL
Project 03 · Incremental loading

Incremental GitHub issue loads that absorb edits, with dlt

Beginner1 weekendLocal & free

dlt's own tutorial for loading data from an API: a Python pipeline that pulls the issues of a GitHub repository into DuckDB. It first loads them all, then switches to incremental loading on updated_at, which fetches only the issues changed since the last run, and to the merge write disposition with id as the primary key, so an issue edited after it was loaded replaces its old row instead of adding a second one.

Most records change after you first load them, the way an issue gets closed or relabelled after the fact. An incremental load keyed on creation time never sees those edits, and one keyed on update time without a merge duplicates every edited row. This project makes you hold both halves: the high-water mark is the largest updated_at loaded, which dlt keeps as pipeline state between runs, and the merge on id makes rereading a record harmless.

Extend it by keeping history. Write every fetched version to an append-only table and only the current version to the merged one, then make the cursor overlap, rereading the last hour on every run so a record committed late on the source is never skipped. The merge absorbs the rereads.

Interview signal

"What happens when a record changes after you loaded it?" is a standard question, and many candidates answer with a full reload. You can answer with a cursor on update time and a merge key, and show a history table holding every version of 1 record.

dltPythonDuckDB
Project 1's monthly load: a staging table and a MERGE on a row key
Schedule
Stage
Publish
orchestrator
monthly trigger
BACKFILL1 run per month from the trigger date
API
monthly trip file
SQL
load staging
PostgreSQL
yellow tripdata staging
SQL
merge on row id
IDEMPOTENCYMERGE on unique row id, insert when not matched
PostgreSQL
yellow tripdata
SLA< 24h

The staging table is truncated at the start of every run, and unique_row_id is an md5 over 7 trip columns, so rerunning or backfilling a month inserts 0 duplicate rows.

The run owns a monthThe file name comes from the trigger date instead of the clock, so a retry loads the same file a backfill would.
The key makes it safeRows already in the target match on their key and are skipped, so running a month twice changes nothing.
Prepare for the interview
01 / Open invite
02min.

Know Pipeline architecture projects the way the interviewer who asks it knows it.

a Pipeline architecture projects query, the same shape a screen would give you.
The diff against expected. Where ties broke. What you missed.
sandbox
1source → bronze → silver → gold
2 ingest : CDC + Kafka
3 transform : dbt + Airflow
4 serve : Snowflake
5
Execute your solution0.4s avg.

Streaming and CDC: process data while it is still an event

Tier 2 · Projects 4-5

Project 4 moves events through a broker into a windowed job, and project 5 turns a database's write-ahead log into events. Both run on a laptop, and both force the questions batch hides: event time vs arrival time, what a late event does, and how a delete travels to every copy.

Project 05 · Change data capture

Change data capture from Postgres with Debezium

Intermediate1 to 2 weekendsLocal & free

The Postgres version of the tutorial in Debezium's official examples: 1 Compose file starts Postgres with an inventory schema alongside Kafka and Kafka Connect, and 1 POST of register-postgres.json registers the connector, which publishes each table's changes to topics such as dbserver1.inventory.customers. You then build what the tutorial leaves to you: a consumer that applies those events to a replica table, and a history table that keeps every version of every row.

Debezium reads the database's write-ahead log instead of querying tables, so it sees deletes and needs no updated_at column. Each event's op field says what happened: r for rows read by the initial snapshot, c for an insert, u for an update and d for a delete. After a delete it also sends a tombstone, an event with a null value, so Kafka's log compaction can drop the older messages for that key. Your consumer handles every case, in order per key, which Kafka keeps because every event for a row carries the same key and lands in the same partition.

Prove it with a script that fires a burst of changes, deletes included, then compare source and replica row for row. Restart the connector halfway through and compare again: it resumes from the log position it stored, so the replica should still match the source.

Interview signal

CDC questions ask how you capture deletes, how you load the history that existed before capture began, and how you keep order. The CDC pipeline interview questions go through them, and this project gives you an event log to answer from.

DebeziumKafkaPostgreSQLDockerPython
Project 4's stream: rides to hourly revenue per zone
Ingest
Process
Serve
API
ride producer
Kafka
rides
Flink
zone revenue 1h
IDEMPOTENCYUpsert on window start and PULocationID
PostgreSQL
zone revenue hourly
SLA< 2h

1 hour tumbling windows on pickup time with a 5 second watermark; the broker is Redpanda, which speaks the Kafka protocol, and each window lands as 1 upserted row per zone.

The broker decouplesThe producer and the job can each fail or scale on their own, and the topic keeps the events until the job restarts.
The watermark decidesA window reports once the watermark passes its end; an event that arrives after that is late for it.

Lakehouse: tables on object storage that keep their history

Tier 3 · Projects 6-7

Both projects run on 1 local stack: a REST catalog and object storage that Spark and Debezium Server write to as Iceberg tables. Project 6 teaches what a table format adds to plain Parquet files, and project 7 keeps such a table in step with a live database, deletes included.

Project 06 · Lakehouse

A local Iceberg lakehouse on taxi data

Intermediate1 to 2 weekendsLocal & free

The Docker environment behind Apache Iceberg's Spark quickstart: Spark with Iceberg, an Iceberg REST catalog and MinIO object storage, with notebooks on port 8888 and months of NYC yellow taxi Parquet files built into the image. Its notebooks create an nyc.taxis table and evolve it in place (renamed columns, a widened type, a new column, a partition field added to a live table), query its snapshots and history, roll it back to an earlier snapshot, and run table maintenance.

A table format is a log of snapshots over plain Parquet files, and every lesson follows from that. Each write adds a snapshot you can query by time travel, which turns a before and after comparison into a query instead of a restore, and a rollback into 1 procedure call. Each small write also adds small files: the maintenance notebook compacts them with rewrite_data_files toward a 50 MiB target and drops old snapshots with expire_snapshots, keeping the last 1.

Extend it into a medallion layout. Bronze holds every month as delivered. Silver deduplicates and types it, with the pickup day as a hidden partition, and gold aggregates daily trips per zone. Then reprocess 1 month, check that no other month's files changed, and count the files each partition holds before and after compaction.

Interview signal

Lakehouse questions go past the diagram: how you reprocess 1 month without rewriting the table, why queries slow down over weeks of small writes, and how long you keep old snapshots. You will have run each of those operations and counted the files before and after.

SparkIcebergMinIODockerSQL
Project 07 · CDC into the lakehouse

Replicating Postgres into Iceberg with Debezium Server

Advanced2 to 3 weekendsLocal & free

Debezium Server Iceberg is a maintained open source sink. It runs Debezium's capture engine and writes change events straight into Iceberg tables without passing them through Kafka or Spark. Point it at the inventory database from project 5 and at the catalog and object store from project 6, turn on upsert mode, and every source table gets an Iceberg table, created on first start when event schemas are enabled, that follows the source row for row.

Upsert mode uses the source table's primary key to delete the old row and insert the new one, and it keeps only the latest event per key within each batch, since a busy key changes several times between commits. Deletes are kept by default as soft deletes, the row's last state with __deleted set to true, so a consumer can tell that a customer went away. Tables without a primary key fall back to append.

Each commit writes small files, so schedule the compaction from project 6 on the replicated tables, and tune the batch size and the batch wait: bigger batches mean larger files, and fewer of them, but the replica lags further behind. Reconcile row counts between Postgres and Iceberg after every scripted burst, and write down what each batch size bought you.

Interview signal

"Replicate an OLTP database into the lake, deletes included" is a common design prompt. You can draw the answer, name the step that collapses repeated changes to 1 row per key, and say what the batch size trades between freshness and file count.

DebeziumIcebergPostgreSQLMinIODocker
memiiso/debezium-server-iceberg332★DataProject 5's inventory database, or any Postgres you runIceberg consumer documentation
The shared lakehouse of projects 6 and 7
Sources
Write
Iceberg tables
Consume
S3
trip parquet files
CDC
inventory postgres
Spark
load taxis
Spark
table maintenance
BACKFILLrewrite data files, then expire snapshots
Iceberg
nyc taxis
Iceberg
inventory customers
SLA< 15min
custom
row reconciliation
MONITORRow count differs from Postgres
Jupyter
notebooks

Spark loads the taxi files in batches; Debezium Server upserts each database change by primary key; 1 maintenance job compacts both tables' small files and expires old snapshots.

Snapshots, not copiesEvery write is a new snapshot, so time travel and rollback both read from the same log, as does an audit branch.
Small files are the costFrequent commits from a change stream leave many small files; compaction on a schedule keeps reads fast.

Six Hours to Refresh Every Number

> We publish credit ratings and financial data to thousands of institutional clients who pay for real-time feeds and historical databases. Our dbt transformation layer has become a bottleneck - full refreshes take 6 hours and we can't meet our real-time client SLAs, but incremental models are giving us stale data when a company's historical prices need retroactive adjustment after a stock split or merger. Design a transformation pipeline that handles both real-time feed delivery and retroactive historical corrections.

+ Source
+ Transform
+ Storage
+ Quality
+ Consumer
+ Queue
Bronze
Silver
Gold
Custom
Pipeline Architecture
Sketch the architecture.

Click or drag a node from the toolbar above. Right-click the canvas for the full menu.

Drag from a node's right port to another node's left port to wire data flow.

Operations: quality gates, freshness, lineage and cost

Tier 4 · Projects 8-10

A pipeline that runs is not yet one people can rely on. These 3 projects add what makes it operable: a gate that keeps a bad batch from readers, a freshness target with an alert and a lineage graph, and a measured cost per query with a ceiling that stops a bad one before it bills.

Project 08 · Quality gate

A write-audit-publish gate on Iceberg branches

Intermediate1 weekendLocal & free

The write-audit-publish notebook in the same Spark and Iceberg environment as project 6. It sets write.wap.enabled on an nyc.permits table of 1,000 NYC film permits, sets spark.wap.branch so the job's writes land on a branch, deletes 1 borough's rows there, audits that exactly 4 boroughs remain, and only then publishes the branch's snapshot to the table with cherrypick_snapshot. It also shows the other ending: a branch that is abandoned is dropped, and readers of the table never see its writes.

Readers of a table should never see a bad batch, not even for the minutes between a load and a failed check. Branches make that a property of the table: writes go to an isolated branch, the checks run against it, and the main table moves only when they pass. A failure leaves the last good snapshot serving every reader, which is late but correct.

Extend it to a load that runs every month: write each taxi month from project 6 to an audit branch, check it with queries (a row count within a band set by earlier months, no null pickup times, pickups inside the month, every zone present in the zone lookup), publish on a pass and alert the owner on a failure.

Interview signal

Interviewers ask where the quality checks sit and what happens when one fails. This project answers both: before publish, on an isolated branch, with the last good snapshot still serving every reader.

SparkIcebergSQLDocker
docker-spark-iceberg, WAP with branches notebook387★DataNYC film permits sample, 1,000 rows (bundled)
Project 09 · Lineage and freshness

Lineage for Airflow with OpenLineage and Marquez, plus freshness targets

Intermediate1 weekendLocal & free

The OpenLineage quickstart for Airflow: Marquez runs in Docker, the OpenLineage provider sends it the datasets each task run reads and writes along with the run's state, and 2 DAGs give it something to track: counter inserts a row into a counts table every minute and sum totals them into sums every 5 minutes. The quickstart then breaks the pipeline the way teams do, renaming a column upstream, and uses Marquez's lineage graph and dataset versions to find what changed and when.

The orchestrator reports lineage as tasks run, so the graph shows what actually ran, including a failed job and every downstream dataset it left stale. Add a freshness target to each table and check it with dbt source freshness using warn_after and error_after: checked every 30 minutes, a 60 minute target catches a stalled table no later than 90 minutes after its last good load.

Then stop the counter DAG for 3 hours and record what each piece reports: the freshness check moving from warn to error, and in Marquez the failed job and the stale sums table downstream of it. The guide to data observability explains the checks behind this in depth.

Interview signal

Operations questions ask how you would know a pipeline is broken before its users do. This project answers with a target, a detection time you computed and a lineage graph that shows how far 1 failed job reaches.

AirflowOpenLineageMarquezdbtPostgreSQLDocker
OpenLineage quickstart for AirflowDataThe quickstart's counts and sums tables (bundled)MarquezProject/marquez
Project 10 · Cost

Partitioned and clustered taxi tables with a BigQuery cost ceiling

Intermediate1 weekendFree cloud tier

Module 3 of the Data Engineering Zoomcamp: create an external table over the taxi CSV files in Cloud Storage, then materialise the trips 3 ways and run the same queries on each. The baseline has no partitions. The second copy is partitioned by DATE(tpep_pickup_datetime), and the third is partitioned and clustered by VendorID. The module's own notes record the difference: 1 query scans 1.6 GB on the unpartitioned table and about 106 MB on the partitioned one, and another scans 1.1 GB on the partitioned table and 864.5 MB once it is also clustered.

With on-demand pricing BigQuery bills the bytes a query reads, so table layout is a cost decision: the partition filter cut 1 scan about 15 times. LIMIT does not help on a table that is not clustered, because BigQuery still bills every byte the query reads. The guard belongs in the pipeline: run each query as a dry run first, and set maximum_bytes_billed so a query that would scan too much fails without a charge.

Extend it into a daily job that appends 1 day to the partitioned table instead of rebuilding it, and log the bytes each run processed. The BigQuery sandbox needs no billing account and gives 10 GiB of storage and 1 TiB of queries a month, but it expires tables and partitions after 60 days and runs no DML, so a MERGE version of the load needs billing turned on.

Interview signal

Senior pipeline interviews ask how you would make a design cheaper. You can answer with the bytes a query scanned before and after partitioning and clustering, and with the guard that stops a bad query before it bills.

BigQueryGCPSQLPython
Project 8's gate on a monthly load, with project 9's lineage and alerts
Write
Audit
Publish
Observe
S3
monthly trip files
Spark
load month
BACKFILLReload 1 month per run
Iceberg
trips audit branch
custom
month checks
ERRORKeep main on the last good snapshot
Slack
owner alert
Spark
publish snapshot
Iceberg
trips main
SLA< 24h
Metabase
zone dashboards
Marquez
lineage graph

Readers only ever query main; a failed check leaves main on its last good snapshot, keeps the branch for inspection and alerts the owner, while every run reports its inputs and outputs to the lineage graph.

Checks before publishThe audit runs on an isolated branch, so a failure costs freshness, never correctness.
Detection has a numberA 60 minute target checked every 30 minutes catches a stall within 90 minutes.
Cost has a ceilingMaximum bytes billed fails an oversized query before it runs, without a charge.

Common mistakes in data pipeline portfolio projects

The most common mistake is a pipeline that has only ever run once. A first run on clean data proves the code works; it says nothing about the architecture, which is how the system behaves on the second run, on a late file, when an event arrives twice or a row is deleted. Rerun every project on purpose at least once and write down what happened.

Appending without a key comes close behind. Retries, backfills and at-least-once delivery all deliver some data twice, so every write needs either an overwrite of a slice it owns or a merge on a key. Rerun the 14:00 partition of an hourly table: an overwrite leaves 1 copy of the hour, an append leaves 2, and every total built on that hour doubles.

4 hourly partitions from 12:00 to 15:00 rerun at 14:00 2 ways: an overwrite leaves 1 copy of the hour, an append stacks a second copy, shown in red as a duplicate4 hourly partitions from 12:00 to 15:00 rerun at 14:00 2 ways: an overwrite leaves 1 copy of the hour, an append stacks a second copy, shown in red as a duplicate

Reading the clock fails more quietly. A task that computes its range from now() instead of its scheduled interval, or a stream that windows by arrival time instead of event time, gives a different answer every time it is replayed, and replays are how pipelines recover.

Portfolios also lean on tool collages and unmodified tutorials: a diagram with 9 logos and no failure story reads as a course followed to the end. 3 tools and a written account of what broke, how you found it and what you changed make the stronger project.

How to turn a tutorial into your own pipeline project

Most of the projects on this page come from courses and official examples, and interviewers know them. What makes one yours is the part the tutorial leaves out. Keep the architecture but point it at a source you care about, and the design decisions become yours to defend, because a new source has its own late files and schema changes, and its own duplicates.

Then give it an operating contract. Write the grain and the freshness target in 1 line, such as 1 partition per month available within a day and safe to rerun, and check the target at least twice per window, as dbt's source freshness guidance recommends: a 60 minute target means a check every 30 minutes.

Last, keep an incident log. Every failure you cause on purpose gets 3 lines: what you did, what the system did, and what you changed. That log is the README section reviewers read first and the source of every answer in the interview.

How to turn a finished pipeline project into interview answers

  1. 01

    Write the contract in 1 line

    Give the grain with its freshness target, such as 1 partition per month available within a day, and then the guarantee, such as safe to rerun.

  2. 02

    Break it on purpose and keep notes

    Kill a run mid-write and replay a month, then send a duplicate and delete a source row, and record what the system did each time and what you changed.

  3. 03

    Measure 3 numbers

    Record the freshness lag at each check, the row difference after a rerun and the bytes or minutes each run costs.

  4. 04

    Draw it the way the round asks

    Redraw the system from sources through queues and transforms into storage, then past the quality gates to consumers. Mark the SLA on the drawing and draw the failure path too.

  5. 05

    Rehearse the 4 follow-ups

    Answer out loud how you would backfill it, handle late data, absorb a schema change and survive 10 times the volume.

From your project to the pipeline design round

A finished project gives you 1 system you know completely, and the design round asks you to reason about systems you have never seen. Practise that transfer. Take a prompt from the data pipeline practice problems, which hold {COUNT:pipeline_architecture} pipeline architecture problems, draw the design with the same parts these projects use, and check it against the follow-ups in the data pipeline interview questions.

The strongest answers borrow from what you built: the key that made a rerun safe, the watermark delay you measured, the branch that kept a bad month out of a dashboard. The system design round guide covers how that round is run and what it scores.

Data pipeline architecture projects FAQ

What is a data pipeline architecture project?+
A project whose deliverable is a running system that moves data for someone else. It ingests from a real source and lands the data, then transforms and checks it before serving it. It has to keep working when a run repeats or data arrives late, and when rows are deleted or a job fails. The architecture is the set of decisions that make that true: how each write stays safe to repeat, how history is loaded, where the quality gate sits and what freshness the consumer can rely on.
Which data pipeline project should I build first?+
A scheduled batch load with backfills, such as the Kestra module of the Data Engineering Zoomcamp (project 1). Every other project assumes you can reload any slice of data without duplicating it, and this is the smallest system that forces you to learn that. With 1 weekend only, the dlt incremental load (project 3) teaches the same lesson for records that change after you load them.
Do pipeline projects need to run in the cloud?+
No. 9 of the 10 projects on this page run on a laptop with Docker for $0, including the Flink stream, Debezium change data capture and the Iceberg lakehouse. The exception is the BigQuery cost project, which fits inside the free sandbox and the free monthly query allowance. Interviewers judge the decisions you can defend, not where the containers ran.
How many pipeline projects do I need for a portfolio?+
2 or 3 finished ones. Between them they should cover a batch pipeline with backfills and either streaming or change data capture, plus 1 operations project such as a quality gate or freshness targets. A project you have broken on purpose and rerun after the repair is worth more than several that have only run once.
Is it fine to build a portfolio project from a course or tutorial?+
Yes, if you extend it. Interviewers recognise the Zoomcamp and the official examples, so an unmodified clone says little. Change the data and add the failure the tutorial skips. Then measure what happened and write it down, because the extension is what the interview conversation will be about.
Should I learn batch or streaming pipelines first?+
Batch. The ideas streaming depends on, such as idempotent writes, event time vs arrival time and replaying history, are easier to see when each run owns a fixed slice of data. Once a batch pipeline survives reruns and backfills, a stream is the same problem with the slices made small and continuous.
What is the difference between a pipeline project and a data engineering project?+
Every pipeline project is a data engineering project, but many data engineering projects are about modelling or analysis. A pipeline architecture project is judged the way the pipeline design round in an interview judges it: on how the system behaves over time, on the second run, when a file arrives late, and when a row is deleted or a batch is bad.
How do I talk about a pipeline project in an interview?+
Lead with the contract, such as 1 partition per month, available within a day, safe to rerun. Then draw the architecture, name the 1 failure you caused on purpose and what the system did, and quote 3 measured numbers: the freshness lag, the row difference after a rerun and the cost per run.
02 / Why practice

The candidate who gets the offer

  1. 01

    Reading a solution is not the same as writing one

    Every engineer who has frozen on a query they had read a dozen times knows the gap. The only preparation that closes it is producing the answer yourself, under time, before the interview does it for you

  2. 02

    76% of hiring managers reject on the coding task, not the resume

    From HackerRank's 2024 Developer Skills Report. Candidates who look strong on paper still fail the live screen if they haven't done timed, executable practice

  3. 03

    System design comes down to the calls you defend out loud

    Ingestion, batch vs streaming, the bronze/silver/gold layers, idempotency, backfill and replay. Sketching the pipeline and naming the failure modes is the signal, not the boxes

Related guides