Every time PySpark hands data to Python, or Python hands data back, the bytes cross a process boundary. For years that meant pickling rows one by one, often the slowest part of a job. Apache Arrow turned the crossing into a per-column copy of buffers that both sides read directly.

This article explains Arrow from the bytes up: what a record batch looks like in memory, how batches are framed on the wire, which Spark operations produce and consume them, how types and timestamps are mapped, and how batch size turns into memory pressure on the driver and on Python workers. The cost of UDFs themselves, and the full table of Arrow configuration flags, are covered in PySpark performance; this page concentrates on the format and the data paths.

Advertisement

Why Spark needs a columnar interchange format

Inside an executor JVM, rows live in UnsafeRow, a compact binary layout managed by Tungsten (see Spark Tungsten). Python runs in separate worker processes. Without Arrow, each row is converted to Python objects, pickled, sent over a socket and unpickled, with interpreter overhead on every value.

Arrow defines a language-independent columnar memory format. A column of one million 64-bit integers is a single contiguous buffer of eight million bytes plus a validity bitmap. The JVM writes that buffer once and pyarrow wraps it without touching individual values. Work moves from per-value interpretation to bulk copies; the gain depends on types and widths, so measure it on your data.

Where Arrow record batches cross process boundaries in SparkExecutor JVMUnsafeRow / InternalRowExecutor JVMtask output partitionsPython workerpyarrow / pandas UDF codeDriver (PySpark)pyarrow.Table, pandasSpark Connect clientgRPC, any languageSpark Connect serverruns the planArrow IPC stream, per taskresult batches backtoPandas / toArrow batchesArrow batchesOne batch = schema-typed columns, at most maxRecordsPerBatch rowsDriver must hold the whole collected resultArrow replaces per-row pickling with column buffers that both sides read without parsing each value.Rows are converted to columns once on the JVM side; Python reads the buffers in place.
Arrow batches flow between executor JVMs and Python workers for UDFs, from executors to the driver for toPandas and toArrow, and from the Spark Connect server to clients.

The columnar format from first principles

An Arrow record batch is a schema plus a set of equal-length column arrays. Each array is made of a few buffers whose meaning depends on the type:

  • Validity bitmap. One bit per row: 1 means the value is present, 0 means null. If a column has no nulls the bitmap may be omitted entirely.
  • Data buffer for fixed-width types. Integers, floats, dates and timestamps are packed back to back, so row i of an int64 column is at byte offset 8 * i. Nulls still occupy a slot; the bitmap says to ignore it.
  • Offsets plus data for variable-width types. Strings and binary use an offsets buffer of n+1 integers and one contiguous data buffer. Value i is the bytes from offsets[i] to offsets[i+1].
  • Child arrays for nested types. A list column has its own offsets into a child array; a struct column is a set of child arrays sharing one validity bitmap at the struct level.

A worked example makes the layout concrete. Take four rows with columns id: int64 and city: string: (1, "Pune"), (2, null), (3, "Oslo"), (4, "Lima"). The id column is 32 bytes of data and no bitmap. The city column has validity 1, 0, 1, 1 for rows 0 to 3 (stored least-significant bit first as the byte 0b00001101), an offsets buffer [0, 4, 4, 8, 12] in which the null row has zero length, and a data buffer PuneOsloLima of 12 bytes. Nothing in that layout needs parsing to find row 3: two offset reads and a slice.

Two consequences matter for Spark users. First, standard string and binary columns use 32-bit offsets, so one column in one batch cannot hold more than about 2 GiB of character data. Arrow has large variants (LargeUtf8, LargeBinary) with 64-bit offsets, and Spark exposes spark.sql.execution.arrow.useLargeVarTypes to choose them; check your version's default before relying on either. Second, Arrow buffers are aligned and padded so that consumers can use vectorised instructions, which is part of why reading them in Python is cheap.

Advertisement

The IPC stream: how batches travel

Arrow defines an inter-process communication format for moving batches between processes. The streaming variant, which Spark uses on sockets, is a sequence of messages: one schema message first, then any number of record batch messages, then an end-of-stream marker. Each record batch message is a small FlatBuffers metadata header describing buffer lengths and null counts, followed by the raw buffers as a body. A reader allocates nothing per value; it records where each buffer starts.

Because the schema is sent once, per-batch overhead is small, and a reader can process batch k while batch k+1 is arriving. That pipelining keeps a Python worker busy, and it is why batch size, not partition size, determines how much memory one step needs.

import pyarrow as pa

batch = pa.record_batch(
    [pa.array([1, 2, 3, 4], pa.int64()),
     pa.array(["Pune", None, "Oslo", "Lima"], pa.string())],
    names=["id", "city"],
)
city = batch.column(1)
print(city.null_count)                      # 1
print([b.size if b else 0 for b in city.buffers()])
# validity, offsets, data -> sizes after padding, e.g. [1, 20, 12] (exact sizes vary by version)

sink = pa.BufferOutputStream()
with pa.ipc.new_stream(sink, batch.schema) as writer:   # schema message
    writer.write_batch(batch)                          # one record batch message
buf = sink.getvalue()                                  # then end-of-stream on close
print(pa.ipc.open_stream(buf).read_all().num_rows)     # 4

The four paths that use Arrow

Spark uses Arrow in four distinct places, and they fail in different ways, so it helps to name them separately.

PathDirectionWho holds memoryTypical failure
toPandas() / toArrow()executors to driverDriver: the whole result as Arrow, then as pandasDriver out of memory, or result size limit
createDataFrame(pandas or pyarrow)driver to executorsDriver: the local data plus its Arrow copySlow plans and large task payloads for big local data
Python UDFs over Arrowexecutor JVM to Python worker and backEach Python worker: one batch at a timeContainer killed for exceeding memory limits
Spark Connect resultsserver to remote clientClient: whatever it collectsClient-side out of memory on large collects

For toPandas(), each executor converts its partition into Arrow batches; the batches are sent to the driver; PySpark assembles a pyarrow.Table and then converts it to pandas. Spark 4.0 added DataFrame.toArrow(), which stops after the Table step and so avoids the pandas copy; the API documentation is explicit that it still collects everything into driver memory. For createDataFrame from pandas, the driver converts the local frame to Arrow, slices it into batches, and ships them as the data of a local relation.

For UDFs, the executor sends the input columns to the worker as an IPC stream and reads results back the same way. Pandas UDFs, mapInPandas, mapInArrow, Arrow-optimised Python UDFs and grouped applyInPandas all ride on this channel; they differ only in what object Python code sees. Spark Connect, finally, returns query results to clients as Arrow batches over gRPC, which is what lets non-Python clients consume results without a JVM; Spark Connect architecture covers the protocol.

Worked example: stay in Arrow end to end

Suppose a job enriches 200 million click events with a URL-normalisation step written in Python, then a downstream analyst wants a 50,000-row sample locally. The naive version uses a row UDF and toPandas() on an unbounded frame. The Arrow version keeps data columnar the whole way and bounds what reaches the driver:

import pyarrow as pa
import pyarrow.compute as pc
from pyspark.sql import SparkSession

spark = (SparkSession.builder
         .config("spark.sql.execution.arrow.pyspark.enabled", "true")          # default true from 4.2
         .config("spark.sql.execution.arrow.pyspark.fallback.enabled", "false")
         .config("spark.sql.execution.arrow.maxRecordsPerBatch", "5000")
         .getOrCreate())

clicks = spark.read.parquet("s3://bucket/clicks/")   # url: string, ts: timestamp, user_id: long

def normalise(batches):
    # Receives an iterator of pyarrow.RecordBatch; yields RecordBatches.
    for b in batches:
        url = pc.utf8_lower(b.column("url"))
        url = pc.replace_substring_regex(url, pattern="[?#].*$", replacement="")
        yield pa.RecordBatch.from_arrays(
            [b.column("user_id"), b.column("ts"), url],
            names=["user_id", "ts", "url"],
        )

out = clicks.mapInArrow(normalise, "user_id long, ts timestamp, url string")
out.write.mode("overwrite").parquet("s3://bucket/clicks_norm/")

sample = out.limit(50_000).toArrow()     # bounded, no pandas copy
print(sample.schema, sample.num_rows)

Three choices carry the benefit. mapInArrow lets pyarrow.compute run vectorised C++ kernels, so no Python runs per row. The batch size is lowered because URLs can be long. And limit before toArrow() means the driver holds 50,000 rows, not 200 million.

Type and timestamp mapping

Arrow types do not map one to one onto Spark SQL types or pandas dtypes, and most surprises come from the edges.

Spark SQL typeArrow typeWhat pandas sees
LongType, IntegerTypeint64, int32int64 without nulls; with nulls, float64 before 4.2 and nullable Int64 in pandas UDFs from 4.2
StringTypeutf8 (or large_utf8)object dtype of Python str
DecimalType(p, s)decimal128(p, s)object dtype of decimal.Decimal
TimestampTypetimestamp in microseconds, UTCdatetime64 nanoseconds, in the session time zone
ArrayType, StructType, MapTypelist, struct, mapobject columns of lists, dicts or tuples

The PySpark 4.2 Arrow guide states the timestamp rules precisely: data moving from Spark to pandas is converted to nanoseconds and to the session time zone set by spark.sql.session.timeZone; data moving to a PyArrow Table stays in microseconds in UTC; and data moving from pandas or PyArrow into Spark becomes UTC microseconds. If the session time zone is unset it defaults to the JVM's local zone, so the same notebook can produce different wall-clock values on two clusters. Set the session time zone explicitly in every job that converts timestamps.

The same guide lists an ArrayType of TimestampType as unsupported by Arrow conversion. Since 4.1, spark.sql.execution.pandas.convertToArrowArraySafely is on by default, so an integer overflow or lossy float conversion raises instead of silently corrupting values. Both changes alter behaviour on upgrade without any code change, which is why UDF-heavy jobs need regression tests pinned to the Spark version.

Batch size, the offset limit and memory

spark.sql.execution.arrow.maxRecordsPerBatch caps rows per batch, and the default is 10,000. Rows are the wrong unit for memory, so translate it. A batch of 10,000 rows with an average of 2 KB of string data per row is about 20 MB of Arrow buffers. A pandas UDF then converts that batch to pandas, often doubling it, and user code may make further copies. An executor with five cores runs up to five Python workers at once, so the Python side of that executor can need several hundred megabytes before any model or lookup table is loaded. That memory is outside the JVM heap and must fit in the container's overhead allowance; Spark memory management shows how the container limit is assembled.

Wide rows make the 32-bit offset limit real. If a single row carries a 5 MB JSON document, 10,000 such rows is 50 GB of string data in one column, far past the 2 GiB a standard string array can address. The fix is either a much smaller batch size for that job or large variable-width types. Prefer the smaller batch size first; it also bounds worker memory.

On the driver, toPandas() has a peak footprint of roughly the Arrow Table plus the pandas frame, because both exist during conversion. The spark.sql.execution.arrow.pyspark.selfDestruct.enabled option lets pyarrow release Arrow buffers as columns are converted, reducing that peak at some cost in speed. spark.driver.maxResultSize also applies to collected Arrow data, so a large collect can fail with a result-size error before it fails with an out-of-memory error.

Fallback, and how failures hide

When Arrow conversion for toPandas or createDataFrame fails early, for example on an unsupported type, spark.sql.execution.arrow.pyspark.fallback.enabled lets PySpark retry with the non-Arrow row path and log a warning. In a notebook that is convenient. In production it converts a clear type error into a job that is quietly ten times slower. Disable fallback in scheduled jobs and fix the type instead; the usual fixes are casting a problem column to a supported type before the conversion, or dropping it.

In the UDF path, a Python worker killed for exceeding its memory limit surfaces as a lost executor with no Python traceback. If UDF tasks die with exit codes and no stack trace, suspect batch size and worker memory first.

Failure modes

SymptomCauseFix
Driver out of memory on toPandasUnbounded result collected as Arrow then pandasWrite to storage; limit first; use toArrow to skip the pandas copy
Job suddenly much slowerSilent fallback to the row pathDisable fallback; read the warning; cast unsupported columns
Containers killed during pandas UDFsLarge batches times concurrent workers exceeds overheadLower maxRecordsPerBatch; raise memory overhead; avoid extra copies
Offset overflow or capacity errors on stringsMore than about 2 GiB of string data in one column of one batchSmaller batches; large var types where supported
Timestamps shifted by hoursSession time zone differs from expectationSet spark.sql.session.timeZone explicitly
UDF output changed after upgradeDefault flips: Arrow on by default, safe conversion, nullable Int64Pin and test versions; read the migration guide

Trade-offs

Arrow is not free: row-to-column conversion costs JVM CPU, and memory becomes lumpy because a whole batch is materialised at once.

The choice between mapInPandas and mapInArrow is the same trade-off in miniature. Pandas is familiar and rich, but each batch pays a conversion and pandas' own type quirks. Arrow with pyarrow.compute skips that conversion and keeps types exact, but its function library is narrower and its API less familiar. Use Arrow when the transformation maps onto compute kernels, and pandas when it needs pandas-specific logic. For pandas-on-Spark semantics, which sit on top of all of this, see the Spark pandas API.

What to do next

  1. Record your Spark and pyarrow versions; PySpark 4.2 requires pyarrow 18.0.0 or later and turns Arrow on by default for conversions and Python UDFs.
  2. Set spark.sql.session.timeZone explicitly in every job that moves timestamps between Spark and Python.
  3. Disable Arrow fallback in production jobs and fix any type it reports.
  4. Replace unbounded toPandas() calls with writes to storage, or with limit plus toArrow().
  5. Size maxRecordsPerBatch from bytes per row, not rows: aim for batches of tens of megabytes, times the number of cores per executor.
  6. Move per-row Python logic to mapInArrow with pyarrow.compute where kernels exist, and measure the change on a real partition.
  7. Add a regression test that runs your UDFs on a fixed input and compares outputs, and rerun it on every Spark upgrade.
Key takeaway: Arrow turns the JVM-Python crossing into a copy of column buffers: a validity bitmap, offsets and data per column, framed as a schema followed by record batches. Know which of the four paths you are on, because toPandas and toArrow load the driver, UDFs load each Python worker, and Connect loads the client. Size batches in bytes, set the session time zone, turn fallback off, and test UDF outputs on every upgrade because the defaults keep moving.