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

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

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

References