Arrow Data Transfer

Apache Arrow and Columnar Data Transfer

Apache Arrow 129 's columnar format (Apache Arrow) moves data between the JVM and Python without per-value conversion: the executor fills Arrow buffers a batch at a time, and Python reads them as pandas 16,086 or PyArrow 129 arrays. Spark 129 uses Arrow in toPandas() and createDataFrame(pandas_df) on the driver, pandas and Arrow UDFs, and, since Spark 4.2 by default, ordinary Python UDFs (spark.sql.execution.pythonUDF.arrow.enabled), which appear in plans as ArrowEvalPython.

How rows reach a Python UDF: pickled rows versus Arrow record batches
How rows reach a Python UDF: pickled rows versus Arrow record batches
Collecting 1.38 million order lines to pandas without and with ArrowJavaScript
import time
sample = lines.select("order_id", "book_id", "qty", F.col("unit_price").cast("double"))
for arrow in ("false", "true"):
    spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", arrow)
    t0 = time.perf_counter()
    pdf = sample.toPandas()                           # 1.38 million rows into the driver
    print(f"arrow {arrow:<5} {len(pdf):,} rows in {time.perf_counter() - t0:5.2f}s")
Output
arrow false 1,380,694 rows in 21.18s
arrow true  1,380,694 rows in  1.80s

Arrow made the transfer about twelve times faster on this shared 4-CPU host. It is also stricter about types: a classic UDF declared "double" that returned a Python int produced NULL for every row, while the Arrow version converted the value, and a string returned for "int" raised an error instead of nulls.