Ingest deliberately
Use spark.read, define a schema, and inspect a sample. CSV is convenient; Parquet/Delta are better long-term formats.
A visual, practical path through distributed data processing—designed to teach both what the code does and why a cluster behaves that way.
A PySpark program can run on your laptop in local mode. The same DataFrame code can later run on a cluster; only the connection and deployment settings change.
1. Install a compatible Java runtime (Spark runs on the JVM). 2. Create a virtual environment so project packages remain isolated. 3. Install PySpark. 4. Write a script. 5. Run it with spark-submit, which is Spark's job launcher.
For learning, master("local[*]") means “use all local CPU cores.” On a managed platform, do not hard-code a master: the platform supplies it.
Use spark.read, define a schema, and inspect a sample. CSV is convenient; Parquet/Delta are better long-term formats.
Use DataFrame operations and assign meaningful intermediate variables. Nothing expensive happens until an action.
Check row counts, nulls, duplicates, totals, and key business rules before writing output.
PySpark is the Python API for Apache Spark. It lets you describe data work in Python while Spark distributes that work across a machine or a cluster.
A huge dataset is split into partitions so several machines can handle it together.
Most PySpark lines build a plan. They do not process data yet.
Calling show(), count(), or write launches a job.
Spark schedules independent partitions on available CPUs and machines.
The driver is your project manager: it reads your code, builds a plan, and assigns tasks. Executors are the workers: they keep data in memory and run those tasks. A partition is one slice of your dataset.
Spark separates planning, resource allocation, task scheduling, and actual data processing. Knowing those layers turns performance tuning from guesswork into diagnosis.
Creates SparkSession, converts your DataFrame calls to logical plans, requests executors, and schedules tasks. Keep large data out of the driver: collect() can crash it.
It decides where executors run and how much CPU/memory they receive. It does not understand your business transformations; Spark's driver does.
Runs many tasks over its lifetime, stores shuffled/cached blocks, and reports status. One executor processes multiple partitions over time according to available cores.
An action such as count() creates a job. Spark splits that job into stages at shuffle boundaries. Every stage is made of tasks; usually one task processes one partition.
A narrow transformation, like filter, lets each output partition be built from one input partition. A wide transformation, like groupBy, needs records redistributed by key—this is a shuffle and becomes a stage boundary.
Use the SQL tab to see query plans and the Stages tab to locate slow tasks. Large “shuffle read/write,” long task tails, spills to disk, and one oversized task point to the usual issues: expensive movement, skew, memory pressure, or unbalanced partitions.
Memory: executor heap is used for execution (sorts, joins, aggregations) and storage (cache). When it fills, Spark may spill intermediate data to disk; that is safer than failure but slower.
A DataFrame is a distributed table with named columns and a schema. Prefer it for nearly all structured data work: it is expressive, optimized, and easy to inspect.
| Choose | When it shines | Remember |
|---|---|---|
| DataFrame | Tables, files, SQL-like analysis | Default choice; Catalyst can optimize it. |
| SQL | Analysts, complex familiar queries | Same engine as DataFrame operations. |
| RDD | Very low-level / unstructured edge cases | More control, fewer automatic optimizations. |
| pandas API on Spark | Gradual pandas migration | Convenient, but learn native PySpark too. |
A schema says which columns exist and what their types mean—string, integer, date, array, and so on. Supplying a schema is safer and faster than asking Spark to guess.
Rule of thumb: inspect early with printSchema() and sample with show(). Bad types silently cause bad aggregations.
Transformations return a new DataFrame; the original stays unchanged. Spark records the lineage and waits. An action is the moment you ask for an answer.
Use built-in functions such as col, when, lower, and to_date. They become part of Spark’s optimized plan.
Each worker first computes partial totals locally. Spark then shuffles matching group keys together and combines the partial totals.
An inner join keeps only matches. A left join keeps every row on the left. Joins often shuffle both datasets by key.
Unlike groupBy, a window keeps the original rows. Use it for rankings, running totals, previous/next values, and “latest record per entity.”
Once the API feels familiar, performance is mostly about where data lives, where it moves, and how much of it exists.
Spark combines compatible transformations, pushes filters close to reads, and only runs on an action. Use df.explain() to see the plan.
groupBy, joins, sorts, and distinct often redistribute data. Reduce shuffles by filtering early and selecting only needed columns.
Use persist() when the same expensive DataFrame feeds multiple actions. Unpersist when done; memory is finite.
Too few: workers idle. Too many: scheduling overhead. repartition() reshuffles; coalesce() reduces partitions cheaply.
Python UDFs are custom escape hatches, but hide logic from the optimizer and add cross-language cost. Favor built-ins or pandas UDFs where appropriate.
Parquet/Delta are columnar and retain types. Partition files by commonly filtered, low-cardinality columns such as event date—not user ID.
Batch processes a bounded dataset on a schedule. Structured Streaming handles new records continuously using the same DataFrame API. MLlib provides distributed feature pipelines and models.
For streaming, understand checkpoints (recovery memory), watermarks (late-data tolerance), and output modes. For all jobs, log metrics and make writes idempotent where possible.