PySpark Projects for Data Engineers

12 Builds With Real Repos

The best PySpark projects for data engineers are real public repositories you can clone and extend, running on data too large for 1 machine, such as a partitioned taxi lake, a TPC-DS warehouse read plan by plan, a Kafka stream and a CDC merge into Delta Lake, with gates on quality and cost. These 12 sit in 4 tiers that build on each other. File layout comes first and joins and skew second, streaming and CDC make up the third tier, and the last keeps jobs healthy, with quality checks and table upkeep alongside cost. Each dossier links the project's own page with its star count and says what it teaches and how to make it yours. It also gives the interview signal and the difficulty, along with the time and cost. 8 of them run on a laptop for $0.

Last updated: Proudly published by: Jeff Wahl24 min read

What makes a PySpark project a data engineering project

A PySpark project belongs in a data engineering portfolio when the data is too large for 1 process and the job has to run again tomorrow. Those 2 conditions force the decisions Spark work is judged on: how files are laid out and partitioned, which join the planner picks and what happens when 1 key holds most of the rows, how a job picks up only new data, how bad data is stopped before anyone reads it, and what each run costs.

Most lists of PySpark projects miss one condition or the other. A 16 million row CSV fits in memory, so Spark adds overhead and teaches nothing a laptop could not, and a churn model in MLlib is a machine learning exercise. Every project here moves and shapes data that someone else reads, and every one is a real public project: a repository you clone, with its star count as its page showed it on 27 September 2026, or the tutorial that hosts its code.

Each dossier says what you build, what it teaches, how to extend it into a project of your own and the signal it sends in an interview. It also gives the difficulty and time for each build, and the cost and stack. The 12 projects sit in 4 tiers that build on each other, and 8 of them run on a laptop for $0. Times are estimates for a first complete build by someone who already writes Python; if DataFrames themselves are new, work through a PySpark tutorial from SparkSession to write first.

The 12 PySpark projects at a glance

#ProjectDifficultyTimeCore stackCost
01NYC taxi trips to Parquet with the Zoomcamp batch moduleBeginner1 to 2 weekendsPySpark · Parquet · GCP$0
02Rerun-safe writes and schema rules with the Delta Lake examplesBeginner+1 weekendPySpark · Delta Lake · Parquet$0
03Joins, shuffles and storage layout with Efficient Data Processing in SparkIntermediate2 weekendsPySpark · Docker · MinIO$0
04A TPC-DS warehouse and the plans behind its queriesIntermediate2 weekendsPySpark · SQL · Parquet$0
05Common Crawl domain statistics from the columnar indexAdvanced2 to 3 weekendsPySpark · AWS · S3Cloud billing
06Event-time streaming with the Spark: The Definitive Guide codeIntermediate1 weekendPySpark · Kafka$0
07Streamify: Kafka to Spark Structured Streaming to BigQueryIntermediate+2 to 3 weeksKafka · PySpark · AirflowCloud billing
08Change data capture into Delta tables with MERGEIntermediate+1 weekendPySpark · Delta Lake · Databricks$0 (free tier)
09Data quality gates with PyDeequIntermediate1 weekendPySpark · Deequ$0
10An Iceberg table kept healthy by scheduled maintenanceIntermediate+1 to 2 weekendsPySpark · Iceberg · Docker$0
11A performance regression gate built on sparkMeasureAdvanced1 to 2 weekendsPySpark · Spark$0
12PySpark jobs on Kubernetes with the Data on EKS Spark stackAdvanced1 to 2 weekendsPySpark · Kubernetes · AWSCloud billing

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

How big the data should be for a PySpark project

Big enough that Spark's choices have consequences: several gigabytes compressed at the least, and ideally more than your machine's memory. Below that, pandas or Polars finishes first, as does DuckDB, and a reviewer knows it. The data on this list clears the bar in different ways: the NYC taxi archive runs monthly from 2009, TPC-H and TPC-DS grow with the scale factor you give the generator, and Common Crawl's corpus holds over 300 billion pages.

Size is not the only thing that makes Spark necessary. A stream that never ends, a table rewritten by thousands of small commits and a join where 1 key holds most of the rows need the same machinery at any size, which is why the streaming and CDC projects can run on generated data. What a project must not be is a single pass over a file that fits in memory with Spark wrapped around it.

Which PySpark project to start with

If your situation is
Pick
Why
You have written PySpark but never run it on real volume
Projects 1 → 2 → 3
Layout, rerun-safe writes and then joins and shuffles on a Docker cluster, in about 4 weekends at no cost
You run Spark at work and want interview material
Projects 3 → 4 → 11
A skew you create and fix, every plan in a benchmark read node by node, and a gate that measures the next change
The roles you want ask for Kafka and streaming, or for CDC
Projects 6 → 7 → 8
Checkpoints and watermarks on local files first, then a Kafka stream, then MERGE, with no broker to debug on day 1
You need 1 project that proves scale and cost
Projects 5 → 12
A monthly web crawl on AWS, then the same questions of executor size and spot capacity on Kubernetes
You're aiming at senior or platform roles
Projects 9 → 10 → 11
Gates on quality and performance, plus table maintenance, show you kept a Spark job healthy after writing it

Running PySpark projects without a cloud bill

Local mode runs the driver and executors in 1 JVM and uses every core of the machine. It runs the same planner, the same adaptive query execution and the same shuffle as a cluster, so a project's plans, and the skewed tasks and spills it's meant to show, all appear in the Spark UI on a laptop. Docker Compose adds the rest: Kafka, MinIO standing in for S3, a History Server, an Iceberg catalog.

What local mode can't show is the network: shuffles over real links, executors lost with their nodes, and the bill. That is why 3 of the projects need a cloud account (Common Crawl on AWS, Streamify on Google Cloud and the Spark on EKS stack) and the CDC tutorial runs in a free Databricks workspace. Take 1 of them once the local projects are done, set a budget alert before the first run, and tear the resources down the same day.

Prepare for the interview
01 / Open invite
02min.

Know PySpark projects the way the interviewer who asks it knows it.

a PySpark projects query, the same shape a screen would give you.
The diff against expected. Where ties broke. What you missed.
sandbox
1SELECT user_id,
2 COUNT(*) AS sessions
3FROM events
4WHERE ts >= NOW() - INTERVAL '7 day'
5
Execute your solution0.4s avg.
CoinbaseInterview question
Solve a PySpark projects problem

File layout and rerun-safe writes

Tier 1 · Projects 1-2

Every later tier assumes 2 habits these projects build: a layout decided on purpose, with an explicit schema, partitions a query can prune and files of a size Spark reads well, and writes that leave the same table however many times a run is retried. Both projects run on a laptop.

Project 01 · Batch foundations

NYC taxi trips to Parquet with the Zoomcamp batch module

Beginner1 to 2 weekendsLocal & free

The batch module of the Data Engineering Zoomcamp. Once Spark is installed, you convert 2020 and 2021 of NYC yellow and green taxi trips from CSV to Parquet under an explicit schema. Revenue and trip counts per hour and pickup zone then come from a groupBy and 2 joins, an outer join of the 2 taxi types followed by a join to the zone lookup. Later lessons run the same jobs on a local standalone cluster, on a Dataproc cluster reading Google Cloud Storage, and into BigQuery.

The preparation notebook declares a StructType for each taxi type instead of trusting inference, then writes each month with repartition(4) into its own year and month folder, so the file layout is decided by hand. The lessons on GroupBy and joins then show what each wide step costs as a shuffle, stage by stage.

Make it yours by extending the layout. Write the trips once with partitionBy on pickup year and month instead of path strings, add the 2025 files with their new cbd_congestion_fee column, and confirm in explain() that a filter on 1 month lists only that month under PartitionFilters. Then size the output files against the 128 MB read partitions Spark packs its input into by default.

Interview signal

"How would you partition this table?" is one of the most common Spark design questions. You answer it with the query pattern you partitioned for and the file counts and sizes you measured, backed by a plan that shows the pruning, and you do it on a course reviewers already know.

PySparkParquetGCPBigQuery
Project 02 · Rerun-safe writes

Rerun-safe writes and schema rules with the Delta Lake examples

Beginner+1 weekendLocal & free

The Delta Lake project's own notebook collection, run in PySpark. Start with the notebooks on append and overwrite and on replaceWhere, then schema enforcement and evolution. After Hive-style partitioning and file sizes comes compaction, and MERGE and the transaction log close the set. Apply each one to the taxi table from project 1 so that every write it makes has a rerun story.

A job that appends blindly writes the same month twice the first time it is retried. replaceWhere atomically overwrites only the rows that match a predicate, so a rerun of March replaces March and leaves the other months alone, and schema enforcement rejects a write whose columns do not match the table instead of quietly widening it. The transaction log notebook shows why both are atomic: each commit is a numbered JSON file in _delta_log, and a reader sees the table before a commit or after it.

Delta Lake 4.3.0 marks the partitionOverwriteMode option legacy in favour of replaceUsing, which matches on key columns, and replaceOn, which matches on a condition. On a current release, rebuild 1 month of the taxi table 3 ways (plain overwrite, replaceWhere, replaceUsing) and record which files each run rewrote.

Interview signal

A backfill can double count a day without any error, so interviewers ask how you would reload 1 bad day without double counting or deleting the rest. You answer with the write you chose, the partitions it touched and the rerun you tested.

PySparkDelta LakeParquet
delta-io/delta-examples242★DataSample data in the notebooks, then your taxi table

Joins, shuffles and query plans

Tier 2 · Projects 3-5

Most Spark performance work is reading what the planner did: which join it chose and where it shuffled, and then which task held the stage up. These 3 projects go from a Docker cluster with a History Server to a benchmark's full query set to a monthly web crawl on AWS.

Project 03 · Joins and shuffles

Joins, shuffles and storage layout with Efficient Data Processing in Spark

Intermediate2 weekendsLocal & free

The code of Joseph Machado's Efficient Data Processing in Spark course. make setup starts a Spark master with 2 workers in Docker, plus a History Server, beside Postgres and MinIO. It then generates TPC-H data with dbgen and loads the tables. The exercises start with query plans and shuffles and how a job splits into stages and tasks. Later ones cover join types, then partition pruning with adaptive query execution, then bucketing and sorting, before table formats and a capstone ETL project with tests.

Because the cluster keeps a History Server, every exercise leaves a run you can open afterwards and read stage by stage: which join the planner chose, where the Exchange nodes sit, which shuffle bucketing removed. Spark tuning depends on that habit of reading the before and after from the job's stages instead of guessing from the code.

Make it yours by breaking it. Regenerate TPC-H at a larger scale factor, skew the join key of lineitem so 1 key holds most of the rows, and run the join with spark.sql.adaptive.skewJoin.enabled off and on. Then fix what adaptive execution does not, a skewed groupBy, with a 2 stage aggregation or a salted key, and record the stage times of each version.

Interview signal

Skew and join strategy are the Spark optimization questions interviewers return to most, and most candidates answer them from memory. You answer with the stage timeline of a skewed task you created, the setting that fixed it and the minutes it saved; the Spark optimization interview questions test exactly this.

PySparkDockerMinIOPostgreSQL
Project 04 · Query plans

A TPC-DS warehouse and the plans behind its queries

Intermediate2 weekendsLocal & free

Build dsdgen from Databricks' fork of the TPC-DS kit, which compiles on Linux and macOS, generate the retail warehouse at a scale factor your machine can hold, and load its fact and dimension tables into Parquet with the large sales facts partitioned by date. Run the benchmark's queries from the set Spark itself tests against, under sql/core/src/test/resources/tpcds in the Spark repository, and publish a results table with 1 row per query, giving the runtime and join strategies under the settings it ran with.

Run explain("formatted") on each query until you can read every node: BroadcastHashJoin against SortMergeJoin, each Exchange as a shuffle, a dynamicpruningexpression in a fact scan where a filtered date dimension skips fact partitions. Then compare the plan Spark started with to the final plan in the SQL tab of the Spark UI, where adaptive execution may have switched a join or coalesced partitions at runtime.

Rerun the slowest queries with the broadcast threshold moved from its 10 MB default and with adaptive execution off, and record what moved. The goal is to explain why a query got faster by pointing at the node that changed.

Interview signal

Interviewers often paste a plan and ask what it does or why a query is slow. After this project, that is a question you've already answered for every query in the benchmark.

PySparkSQLParquet
databricks/tpcds-kit107★DataTPC-DS data generated by dsdgen
Project 05 · Cloud scale

Common Crawl domain statistics from the columnar index

Advanced2 to 3 weekendsCloud account

Common Crawl's own PySpark jobs: count HTML tags and web server names, extract links into a host-level web graph, and query the columnar URL index with SQL. Use the index job to publish a table of domain statistics for 1 monthly crawl, pages per registered domain by HTTP status with language and content type as further splits, and the previous crawl's numbers alongside.

The corpus holds over 300 billion pages in the commoncrawl S3 bucket in us-east-1. The repository says index queries need authenticated S3 access, with no HTTP fallback, and recommends running in that region to avoid transfer costs. The index is Parquet partitioned by crawl and subset, so a filter on both, and a select of only the columns you need, decides how much of it the scan reads; the scan's ReadSchema in the plan lists them.

Test on a sample of files before the full crawl, run the executors on spot capacity, and write down what the finished run cost. The project's result is the table plus the row count and runtime, and the bill for the cluster size you used.

Interview signal

"What is the largest dataset you have processed?" comes up in most data engineering loops. This project turns it into an answer with a row count, a runtime, the cluster it ran on and what the run cost.

PySparkAWSS3Parquet

Incremental, streaming and CDC jobs

Tier 3 · Projects 6-8

These projects process only what's new: files that arrived since the last run, Kafka offsets past the last checkpoint, rows that changed in a source database. Each one has to answer what happens when a batch runs twice or an event arrives late.

Project 06 · Incremental files

Event-time streaming with the Spark: The Definitive Guide code

Intermediate1 weekendLocal & free

The streaming chapters of the book's code repository, in Python. Read a folder of activity JSON files as a stream with maxFilesPerTrigger and write it to the memory and console sinks, then to Kafka. Next, count events in 10 minute tumbling and sliding windows on event time. Once withWatermark is in place, drop duplicate events on user and event time.

The code predates 2 changes worth making as you go. trigger(once=True) is deprecated in favour of availableNow=True, which processes everything available in several batches and then stops, so a scheduled run picks up only the files that arrived since the last one and a failed run resumes from its checkpoint. And the windowed counts run in complete output mode, which keeps every window's state for good, because the watermark cleans state only in append or update mode.

Make it yours by pointing the job at a folder that grows. Drop new files in on a schedule and run it with availableNow. When you delete the checkpoint, it reprocesses everything. Then switch the windowed count to append mode and watch the state stop growing.

Interview signal

Streaming interviews ask about late data and state size, and about what a checkpoint guarantees. Your answer is the watermark threshold you chose and the state you watched grow in complete mode and level off in append mode. The restart you tested covers the checkpoint.

PySparkKafka
databricks/Spark-The-Definitive-Guide3.2k★DataActivity data JSON files, in the repository
Project 07 · Kafka streaming

Streamify: Kafka to Spark Structured Streaming to BigQuery

Intermediate+2 to 3 weeksCloud account

A music streaming service's analytics pipeline on Google Cloud. Eventsim writes listen and page view events, along with authentication events, into 3 Kafka topics; a PySpark Structured Streaming job reads them from the earliest offset and writes Parquet to Cloud Storage every 120 seconds, in hourly partitions under month and day folders; hourly dbt models orchestrated by Airflow build the BigQuery tables a dashboard reads. Terraform provisions the cloud resources.

The streaming job shows the plumbing every Kafka to lake pipeline needs: a declared schema per topic, millisecond timestamps cast to event time, a checkpoint per stream, a processing-time trigger that sets how often files land, and partitions derived from the event rather than from its arrival. Count the files a day of 2 minute micro-batches leaves in each hour's folder and you've a small files problem to fix.

Make it yours where the job stops short. It sets no watermark, so add a 10 minute event-time window of listens per song with withWatermark, restart the job from its checkpoint mid-stream, and check that no window was counted twice. Then compact each finished hour's folder.

Interview signal

1 repository that runs Kafka and Spark under Airflow, with dbt models and Terraform provisioning, covers most of the tools a mid-level posting names, and reviewers have seen the design before. What sets yours apart is the watermark and the restart test you added.

KafkaPySparkAirflowdbtTerraformBigQuery
ankurchavda/streamify919★DataEventsim events over a Million Song Dataset subset
Project 08 · Change data capture

Change data capture into Delta tables with MERGE

Intermediate+1 weekendFree cloud tier

Databricks' CDC pipeline tutorial, installed into a free Databricks workspace with dbdemos.install('cdc-pipeline'). Auto Loader lands change files in a bronze table on an availableNow trigger. A foreachBatch function keeps the latest change per customer id and applies it to a silver table with MERGE, deletes included. The silver table's change data feed then drives a gold table downstream. A second notebook runs the same pattern across many tables at once.

2 details decide correctness, and the notebook handles both. A Delta MERGE can fail when several source rows match the same target row, so each batch is first cut to 1 change per key with ROW_NUMBER() OVER (PARTITION BY id ORDER BY operation_date DESC). And a restarted stream can apply the same batch twice, so the MERGE inside foreachBatch has to leave the same table either way.

With delta.enableChangeDataFeed = true, each changed row carries a _change_type that's one of insert, update_preimage, update_postimage and delete. The gold job reads only those rows instead of rescanning the table. Make it yours by feeding it real change events: run Debezium's Postgres tutorial locally, land its changes as files, and create the awkward cases on purpose, such as a row updated twice in 1 second or a delete followed by a reinsert.

Interview signal

Change data capture comes up in nearly every pipeline design round. You answer it with the out of order updates and deletes you handled, the batch you replayed and the change feed that kept the downstream job incremental.

PySparkDelta LakeDatabricks
Databricks tutorial: CDC pipeline with DeltaDataCustomer change files, generated by the tutorialThe pipeline notebook
The CDC design of project 8, with a database as its source
Capture
Bronze
Silver
Gold
Serving
PostgreSQL
customers db
CDC
change capture
S3
change files
Spark
load bronze
SLA< 1h
Delta Lake
clients cdc
Spark
merge silver
ERRORFail the batch and retry from the checkpointIDEMPOTENCYLatest change per id, then MERGE keyed on id
Delta Lake
retail client silver
custom
silver checks
ERRORStop the gold refresh
Spark
read change feed
Delta Lake
retail client gold
Jupyter
analyst notebooks

The tutorial starts at the change files; the database and capture stages are the extension. Bronze keeps every change, silver keeps 1 current row per customer, and gold reads only the rows the change data feed marks as changed.

Dedupe before MERGEA MERGE fails when 2 source rows match 1 target row, so each batch keeps only the latest change per key first.
Replays are normalA restarted stream can apply a batch twice. The MERGE must leave the same table either way.
Read only what changedThe change data feed hands the gold job only the rows that were inserted or updated, plus deletes, instead of the whole table.

Quality, upkeep and cost

Tier 4 · Projects 9-12

These 4 projects keep a Spark pipeline working after it ships: checks that stop bad batches, maintenance for a table that collects small files, and measurements of each run's time and cost. Each one attaches to a pipeline from an earlier tier.

Project 09 · Data quality

Data quality gates with PyDeequ

Intermediate1 weekendLocal & free

AWS Labs' Python API for Deequ, with a tutorial notebook per module: profile a table, let the constraint suggester propose checks, verify a set of constraints against a DataFrame, and keep every run's metrics in a repository. Put the result between the raw and published tables of project 1 or 3, so a batch is published only when its checks pass and a failing batch moves to quarantine with the reasons.

Each check is an aggregate Spark computes in a pass over the data, not a loop over rows, which is why checks like these stay cheap on hundreds of millions of rows. For the taxi trips, require a pickup time inside the month being loaded, a non-negative fare, a known pickup zone and a row count within a band of recent runs, then decide per check whether a failure stops the pipeline or quarantines the batch, and write down why.

The metrics repository turns checks into history: when a check fires, the stored metrics show the day the distribution started to drift. PyDeequ 1.7.0 and later run on Spark 3.5 and Spark 4.1; Spark 3.1 to 3.4 need PyDeequ 1.6.0.

Interview signal

Data quality comes up in almost every data engineering interview, usually as how you would know the data is wrong. You answer with the checks you wrote and 1 that fired, backed by the metric history that shows when the problem started.

PySparkDeequ
awslabs/python-deequ826★DataYour tables from projects 1 to 3
Project 10 · Table maintenance

An Iceberg table kept healthy by scheduled maintenance

Intermediate+1 to 2 weekendsLocal & free

A Docker Compose environment that starts a Spark notebook server, an Iceberg REST catalog and MinIO. Its Getting Started notebook loads NYC taxi trips into an Iceberg table in PySpark and walks through schema and partition evolution, then snapshots and rollback. Its Table Maintenance notebook runs the rewrite_data_files and expire_snapshots procedures along with rewrite_manifests, then reads the files and snapshots metadata tables.

Append to the table every hour for a week and it collects small files and a snapshot per commit, and query planning slows as both pile up. Count them in the files and snapshots tables. Then compact toward Iceberg's default target of 512 MB per file (write.target-file-size-bytes) and expire snapshots past the default maximum age of 5 days, timing planning and scans before and after each step.

The maintenance notebook calls the procedures from Scala cells, but they are SQL CALL statements, so the same job runs from PySpark through spark.sql. Make it yours by scheduling it so that each run compacts only the partitions written since the last one and then expires snapshots. remove_orphan_files then clears the orphaned files.

Interview signal

Table formats appear in most lakehouse job descriptions. If you ran a maintenance schedule and counted files and snapshots before and after, you can show you've kept a table healthy after writing it.

PySparkIcebergDockerMinIO
Project 11 · Performance gate

A performance regression gate built on sparkMeasure

Advanced1 to 2 weekendsLocal & free

Wrap the pipeline from project 1 or 3 in sparkMeasure's StageMetrics with begin() and end(), run it on a fixed input, and store each run's stage metrics in a table. Elapsed time and executor run time go in it next to shuffle bytes and spill. A gate fails a change that makes the job meaningfully slower or more expensive than its recent runs, and event logs kept for every run let the Spark History Server rebuild its UI after it ends.

Event logs are off by default; spark.eventLog.enabled turns them on. Run the job 5 times on the same input to learn its normal variance, then set the gate outside that band, for example failing when shuffle bytes or disk spill rise more than 20% above the median of the last 5 runs. Executor run time, not wall-clock time, is the number closest to the bill, because it counts every core the job held.

Break a regression down by stage so the gate names the step that got worse, and use TaskMetrics when a stage hides 1 straggler task. sparkMeasure 0.28 supports Spark 3 and 4, and its flight recorder mode collects the same metrics without code changes, writing them to files or InfluxDB, or pushing them to Kafka or a Prometheus Pushgateway.

Interview signal

Senior interviews ask how you would know a change made a job slower or more expensive before the bill arrives. You answer with a gate you built and the regression it caught.

PySparkSpark
LucaCanali/sparkMeasure829★DataRuns of your own jobs
Project 12 · Cost at scale

PySpark jobs on Kubernetes with the Data on EKS Spark stack

Advanced1 to 2 weekendsCloud account

AWS Labs' Spark on EKS stack: 1 deploy script provisions an EKS cluster with Karpenter, the Spark Operator, a Spark History Server and ArgoCD, and can add Prometheus and Grafana or the YuniKorn scheduler. Submit its PySpark example that keeps the driver on on-demand nodes and runs 4 executors of 2 cores and 4 GB on spot capacity. Then run the same job from the on-demand example and from the Graviton and NVMe ones before you tear it all down.

Every choice in the manifests is a cost decision you can measure: executor size against instance size, spot against on-demand, local NVMe for shuffle against network storage. Run the project 1 or 3 pipeline under 2 of them, read executor run time and shuffle spill from the History Server, and turn the difference into dollars per run.

The repository's benchmark folders go further, with a remote shuffle service (Celeborn), native execution engines and Graviton comparisons; pick 1 once the basic stack runs. Keep the cleanup script close, because the cluster bills for every hour it's up.

Interview signal

Platform and senior roles ask how you size executors, why spot capacity suits executors but not the driver, and what a job costs to run. You answer with the node pools you ran, a cost per run you measured and the change that lowered it.

PySparkKubernetesAWSTerraformGrafana

Salt the Hot Merchant

> The daily payment reconciliation Spark job joins 1.2 billion transactions against a 500K-row merchants dimension on merchant_id. It has been failing for three days. Spark UI shows one task processing 38% of all rows while the other 199 finish in seconds. The hot merchant is your company's internal payment processor that handles all driver payouts. You cannot broadcast merchants because a downstream join adds a 2 GB enrichment table. Propose and implement a salting strategy.

How a skewed join shows up in the Spark UI

In a sort-merge join every row for 1 key lands in the same shuffle partition, so a key that holds most of the rows makes 1 task far bigger than the rest. Say 1 join stage has 9 shuffle partitions with a median of 50 MB and 1 partition of 1,210 MB. That task runs about 24 times longer than the others, the stage waits for it, and the stage's task summary in the Spark UI shows a maximum far above the median.

Adaptive query execution, on by default since Spark 3.2.0, treats a partition as skewed when it's larger than 5 times the median and also larger than 256 MB, and splits it into smaller tasks, as the Spark SQL performance tuning guide documents. Here 5 times the median is 250 MB, so the 256 MB floor is the line that counts, and the 1,210 MB partition clears both.

9 shuffle partitions of 1 join stage drawn to scale: 8 sit near the 50 MB median, and 1 of 1,210 MB, 24 times the median, rises past the 256 MB skew line that adaptive query execution uses to split a partition9 shuffle partitions of 1 join stage drawn to scale: 8 sit near the 50 MB median, and 1 of 1,210 MB, 24 times the median, rises past the 256 MB skew line that adaptive query execution uses to split a partition

That handling applies to joins, so a skewed groupBy still needs a 2 stage aggregation or a salted key. A dimension small enough to fall under the 10 MB broadcast threshold removes the shuffle altogether; the guide to PySpark joins, from broadcast to anti joins shows how to request one. Projects 3 and 4 are where to create this skew on purpose and measure the fix.

Common mistakes in PySpark portfolio projects

Small data on a big engine. A dataset that fits in memory runs faster in pandas or DuckDB, and a reviewer knows it. If the input is under a few gigabytes, the project shows you can install Spark, not that you can use it.

Pulling data to the driver. collect() and toPandas() on a large DataFrame move every row to 1 process and fail or crawl. Aggregate in Spark, bring back only the result, and write outputs from the executors.

1 output file. coalesce(1) before a write funnels the final stage through 1 task to produce a single file, which throws away the parallelism the project was meant to show. Size output files with partitioning and maxRecordsPerFile instead.

Appends that double count on rerun. A job that appends blindly writes the same day twice the first time it is retried. Every project on this list has a rerun story: overwrite a partition, MERGE on a key, or read from a checkpoint.

Tuning without a baseline. Changing spark.sql.shuffle.partitions from its default of 200, or adding a cache, proves nothing until you have a before and after from the Spark UI on the same input, and adaptive query execution already adjusts several such settings at runtime.

A clone with nothing changed, or machine learning as the centrepiece. Interviewers recognise the popular repositories and ask what you changed, and a classifier in MLlib is a data science project. For data engineering roles, the valuable part is the pipeline that feeds the model.

A rerun-safe monthly write in PySpark

from pyspark.sql import SparkSession, functions as F

spark = SparkSession.builder.appName("trips_daily").getOrCreate()

def build_month(year: int, month: int) -> None:
    first_day = F.lit(f"{year}-{month:02d}-01").cast("date")
    daily = (
        spark.read.parquet(f"/data/pq/yellow/{year}/{month:02d}/")
        .withColumn("pickup_date", F.to_date("tpep_pickup_datetime"))
        # Keep rows stamped with another month out of this run, or they
        # would replace that day's whole partition
        .where(F.trunc("pickup_date", "month") == first_day)
        .groupBy("pickup_date", F.col("PULocationID").alias("zone"))
        .agg(F.count("*").alias("trips"), F.sum("total_amount").alias("revenue"))
    )
    # Dynamic mode replaces only the partitions this run writes, so a rerun
    # of 1 month rewrites that month and leaves the rest of the table alone
    (
        daily.write.mode("overwrite")
        .option("partitionOverwriteMode", "dynamic")
        .partitionBy("pickup_date")
        .parquet("/lake/trips_daily")
    )

build_month(2021, 3)

The job rebuilds 1 month of daily trip counts per pickup zone from the Parquet the Zoomcamp batch module writes, and replaces only that month's date partitions. Run it for any month, any number of times, and the table holds 1 copy of that month.

How to make a cloned PySpark project your own

Interviewers recognise the popular repositories, so the question after "tell me about this project" is "what did you change?" A clone that runs unmodified answers nothing. Each dossier names an extension that turns the repository into your own project: a skew you created and fixed, a watermark the original job lacks, a real change feed where the tutorial ships files, a gate on a pipeline you built earlier.

A good extension changes 1 decision and measures the result. Swap in data with more volume or messier keys, or add the incremental load or the check the original skipped. Then record before and after numbers from the Spark UI, such as stage time and shuffle bytes, files written and cost per run. Write down what broke on the way, because that note becomes your best interview answer.

How a watermark decides which late events count

A streaming aggregation has to decide how long to wait for late events, and the watermark is that decision. The Structured Streaming programming guide defines it: state for a window ending at T is kept until the latest event time seen, minus the threshold, passes T.

With a 10 minute threshold and a latest event at 12:21, the watermark is 12:11. An event stamped 12:07 arriving now belongs to the 12:00 to 12:10 window, whose state is gone, so it cannot change that count. An event stamped 12:13 still counts in the 12:10 to 12:20 window.

Event-time timeline with 10 minute windows: the latest event at 12:21 and a 10 minute threshold put the watermark at 12:11, so the 12:00 to 12:10 window is closed, an event stamped 12:07 is too late and one stamped 12:13 is countedEvent-time timeline with 10 minute windows: the latest event at 12:21 and a 10 minute threshold put the watermark at 12:11, so the 12:00 to 12:10 window is closed, an event stamped 12:07 is too late and one stamped 12:13 is counted

The watermark frees state only in append or update output mode; in complete mode Spark keeps every window for good, which is how the Definitive Guide's streaming examples run. In append mode each window's result is written once the watermark passes its end, about 10 minutes after the fact, so the threshold is a trade between completeness and latency that you choose and then defend.

How to turn a finished PySpark project into interview answers

  1. 01

    Record the numbers while you build

    Before and after each change, write down the input size and row counts with the runtimes, and the shuffle bytes and file counts, each backed by a screenshot of the Spark UI stage that shows it

  2. 02

    Write the README around 1 decision

    Open with the problem and the data size, then the design decision the project turns on, such as the partition scheme or the skew fix, and the measured result

  3. 03

    Draw the pipeline

    Put a diagram at the top of the README, from sources through jobs to tables and their checks, because it's what an interviewer will ask you to redraw and defend

  4. 04

    Prepare 3 stories

    For each project, prepare 1 thing that broke, 1 thing you made faster or cheaper, and 1 trade-off you would make differently, each with its numbers

  5. 05

    Say them out loud under time pressure

    Explaining a plan or a skew fix while someone interrupts is a separate skill from building it, so rehearse each story before the loop

Practise the questions these projects answer

A finished project gets you through the resume screen and gives you material for the loop, but the loop still tests Spark directly: writing a transformation under time pressure, explaining a plan, choosing a join, handling a late event. Those need timed practice on top of the projects.

Work through the PySpark interview questions with a project in mind for each one, and solve the PySpark practice problems to keep the DataFrame API fast in your hands. When you can tell each project's stories cleanly, run a Spark mock interview to hear how they land with an interviewer asking follow-up questions.

PySpark projects FAQ

What are good PySpark projects for beginners?+
Start with the batch module of the Data Engineering Zoomcamp, which converts 2 years of NYC taxi trips to Parquet under an explicit schema and aggregates them with groupBy and joins, then work through the Delta Lake examples to make every write safe to rerun. Both run on a laptop for $0 and teach the file layout and rerun habits every later Spark project depends on.
How big should the data be for a PySpark project?+
Big enough that the choices Spark makes you take have consequences: several gigabytes compressed at the least, and ideally more than your machine's memory. TPC-DS and TPC-H grow with the scale factor you give the generator, and the NYC taxi archive runs monthly from 2009. A 1 million row CSV fits in pandas, and running it on Spark shows only that you can install Spark.
Can I build PySpark projects without a cloud account?+
Yes. 8 of the 12 projects on this page run on a laptop in local mode or in Docker, including the streaming and Iceberg projects and the performance gate. Local mode uses every core of 1 machine, which is enough to read query plans and watch shuffles and skew happen. Common Crawl and Streamify need a cloud account, as does the Spark on EKS stack, and the CDC tutorial runs in a free Databricks workspace.
Should a PySpark portfolio project include machine learning?+
Not for a data engineering role. Interviewers for these roles ask how the data was partitioned and joined, how it was loaded incrementally, how it was checked and what it cost, and an MLlib model answers none of those questions. If you want an ML angle, build the feature table a model would read and make it correct and fresh without being expensive to recompute.
Is it fine to put a cloned GitHub project on my resume?+
Only after you change it. Interviewers recognise the popular repositories and ask what you did differently, so extend the clone first. Create a skew and fix it, or add the watermark or quality gate the original lacks. Feeding it real data also counts. Then record the before and after numbers. The extension is what the interview conversation will be about.
How do I describe a PySpark project on my resume?+
Lead with the data and the decision, then the measured result: the input size, what you changed (a partition scheme, a broadcast, a skew fix, a MERGE instead of a full reload) and the before and after numbers from the Spark UI. The usual ones are stage time and shuffle bytes, plus the count of files scanned. A line with a number you measured is worth more than a list of tools.
Is local mode enough to learn Spark performance tuning?+
For most of it. Local mode runs the same planner and adaptive query execution as a cluster and shuffles data the same way, so skewed tasks and spills show up in the Spark UI next to the plans. What it can't show is network cost and executor sizing across machines, which is why 1 cloud project, such as Common Crawl on AWS or the Spark on EKS stack, is worth adding once the local ones are done.
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

    Spark rounds are diagnosis rounds

    Skew, shuffles, broadcast judgment, partition pruning, exactly-once streaming. Reading the task table and naming the failure class is the signal, not the API

Related guides