Accepted answer
toPandas() collects everything to the driver as Arrow batches. If Arrow-based collection isn't enabled (spark.sql.execution.arrow.pyspark.enabled), it falls back to a much less memory-efficient row-by-row path.
Calling df.toPandas() on what I estimated as a 2GB DataFrame kills the driver with an OOM. Driver has 8GB.
pdf = df.toPandas()The estimate was rows times avg_row_size, but that's clearly wrong somewhere.
Accepted answer
toPandas() collects everything to the driver as Arrow batches. If Arrow-based collection isn't enabled (spark.sql.execution.arrow.pyspark.enabled), it falls back to a much less memory-efficient row-by-row path.
Document the grain decision, most BI bugs turn out to be grain bugs.
In our case the root cause was an implicit cast preventing pushdown.
Another path: push the compute to the warehouse if the data's already there.
In our case the root cause was an implicit cast preventing pushdown.
Document the grain decision, most BI bugs turn out to be grain bugs.
Check whether AQE is disabled in your Spark conf, skew join handling helped us a lot here.
Sign in to reply.
© 2026 Lakebench, operated by Hunnurji Rao. Bengaluru, Karnataka, India.
No cluster. No install. Just the tab.