Batch Processing Optimizer
Optimize batch data processing workloads using Spark, Polars, DuckDB, and pandas with focus on memory efficiency, parallelism, and cost reduction.
Activation Triggers
Activate on: "batch processing", "Spark optimization", "Polars", "DuckDB", "pandas performance", "data frame", "shuffle optimization", "partition skew", "memory optimization", "out of memory"
NOT for: Real-time streaming → streaming-pipeline-architect | Warehouse SQL tuning → data-warehouse-optimizer | Pipeline orchestration → airflow-dag-orchestrator
Quick Start
- Choose the right tool — DuckDB for single-node analytics, Polars for DataFrames, Spark for distributed
- Profile first — identify bottlenecks (shuffle, skew, memory) before optimizing
- Reduce data early — filter and select columns as early as possible in the pipeline
- Avoid shuffles — broadcast small tables, pre-partition data, use map-side joins
- Right-size resources — match executor memory/cores to actual data size
Core Capabilities
| Domain |
Technologies |
| Distributed |
Apache Spark 3.5+, Dask, Ray |
| Single-Node |
DuckDB 1.1+, Polars 1.x, pandas 2.2+ |
| File Formats |
Parquet, Arrow IPC, Delta Lake, Iceberg |
| Optimization |
AQE (Spark), lazy evaluation (Polars), columnar scans |
| Cloud |
Databricks, EMR, Dataproc, serverless Spark |
Architecture Patterns
Tool Selection Decision Tree
Data Size?
├─ < 10 GB → DuckDB (SQL) or Polars (DataFrame)
│ Single machine, zero setup, fastest iteration
│
├─ 10-100 GB → Polars (lazy) or DuckDB (out-of-core)
│ Still single machine with spill-to-disk
│
└─ > 100 GB → Spark (distributed)
Multi-node cluster, shuffle-based joins
Complexity?
├─ SQL-centric → DuckDB (fastest SQL engine for analytics)
├─ DataFrame → Polars (10x faster than pandas, lazy evaluation)
└─ Complex ML → Spark + MLlib or Spark + Ray
Spark Optimization Patterns
from pyspark.sql import SparkSession
import pyspark.sql.functions as F
spark = SparkSession.builder \
.config("spark.sql.adaptive.enabled", "true") \
.config("spark.sql.adaptive.coalescePartitions.enabled", "true") \
.config("spark.sql.adaptive.skewJoin.enabled", "true") \
.getOrCreate()
# GOOD: broadcast small dimension table (< 100MB)
from pyspark.sql.functions import broadcast
result = large_df.join(broadcast(small_dim_df), "key")
# GOOD: predicate pushdown — filter before join
orders = spark.read.parquet("s3://data/orders/") \
.filter(F.col("order_date") >= "2026-01-01") \
.select("order_id", "customer_id", "amount") # column pruning
# BAD: collect() on large dataset — causes OOM on driver
# all_data = large_df.collect() # NEVER do this
# GOOD: write partitioned output
result.repartition(200) \
.write.mode("overwrite") \
.partitionBy("order_date") \
.parquet("s3://output/results/")
Polars Lazy Evaluation
import polars as pl
# Lazy mode: builds query plan, optimizes, then executes
result = (
pl.scan_parquet("data/orders/*.parquet") # lazy scan
.filter(pl.col("order_date") >= "2026-01-01")
.join(
pl.scan_parquet("data/customers/*.parquet"),
on="customer_id",
how="inner"
)
.group_by("region")
.agg([
pl.col("amount").sum().alias("total_revenue"),
pl.col("order_id").n_unique().alias("order_count"),
])
.sort("total_revenue", descending=True)
.collect() # executes optimized plan
)
# Polars optimizes: predicate pushdown, projection pushdown,
# join reordering — all automatically via lazy evaluation
Anti-Patterns
- pandas for >5GB — pandas loads everything into memory; use Polars (lazy) or DuckDB for medium data, Spark for large
- Collect to driver —
df.collect() or df.toPandas() on large Spark DataFrames causes OOM; aggregate first
- Ignoring partition skew — one partition with 10x more data than others bottlenecks the entire job; use AQE or salting
- Reading all columns — always select only needed columns; Parquet columnar format skips unused columns entirely
- Tiny output files — too many small output files (< 128MB) slow downstream reads; coalesce before writing
Quality Checklist
1---2name: batch-processing-optimizer3description: Spark, pandas, polars, DuckDB optimization for batch data processing. Activate on: batch processing, Spark optimization, polars, DuckDB, pandas performance, data frame, shuffle, partition, memory optimization. NOT for: streaming pipelines (use streaming-pipeline-architect), warehouse queries (use data-warehouse-optimizer).4license: Apache-2.05---6
7# Batch Processing Optimizer
8
9Optimize batch data processing workloads using Spark, Polars, DuckDB, and pandas with focus on memory efficiency, parallelism, and cost reduction.
10
11## Activation Triggers
12
13**Activate on:** "batch processing", "Spark optimization", "Polars", "DuckDB", "pandas performance", "data frame", "shuffle optimization", "partition skew", "memory optimization", "out of memory"
14
15**NOT for:** Real-time streaming → `streaming-pipeline-architect` | Warehouse SQL tuning → `data-warehouse-optimizer` | Pipeline orchestration → `airflow-dag-orchestrator`
16
17## Quick Start
18
191. **Choose the right tool** — DuckDB for single-node analytics, Polars for DataFrames, Spark for distributed
202. **Profile first** — identify bottlenecks (shuffle, skew, memory) before optimizing
213. **Reduce data early** — filter and select columns as early as possible in the pipeline
224. **Avoid shuffles** — broadcast small tables, pre-partition data, use map-side joins
235. **Right-size resources** — match executor memory/cores to actual data size
24
25## Core Capabilities
26
27| Domain | Technologies |
28|--------|-------------|
29| **Distributed** | Apache Spark 3.5+, Dask, Ray |
30| **Single-Node** | DuckDB 1.1+, Polars 1.x, pandas 2.2+ |
31| **File Formats** | Parquet, Arrow IPC, Delta Lake, Iceberg |
32| **Optimization** | AQE (Spark), lazy evaluation (Polars), columnar scans |
33| **Cloud** | Databricks, EMR, Dataproc, serverless Spark |
34
35## Architecture Patterns
36
37### Tool Selection Decision Tree
38
39```
40Data Size?
41 ├─ < 10 GB → DuckDB (SQL) or Polars (DataFrame)
42 │ Single machine, zero setup, fastest iteration
43 │
44 ├─ 10-100 GB → Polars (lazy) or DuckDB (out-of-core)
45 │ Still single machine with spill-to-disk
46 │
47 └─ > 100 GB → Spark (distributed)
48 Multi-node cluster, shuffle-based joins
49
50Complexity?
51 ├─ SQL-centric → DuckDB (fastest SQL engine for analytics)
52 ├─ DataFrame → Polars (10x faster than pandas, lazy evaluation)
53 └─ Complex ML → Spark + MLlib or Spark + Ray
54```
55
56### Spark Optimization Patterns
57
58```python
59from pyspark.sql import SparkSession
60import pyspark.sql.functions as F
61
62spark = SparkSession.builder \
63 .config("spark.sql.adaptive.enabled", "true") \
64 .config("spark.sql.adaptive.coalescePartitions.enabled", "true") \
65 .config("spark.sql.adaptive.skewJoin.enabled", "true") \
66 .getOrCreate()
67
68# GOOD: broadcast small dimension table (< 100MB)
69from pyspark.sql.functions import broadcast
70result = large_df.join(broadcast(small_dim_df), "key")
71
72# GOOD: predicate pushdown — filter before join
73orders = spark.read.parquet("s3://data/orders/") \
74 .filter(F.col("order_date") >= "2026-01-01") \
75 .select("order_id", "customer_id", "amount") # column pruning
76
77# BAD: collect() on large dataset — causes OOM on driver
78# all_data = large_df.collect() # NEVER do this
79
80# GOOD: write partitioned output
81result.repartition(200) \
82 .write.mode("overwrite") \
83 .partitionBy("order_date") \
84 .parquet("s3://output/results/")
85```
86
87### Polars Lazy Evaluation
88
89```python
90import polars as pl
91
92# Lazy mode: builds query plan, optimizes, then executes
93result = (
94 pl.scan_parquet("data/orders/*.parquet") # lazy scan
95 .filter(pl.col("order_date") >= "2026-01-01")
96 .join(
97 pl.scan_parquet("data/customers/*.parquet"),
98 on="customer_id",
99 how="inner"
100 )
101 .group_by("region")
102 .agg([
103 pl.col("amount").sum().alias("total_revenue"),
104 pl.col("order_id").n_unique().alias("order_count"),
105 ])
106 .sort("total_revenue", descending=True)
107 .collect() # executes optimized plan
108)
109
110# Polars optimizes: predicate pushdown, projection pushdown,
111# join reordering — all automatically via lazy evaluation
112```
113
114## Anti-Patterns
115
1161. **pandas for >5GB** — pandas loads everything into memory; use Polars (lazy) or DuckDB for medium data, Spark for large
1172. **Collect to driver** — `df.collect()` or `df.toPandas()` on large Spark DataFrames causes OOM; aggregate first
1183. **Ignoring partition skew** — one partition with 10x more data than others bottlenecks the entire job; use AQE or salting
1194. **Reading all columns** — always select only needed columns; Parquet columnar format skips unused columns entirely
1205. **Tiny output files** — too many small output files (< 128MB) slow downstream reads; coalesce before writing
121
122## Quality Checklist
123
124- [ ] Tool matches data size (DuckDB/Polars < 100GB, Spark > 100GB)
125- [ ] Columns pruned early (select only what is needed)
126- [ ] Filters pushed down to scan level (predicate pushdown)
127- [ ] Small tables broadcast in joins (< 100MB)
128- [ ] Spark AQE enabled (adaptive query execution)
129- [ ] No `collect()` on large datasets (aggregate before collecting)
130- [ ] Output files sized 128MB-1GB (coalesce/repartition before write)
131- [ ] Partition skew monitored and mitigated (salting or AQE)
132- [ ] Job profiled: Spark UI stages, Polars `.explain()`, DuckDB `EXPLAIN ANALYZE`
133- [ ] Memory sized appropriately: executor memory >= 2x largest partition