Mergeschema and OverwriteSchema
AQE and what it does
- aqe and cached data
- where aqe failes wtih cache
- where does aqe not work
- aqe and cached data
- checkpointing, and its internal process
inferschema - how it works
transformations vs actions
- job, stage, task
- udf
- how they work
- why are they slow (python tax)
row by row vs vectored execution
reducebykey vs groupbykey
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.
- Column added
- A new column is added to destination schema. Older rows get
Nullvalue in the newly added column.
- A new column is added to destination schema. Older rows get
- Column modified (datatype is changed)
- Will fail with schema mismatch
- Will apply type widening (existing datatype is changed to wider datatype)
- Column deleted
- Doesn't delete col in existing table. Just puts value NULL for the deleted column.
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.
- Marking - Marks df as "to be checkpointed"
- Materialization trigger - When action is triggered on this df, spark executes plan upto that point
- Write to storage - Once df is computed, it is written to physical files on disk
- Lineage truncation - delete logical plan that led upto that dataframe
- Reference update - future operation on df will point directly to files on disk, ignoring original files and transformations
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:
- Serialize: Convert JVM rows into Python-readable bytes (Pickling).
- Transfer: Send those bytes across a socket to a Python worker process.
- Deserialize: The Python worker un-pickles the bytes back into Python objects.
- Process & Repeat: After your logic runs, the whole process happens in reverse to get the result back to the JVM.
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 - compaction for small files
- data skipping - file level statistics like min/max values, so that queries can skip files that don't have required rows
- z-order - colocates related value, to improve data skipping and joins
Optimize data layout using OPTIMIZE and liquid clustering.
Liquid clustering
- Replaces partitioning and bucketing.
- Doesn't replace OPTIMIZE
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
- on an older version of delta that doesn't support liquid clustering.
- data is static, only needs to be optmized once
- filtering is done only on one or two specific columns
MergeSchema - leads to schema evolution. Used when appending data to an existing location. Combines existing data with new data.
- Column added
- A new column is added to destination schema. Older rows get
Nullvalue in the newly added column.
- A new column is added to destination schema. Older rows get
- Column modified (datatype is changed)
- Will fail with schema mismatch
- Will apply type widening (existing datatype is changed to wider datatype)
- Column deleted
- Doesn't delete col in existing table. Just puts value NULL for the deleted column.
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:
- Enable AQE
- Optimize joins
- Reduce the amount of data being processed
- Use caching and buffering to re-use lineages
- Use better data serialization formats
- AQE
- Data persistence - Caching and buffering
- Join optimization
One time setups
- Data serialization
- Catalyst optimizer, tungsten
AQE
Is a runtime optmizer for batch optimization tool where data is large. Won't work for streaming pipelines.
Does 3 things:
- Coalesce shuffle partitions
- Change join strategy
- Handle skewed joins
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:
- DISK_ONLY - full awareness. AQE handles everything perfectly. Can coalesce and change join type.
- MEMORYANDDISK - partial awareness. Cannot coalesce, can change join type
- Persisted delta tables - unaware during write
Doesn't work for:
- before first shuffle (as it is a runtime optimizer)
- in-memory cached data
- non equi joins
- threshold mismatches (config issues - workload didn't meet thresholds of skew or small)
- UDFs (non vectorized) are used
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:
spark.sql.adaptive.minPartitionSizespark.sql.adaptive.advisoryPartitionSizeInBytes
CHANGE JOIN STRATEGY
Looks at the sizes of data.
Uses variables in the scope spark.sql.adaptive
- Convert Shuffle sort merge join to Broadcast hash join
- When any join side is smaller than
autoBroadcastJoinThreshold
- When any join side is smaller than
- Convert Shuffle sort merge join to Shuffle hash join
- Partitions are smaller than
maxShuffledHashJoinLocalMapThreshold
- Partitions are smaller than
Note: AQE will ignore the preferSortMergeJoin setting, and go ahead with shuffle hash join when:
advisoryPartitionSizeInBytesis less thanmaxShuffledHashJoinLocalMapThreshold- All partitions are smaller than
advisoryPartitionSizeInBytes
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
- Project = select column
- Predicate = filter condition
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:
- UDFs
- complex expressions
- filtering after aggregation (eg -
HAVING COUNT(*) > 100;) - derives values
WHERE amount * 12 > 1000000;
- filtering after aggregation (eg -
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
- Recomputation is expensive
- The same df is used for many different downstream tasks, so repeating the lineage is necessary
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:
- DISK_ONLY - full awareness. AQE handles everything perfectly. Can coalesce and change join type.
- MEMORYANDDISK - partial awareness. Cannot coalesce, can change join type
- Persisted delta tables - unaware during write
WHAT TO DO?
- Let AQE optimize before cache. Force a wide transformation before caching using an explicit repartition() etc.
.unpersist()dataframes as soon as their use is over
driver oom vs executor oom
Driver OOM symptoms:
- Error messages -
java.lang.OutOfMemoryError: Java heap spaceorGC overhead limit exceeded - collect() shows failure
- broadcasting causes failure (crashes the driver)
- sparkui disappears/crashes (driver is responsible of it)
Executor OOM:
- Error messages -
ExecutorLostFailure,Exit status 137(means os killed process),Container killed by YARN for exceeding memory limits - Skew
- Spark job continues running - gets re-run on another executor and fails repeatedly
- Memory spill - seen in sparkui statistics
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.
- How it works: It merges data locally on each executor before sending anything across the network. If an executor has ten "Apple" records, it reduces them to one "Apple" record with a sum and then shuffles that single record.
- Benefit: Massive reduction in network traffic and memory usage.
groupByKey is a "brute force" operation.
- How it works: It shuffles every single record across the network to the reducing executor. All values for a key are collected into a single list in memory.
- Risk: If you have a skewed key (e.g., millions of records for one "ID"), the executor will run out of memory (
OutOfMemoryError) because it tries to hold that entire list in RAM.
| - | 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:
- Check spark ui's summary metrics
- Duration - look at max vs median. If max is vastly largen than median, there is skew
- Shuffle read size/records - if 75th percentile is 10mb, but max is 100mb, there is skew
- Event timeline - healthy execution has dense block of short bars. Skew will show a few very long bars.
- Disk spills - the large partitions will spill to disk. In tasks, look at "Peak execution memory" and "disk spill"
Solutions
- Enable AQE
- Manually handle the skew - salting (when AQE isn't available, or isn't working well)
- Repartition df to get more even dataframes
Warehouse selection
Use serverless for BI and SQL tasks
2 options:
- Persistent cluster
- Serverless (per job cluster)
Use persistent cluster when workload is predictable and continuously running.
Otherwise, just use serverless compute
Cluster types
- all purpose compute
- For interactive development and exploration (notebooks, dev exploration etc)
- job compute - for production/automated workloads. Terminated when the job is finished.