Spark DataFrames & Spark SQL
The DataFrame is the primary API of Apache Spark: a distributed table of rows with a named, typed schema, partitioned across a cluster. DataFrames and Spark SQL are two faces of the same engine. Whether you write df.groupBy("country").agg(sum("amount")) in Python or Scala, or SELECT country, SUM(amount) FROM orders GROUP BY country in SQL, Spark builds the same logical plan, optimizes it with the Catalyst optimizer, and executes it with generated code across the cluster.
That shared engine is why DataFrames replaced the older RDD API for almost all work: Spark understands the structure of your data and operations, so it can prune columns, push filters into file scans, reorder joins, and pick join strategies automatically. Writing good Spark code is largely about expressing logic in ways Catalyst can optimize.
TL;DR
- DataFrames are immutable, distributed, schema'd tables, and the SQL and DataFrame APIs compile to the same plans.
- Transformations (
select,filter,join,groupBy) are lazy. Actions (count,collect,write,show) trigger execution. - Catalyst optimizes plans (predicate pushdown, column pruning, join selection), and Tungsten executes them efficiently.
- Prefer columnar formats (Parquet, Delta, Iceberg) with explicit schemas over CSV and JSON inference.
- Use built-in functions over Python UDFs. If you need custom logic, use pandas UDFs (vectorized, Arrow-based).
- Read plans with
explain()and the Spark UI to understand shuffles and scans.
Quick Example
The same logic in Spark SQL:
Core Concepts
Transformations, Actions, and Lazy Evaluation
Transformations build up a plan without computing anything. Actions trigger a job, which Spark splits into stages at shuffle boundaries, and each stage into tasks, one per partition. Laziness lets Catalyst optimize the entire pipeline at once, but it also means errors (bad paths, type issues) may surface only at the action.
Catalyst and Tungsten
Catalyst turns your code into an unresolved logical plan → analyzed plan (resolving columns and types against the catalog) → optimized logical plan (rule-based rewrites: predicate pushdown, constant folding, column pruning) → physical plans (join strategies chosen with statistics), and one plan gets selected. Tungsten then executes it with whole-stage code generation and compact binary memory layouts. Adaptive Query Execution (AQE) re-optimizes at runtime using actual shuffle statistics. See Spark performance tuning.
Schemas and Data Sources
- Parquet / ORC: columnar, compressed, with schema and statistics. They enable column pruning and predicate pushdown.
- Table formats (Delta Lake, Apache Iceberg, Hudi): add ACID transactions, schema evolution, time travel, and efficient upserts on top of Parquet. They're the basis of the data lakehouse.
- CSV / JSON: avoid schema inference in production (it's slow, and types can be wrong). Supply explicit schemas.
Column Expressions and Functions
pyspark.sql.functions offers hundreds of optimized functions: dates, strings, arrays, maps, JSON (from_json, get_json_object), conditionals (when/otherwise, coalesce), and aggregates (including approximate ones like approx_count_distinct and percentile_approx). Expressions are evaluated inside the JVM with code generation. They're far faster than row-by-row Python.
Joins
Spark picks among broadcast hash join (a small side is copied to every executor, with no shuffle of the large side), sort-merge join (both sides shuffled and sorted by key, the default for large-large joins), and shuffle hash join. Join types include inner, left/right/full outer, left semi, left anti, and cross. Watch for skewed keys and exploding many-to-many joins. See Spark partitioning.
Window Functions
Windows compute rankings, running totals, and lags per group without collapsing rows, exactly like SQL window functions. Window.partitionBy(...).orderBy(...) requires a shuffle by partition keys, so be careful with very large or skewed partitions.
UDFs
Best Practices
Filter and Select Early
Keep only the needed columns and rows as early as possible. Catalyst pushes many filters down automatically, but explicit selects clarify intent and help when plans become complex (after UDFs, for example, which block pushdown).
Define Schemas Explicitly
Explicit schemas make reads faster, avoid type surprises, and catch upstream changes loudly. Table formats with enforced schemas help further.
Avoid collect() on Large Data
collect() pulls all data to the driver and can crash it. Use show(), take(n), limit, or write results to storage.
Check Plans for Unexpected Shuffles
df.explain("formatted") shows Exchange nodes (shuffles), broadcast joins, and pushed filters. Verify that expensive operations happen where you expect them.
Common Mistakes
Using Python UDFs for Simple Logic
Recomputing the Same DataFrame
Each action re-executes the full lineage. If a DataFrame feeds several actions, cache() or persist() it (and unpersist() when done), or write an intermediate table.
Case and Null Surprises in Joins
Joining on columns with different types or null keys silently drops rows (nulls never match in equality joins). Normalize types, and handle nulls explicitly, as in SQL null handling.
FAQ
What's the difference between a Spark DataFrame and an RDD?
RDDs are low-level distributed collections of arbitrary objects, without schema information, so Spark can't optimize operations on them. DataFrames have schemas and declarative operations, which Catalyst optimizes and Tungsten executes efficiently. Use DataFrames (or Datasets in Scala) for almost everything.
Should I use Spark SQL or the DataFrame API?
Both compile to identical plans with the same performance. Use SQL for analyst-friendly, declarative queries, and the DataFrame API for programmatic composition, reusable functions, and testing. Many pipelines mix both.
How is a Spark DataFrame different from a pandas DataFrame?
Spark DataFrames are distributed across a cluster, lazy, and immutable, and they're designed for data far larger than one machine's memory. pandas DataFrames are in-memory on a single machine and execute eagerly. The pandas API on Spark (pyspark.pandas) offers pandas-like syntax on Spark's engine.
Why is my Spark job slow even though the code is simple?
Usually it's large shuffles, data skew, too many small files, Python UDFs, or poor partition counts. Check the Spark UI for stages with long tails or spills, and see Spark performance tuning.
Related Topics
- Apache Spark — Pillar overview
- Spark Partitioning — Shuffles, skew, and file layout
- Spark Performance Tuning — AQE, memory, and joins
- PySpark — Python-specific APIs and pitfalls
- SQL Window Functions — The same semantics in SQL
- Data Lakehouse — Table formats and architecture