spark - theory

10m read · 1955 words

Theory

todo - 5 stages of generation

MergeSchema vs OverwriteSchema

Overwriteschema - will delete existing schema, and use the new incoming schema only. May corrupt existing table. Older data remains recoverable through time travel.

MergeSchema - leads to schema evolution. Used when appending data to an existing location. Combines existing data with new data.

Type widening formats: link

Source type Supported wider types
BYTE SHORT, INT, BIGINT, DECIMAL, DOUBLE
SHORT INT, BIGINT, DECIMAL, DOUBLE
INT BIGINT, DECIMAL, DOUBLE
BIGINT DECIMAL
FLOAT DOUBLE
DECIMAL DECIMAL with greater precision and scale
DATE TIMESTAMP_NTZ
VOID Any type

Checkpointing

Checkpoint materializes the data and cuts off the lineage.
Persist caches the computed data but keeps the lineage.

Checkpoint is primarily about lineage truncation.

TODO: Role of gc

UDF

Why are UDFs slow? because of the serialization bottleneck

Moving data between jvm (spark) and python is very expensive

To run python code on data:

SOLUTION

Use arrow. Arrow memory format is readable by both jvm and python. So no serde is required.

Arrow is faster because arrow is column based, while pickling is row based. So arrow takes advantage of SIMD etc and gives a tremendous performance improvement.

spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "true")

Delta lake

Schema evolution

Done when a table's MergeSchema is enabled

Optimization

Optimize data layout using OPTIMIZE and liquid clustering.

Liquid clustering

Continuously re-organizes data around clustering columns.

Older approach Liquid clustering
Static partitioning Dynamic clustering
Z-Ordering Dynamic clustering
Need to choose partitions upfront Can change clustering columns
Risk of too many/small partitions Avoids rigid directory partitioning

partition+zorder

Use when

MergeSchema - leads to schema evolution. Used when appending data to an existing location. Combines existing data with new data.

Type widening formats: link

Source type Supported wider types
BYTE SHORT, INT, BIGINT, DECIMAL, DOUBLE
SHORT INT, BIGINT, DECIMAL, DOUBLE
INT BIGINT, DECIMAL, DOUBLE
BIGINT DECIMAL
FLOAT DOUBLE
DECIMAL DECIMAL with greater precision and scale
DATE TIMESTAMP_NTZ
VOID Any type

Optimization

Overall strategy is:

  1. Enable AQE
  2. Optimize joins
  3. Reduce the amount of data being processed
  4. Use caching and buffering to re-use lineages
  5. Use better data serialization formats

One time setups

AQE

Is a runtime optmizer for batch optimization tool where data is large. Won't work for streaming pipelines.

Does 3 things:

AQE looks at statistics of data as it is written to shuffle files.

WHEN DOES AQE NOT WORK
So it doesn't work well if data is cached, because the current state of the df is frozen. So AQE cannot optimize it.

Once data is cached, its execution plan in considered to be fixed. So AQE cannot modify it further.

AQE's effectiveness changes based on where data is stored. The shuffle data write has varying levels of AQE awareness:

Doesn't work for:

COALESCE SHUFFLE PARTITIONS

After a shuffle stage is completed, AQE looks at the partition sizes, and coalesces the partitions smaller than minPartitionNum to achieve target size specified in advisoryPartitionSizeInBytes

Variables are set in:

CHANGE JOIN STRATEGY

Looks at the sizes of data.

Uses variables in the scope spark.sql.adaptive

Note: AQE will ignore the preferSortMergeJoin setting, and go ahead with shuffle hash join when:

HANDLE SKEWED JOINS

AQE detects oversized shuffle partitions, splits them into smaller pieces, and replicates the matching data from the other side so the join work can run in parallel.
It doesn't use salting. It simply breaks down dataframes using split points directly on the dataframes.

Detects skew using variables in scope spark.sql.adaptive.skewJoin

A partition considered to be skewed when:

partition_size > median_partition_size × skewedPartitionFactor
AND
partition_size > skewedPartitionThresholdInBytes

Projection pruning, predicate pushdown

Cut down the data as much as possible before processing it.

The words come from relational algebra

Projection pruning

Reduce the number of columns read.
Spark does this automatically via catalyst optimizer. If a column isn't used in any calculation, then it is simply not even read.

Predicate pushdown

Reduce the number of rows read, by pushing filters as close to data source as possible. Also done automatically by catalyst.

Shows up in df.explain(True) as:

PushedFilters: [GreaterThan(salary,100000)]

Pushdown depends on the data source. It won't work when filters use:

Join optimization

USE BUCKETING

Bucketing improves joins by pre-organizing data into buckets based on a join key, reducing or eliminating the need for a full shuffle during the join.

If the bucketed column is also the basis of the join, then the related rows are already within the same executor. This eliminates the need of a data shuffle.

FORCE BROADCAST JOIN

Using hint when the df is larger than autobroadcast threshold, but will still fit in memory.

from pyspark.sql.functions import broadcast

result = large_df.join(
    broadcast(small_df),
    "customer_id"
)

Cache and persist

Caching will cause AQE to fail.
Use it anyways when

AQE looks at statistics of data as it is written to shuffle files.

So it doesn't work well if data is cached, because the current state of the df is frozen. So AQE cannot optimize it.

Once data is cached, its execution plan in considered to be fixed. So AQE cannot modify it further.

AQE's effectiveness changes based on where data is stored. The shuffle data write has varying levels of AQE awareness:

WHAT TO DO?

driver oom vs executor oom

Driver OOM symptoms:

Executor OOM:

reduceByKey vs groupByKey

groupByKey shuffles all values belonging to a key, whereas reduceByKey performs local aggregation before the shuffle, reducing network traffic and memory usage. Therefore, use reduceByKey when you can express the operation as an associative and commutative reduction.

reduceByKey performs a Map-side Combine.

groupByKey is a "brute force" operation.

- reduceByKey groupByKey
Shuffle data Less More
Map-side combine ✅ ❌
Memory usage Lower Higher
Requires aggregation function ✅ ❌
Typical use Aggregations Need all values

Pro Tip: In the DataFrame API, when you use df.groupBy("key").sum(), Spark is smart enough to use the reduceByKey logic under the hood. You don't have to worry about this choice unless you are working with the low-level RDD API.

Scenarios

Data skews, and how to handle them

To detect skew:

Solutions

Warehouse selection

Use serverless for BI and SQL tasks

2 options:

  1. Persistent cluster
  2. Serverless (per job cluster)

Use persistent cluster when workload is predictable and continuously running.

Otherwise, just use serverless compute

Cluster types

Colophon 1955 words · 10m read
Written as a markdown note in Obsidian. Built into this page by a Python script on 2026-10-02.

Pages