| name | fabric-pandas-perf-remediate |
| description | Troubleshoot and optimize pandas performance in Microsoft Fabric Spark notebooks. Use when diagnosing slow pandas operations, toPandas() out-of-memory errors, pandas API on Spark (pyspark.pandas) bottlenecks, DataFrame conversion failures, collect() memory issues, driver memory exhaustion, notebook cell timeouts, or when optimizing pandas workloads for Fabric capacity. Covers pandas vs Spark DataFrame conversion, memory profiling, broadcast joins, shuffle tuning, resource profiles, and Native Execution Engine integration. |
| license | Complete terms in LICENSE.txt |
Fabric Pandas Performance Troubleshooting
Diagnose and resolve pandas-related performance issues in Microsoft Fabric Spark notebooks, including memory exhaustion, slow conversions, and suboptimal pandas API on Spark usage.
When to Use This Skill
- Notebook cells hang or timeout during pandas operations
toPandas() fails with OutOfMemoryError or Java heap space errors
collect() crashes the driver node
- Pandas API on Spark (
pyspark.pandas / ps) runs slower than expected
- DataFrame conversion between Spark and pandas causes memory spikes
- Notebook kernel restarts unexpectedly during data processing
- Large dataset operations exhaust driver memory on Fabric capacity
- Need to choose between pandas, Spark DataFrame, or pandas API on Spark
Prerequisites
- Microsoft Fabric workspace with Data Engineering experience
- Fabric capacity F2 or higher (F64+ recommended for large datasets)
- PySpark notebook with Spark session active
- Basic familiarity with pandas and PySpark DataFrames
Quick Diagnosis
Symptom-to-Solution Map
Right-Size Your Approach
Critical Decision: Choose the right DataFrame API for your data size and workload.
Dataset Size Decision Tree:
─────────────────────────────────────────────────────────
< 100 MB → Native pandas (pd.DataFrame)
100 MB - 1 GB → pandas API on Spark (ps.DataFrame)
> 1 GB → PySpark DataFrame (spark.DataFrame)
> 10 GB → PySpark + partitioning + Delta optimization
Mixed workload? → Process in Spark, convert final aggregation to pandas
Visualization? → Aggregate in Spark first, toPandas() on summary only
ML feature eng? → Spark for transforms, pandas for final model input
API Comparison
| Operation | Native pandas | pandas API on Spark | PySpark DataFrame |
|---|
| Memory model | Single-node (driver) | Distributed | Distributed |
| Max practical size | ~2-4 GB | 10s-100s GB | TB+ |
| Startup overhead | None | Spark session | Spark session |
| groupby speed (small) | Fast | Slower (shuffle) | Slower (shuffle) |
| groupby speed (large) | OOM risk | Fast | Fast |
| Interop with Spark | .toPandas() | .to_spark() | Native |
toPandas Optimization
Problem
toPandas() collects the entire distributed DataFrame to the single driver node. This is the #1 cause of OOM in Fabric notebooks.
Solutions (Progressive)
1. Reduce data BEFORE conversion
pdf = spark_df.toPandas()
pdf = (spark_df
.filter("date >= '2024-01-01'")
.select("customer_id", "revenue", "region")
.toPandas())
2. Aggregate in Spark, convert summary
pdf = spark_df.toPandas()
result = pdf.groupby('region')['revenue'].sum()
summary = spark_df.groupBy("region").agg(F.sum("revenue").alias("total_revenue"))
pdf = summary.toPandas()
3. Enable Apache Arrow for faster conversion
spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "true")
spark.conf.set("spark.sql.execution.arrow.pyspark.fallback.enabled", "true")
pdf = spark_df.toPandas()
4. Use sampling for exploration
pdf = spark_df.sample(fraction=0.01, seed=42).toPandas()
pdf = spark_df.limit(100000).toPandas()
5. Chunk large conversions
def process_in_chunks(spark_df, chunk_col="date", process_fn=None):
"""Convert Spark DF to pandas in manageable chunks."""
chunks = [row[chunk_col] for row in spark_df.select(chunk_col).distinct().collect()]
results = []
for chunk_val in chunks:
chunk_pdf = spark_df.filter(F.col(chunk_col) == chunk_val).toPandas()
if process_fn:
chunk_pdf = process_fn(chunk_pdf)
results.append(chunk_pdf)
return pd.concat(results, ignore_index=True)
Driver Memory Tuning
Fabric Driver Memory by Node Size
| Node Size | vCores | Memory | Recommended Max toPandas() |
|---|
| Small | 4 | 32 GB | ~4-6 GB |
| Medium | 8 | 64 GB | ~10-12 GB |
| Large | 16 | 128 GB | ~20-25 GB |
| X-Large | 32 | 256 GB | ~40-50 GB |
Rule of thumb: toPandas() safe limit ≈ 15-20% of total driver memory (pandas creates copies during operations).
Configure Driver Memory
print(f"Driver memory: {spark.conf.get('spark.driver.memory', 'default')}")
%%configure
{
"driverMemory": "28g",
"driverCores": 8
}
Resource Profile Selection for Pandas Workloads
spark.conf.set("spark.fabric.resourceProfile", "readHeavyForSpark")
Shuffle Optimization
Tune for pandas API on Spark Operations
spark.conf.set("spark.sql.shuffle.partitions", "auto")
data_size_gb = 5
optimal_partitions = max(1, int(data_size_gb * 1024 / 128))
spark.conf.set("spark.sql.shuffle.partitions", str(optimal_partitions))
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
Broadcast Join Optimization
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "100m")
from pyspark.sql.functions import broadcast
result = large_df.join(broadcast(small_lookup_df), "key_col")
Enable Autotune
spark.conf.set("spark.ms.autotune.enabled", "true")
Arrow Conversion Fixes
Common Errors and Solutions
| Error | Cause | Fix |
|---|
ArrowInvalid: Could not convert X | Unsupported type | Cast column before conversion |
ArrowNotImplementedError | Nested types | Flatten struct/array columns |
pyarrow.lib.ArrowMemoryError | OOM during Arrow transfer | Reduce data size or increase memory |
| Null handling mismatch | Pandas NaN vs Spark null | Use spark.sql.execution.arrow.pyspark.fallback.enabled |
from pyspark.sql.types import StringType, DoubleType
spark_df = spark_df.withColumn("mixed_col", F.col("mixed_col").cast(StringType()))
spark_df = spark_df.select(
"simple_col",
F.col("struct_col.field1").alias("field1"),
F.col("struct_col.field2").alias("field2")
)
spark_df = spark_df.fillna({"numeric_col": 0, "string_col": ""})
Incremental Processing
Pattern: Spark Processing with Pandas Finish
import pyspark.sql.functions as F
import pandas as pd
aggregated = (spark_df
.filter(F.col("status") == "active")
.groupBy("category", "month")
.agg(
F.sum("amount").alias("total"),
F.count("*").alias("cnt"),
F.avg("score").alias("avg_score")
))
row_count = aggregated.count()
print(f"Rows to convert: {row_count:,}")
assert row_count < 1_000_000, f"Too many rows ({row_count:,}) for toPandas()"
pdf = aggregated.toPandas()
pivot = pdf.pivot_table(index='category', columns='month', values='total')
Memory Profiling
Monitor Driver Memory in Notebook
import os, psutil
def check_memory():
"""Report current driver memory usage."""
process = psutil.Process(os.getpid())
mem_info = process.memory_info()
print(f"RSS Memory: {mem_info.rss / 1024**3:.2f} GB")
print(f"VMS Memory: {mem_info.vms / 1024**3:.2f} GB")
sys_mem = psutil.virtual_memory()
print(f"System Used: {sys_mem.used / 1024**3:.2f} / {sys_mem.total / 1024**3:.2f} GB ({sys_mem.percent}%)")
check_memory()
pdf = spark_df.toPandas()
check_memory()
Reduce pandas Memory Footprint
def optimize_pandas_dtypes(df):
"""Downcast pandas DataFrame dtypes to reduce memory."""
for col in df.select_dtypes(include=['int64']).columns:
df[col] = pd.to_numeric(df[col], downcast='integer')
for col in df.select_dtypes(include=['float64']).columns:
df[col] = pd.to_numeric(df[col], downcast='float')
for col in df.select_dtypes(include=['object']).columns:
if df[col].nunique() / len(df) < 0.5:
df[col] = df[col].astype('category')
return df
pdf = optimize_pandas_dtypes(pdf)
print(f"Memory after optimization: {pdf.memory_usage(deep=True).sum() / 1024**2:.1f} MB")
pandas API on Spark Best Practices
Use pyspark.pandas Instead of Conversion
import pyspark.pandas as ps
psdf = ps.read_delta("Tables/my_table")
psdf = spark_df.pandas_api()
result = psdf.groupby("region")["revenue"].sum()
pdf = result.to_pandas()
Common Pitfalls
for idx, row in psdf.iterrows():
process(row)
psdf["new_col"] = psdf["col_a"] * psdf["col_b"]
psdf["result"] = psdf["col"].apply(lambda x: complex_fn(x))
from pyspark.sql.functions import pandas_udf
@pandas_udf("double")
def optimized_fn(series: pd.Series) -> pd.Series:
return series * 2 + 1
Troubleshooting Checklist
- Check data size before any
toPandas() / collect() call
- Enable Arrow transfer:
spark.sql.execution.arrow.pyspark.enabled = true
- Filter/aggregate in Spark before converting to pandas
- Match node size to workload (Medium 64 GB minimum for pandas-heavy notebooks)
- Use pandas API on Spark for distributed pandas-like operations
- Monitor memory with
psutil before/after conversions
- Set resource profile to
readHeavyForSpark for notebook-heavy workloads
- Enable autotune for automatic shuffle and partition optimization
- Downcast dtypes after conversion to reduce pandas memory footprint
- Chunk processing for datasets that exceed single-node memory
Automation & Diagnostics
Run the diagnostic script to collect Spark session configuration, memory settings, and environment details for troubleshooting.
See the detailed reference guide for advanced patterns including pandas UDFs, Koalas migration, Native Execution Engine integration, and capacity planning formulas.
Use the notebook template as a starting point for memory-safe pandas workflows in Fabric notebooks.