| --- |
| layout: global |
| title: Apache Arrow Cache Format |
| displayTitle: Apache Arrow Cache Format |
| license: | |
| Licensed to the Apache Software Foundation (ASF) under one or more |
| contributor license agreements. See the NOTICE file distributed with |
| this work for additional information regarding copyright ownership. |
| The ASF licenses this file to You under the Apache License, Version 2.0 |
| (the "License"); you may not use this file except in compliance with |
| the License. You may obtain a copy of the License at |
| |
| http://www.apache.org/licenses/LICENSE-2.0 |
| |
| Unless required by applicable law or agreed to in writing, software |
| distributed under the License is distributed on an "AS IS" BASIS, |
| WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
| See the License for the specific language governing permissions and |
| limitations under the License. |
| --- |
| |
| ## Overview |
| |
| Apache Spark supports using Apache Arrow as an alternative cache format for in-memory Dataset caching. This format provides improved performance for certain workloads, especially when working with columnar data sources like Parquet and ORC. |
| |
| ## Benefits |
| |
| The Arrow cache format offers several advantages over the default cache format: |
| |
| - **Zero-copy reads** when input is already in Arrow format (e.g., Arrow-based data sources, re-caching Arrow cached data) |
| - **Better filter pushdown** with min/max statistics for partition pruning |
| - **Compact columnar layout** with zstd compression support |
| |
| Note that the cached bytes are an internal format (a schema-less Arrow RecordBatch payload), not |
| a complete Arrow IPC stream, so they are not directly readable by external Arrow tooling. |
| |
| **Note**: Spark's built-in Parquet/ORC readers use internal column vectors (`OnHeapColumnVector`/`OffHeapColumnVector`), not Arrow format, so they don't benefit from zero-copy optimization. |
| |
| ## Configuration |
| |
| `spark.sql.cache.serializer` is a static SQL configuration, so it must be set when the |
| SparkSession is built and cannot be changed on a running session (`spark.conf.set` rejects static |
| keys with `CANNOT_MODIFY_CONFIG`): |
| |
| ```scala |
| val spark = SparkSession.builder() |
| .appName("MyApp") |
| .config("spark.sql.cache.serializer", |
| "org.apache.spark.sql.execution.columnar.ArrowCachedBatchSerializer") |
| .getOrCreate() |
| ``` |
| |
| **Note**: This config selects the cache serializer for the whole session; once set, this |
| serializer handles every cached relation. There is no automatic per-relation fallback to another |
| cache serializer based on the data types involved (see |
| [Supported Data Types](#supported-data-types) for how unsupported types are handled). The chosen |
| serializer is also cached process-wide on first use, so switching cache formats within a JVM that |
| has already materialized a cache requires a fresh JVM (see |
| [Migration from Default Cache](#migration-from-default-cache)). |
| |
| ## Usage |
| |
| Once configured, use cache operations as normal: |
| |
| ```scala |
| // Cache a DataFrame |
| val df = spark.read.parquet("data.parquet") |
| df.cache() |
| |
| // Use cached data |
| df.filter("age > 30").count() |
| |
| // Uncache when done |
| df.unpersist() |
| ``` |
| |
| ## Compression |
| |
| Arrow cache supports multiple compression codecs. Configure compression with: |
| |
| ```scala |
| spark.conf.set("spark.sql.execution.arrow.compression.codec", "zstd") |
| ``` |
| |
| Available options: |
| - `none` - No compression (fastest, largest size, **default**) |
| - `zstd` - Zstandard compression (best compression; tunable level) |
| - `lz4` - LZ4 compression. Not recommended: Arrow's Java LZ4 codec is implemented with the |
| pure-Java Commons Compress framed LZ4 streams and is much slower than zstd |
| |
| For zstd, you can also configure the compression level. Positive values (up to 22) give better |
| compression but slower speed; negative values give ultra-fast compression with lower ratios: |
| |
| ```scala |
| spark.conf.set("spark.sql.execution.arrow.compression.zstd.level", "3") // Default: 3 |
| ``` |
| |
| ## Vectorized Reader |
| |
| Enable vectorized reading for better performance with primitive types: |
| |
| ```scala |
| spark.conf.set("spark.sql.inMemoryColumnarStorage.enableVectorizedReader", "true") |
| ``` |
| |
| When enabled, cached data is read as columnar batches instead of rows, which can significantly improve performance for columnar operations. |
| |
| ## Performance Characteristics |
| |
| In our benchmarks, the Arrow cache format performs best on the following workloads. Actual |
| results depend on data types, compression settings, and hardware, and the default cache format |
| can be faster in some cases (for example, with higher compression levels): |
| |
| 1. **Filter-Heavy Workloads**: Queries with selective filters benefit from min/max statistics. |
| 2. **Columnar Operations**: Aggregations and projections on cached data benefit from the Arrow format. |
| 3. **Parquet/ORC Caching**: Arrow's batch processing helps even without the zero-copy path. |
| 4. **Re-caching with Column Projection**: Dropping columns from Arrow-cached data preserves the |
| `ArrowColumnVector` format, enabling true zero-copy extraction and the largest gains. |
| |
| ### Benchmark Results |
| |
| The numbers below are illustrative results from one run on an Apple M4 Max (OpenJDK 21.0.8) and |
| will vary with hardware, JDK, and compression settings. They are not a guarantee. For the |
| authoritative, regularly regenerated numbers, see |
| `sql/core/benchmarks/ArrowCacheBenchmark-jdk21-results.txt` and the `ArrowCacheBenchmark` suite. |
| |
| | Workload | Default Cache | Arrow Cache | Speedup | |
| |----------|--------------|-------------|---------| |
| | Write + Read (5M rows, 3 primitive columns) | 153.7 ns/row | 74.2 ns/row | **~2X faster** | |
| | Cache then filter (5M rows) | 100.1 ns/row | 70.8 ns/row | **~1.4X faster** | |
| | Columnar input from Parquet (2M rows, 3 primitive columns) | 195.3 ns/row | 113.1 ns/row | **~1.7X faster** | |
| | Re-cache with zero-copy (2M rows, 2 columns) | 123.3 ns/row | 38.5 ns/row | **~3.2X faster** | |
| |
| **Notes**: |
| - **Write + Read**: Significant improvement from efficient Arrow serialization and vectorized operations |
| - **Cache then filter**: This measures end-to-end cache build plus a filtered scan, comparing the two cache formats. Both formats collect min/max statistics and can prune batches, so the difference reflects overall cache+scan throughput rather than pruning unique to Arrow |
| - **Parquet caching**: Shows improvement despite Spark's Parquet reader producing `OnHeapColumnVector`/`OffHeapColumnVector` rather than `ArrowColumnVector`, due to Arrow's efficient batch processing |
| - **Re-cache with zero-copy**: When caching a subset of columns from Arrow-cached data (e.g., `df.drop("column")`), the remaining columns preserve their `ArrowColumnVector` format, enabling true zero-copy extraction and achieving the best performance |
| - **Zero-copy benefits** only apply when input is already `ArrowColumnVector` (e.g., Python Arrow sources, re-caching Arrow cached data with column projection) |
| |
| ## Supported Data Types |
| |
| Arrow cache supports the following data types: |
| |
| ### Primitive Types |
| - BooleanType |
| - ByteType, ShortType, IntegerType, LongType |
| - FloatType, DoubleType |
| - DecimalType (all precision/scale combinations) |
| - NullType |
| |
| ### Temporal Types |
| - DateType |
| - TimestampType |
| - TimestampNTZType |
| - Nanosecond-precision timestamps (`TIMESTAMP_NTZ(p)` / `TIMESTAMP(p)` with `p` in 7..9) |
| - TimeType |
| |
| ### Interval Types |
| - YearMonthIntervalType |
| - DayTimeIntervalType |
| - CalendarIntervalType |
| |
| Nanosecond-precision timestamps and `CalendarIntervalType` are stored in lossless internal |
| Arrow representations (structs of the types' own components) rather than the standard Arrow |
| interchange encodings, so their full Spark value domains round-trip through the cache -- the |
| same domains the default cache serializer supports. The standard interchange encodings pack |
| these values into a single int64 of nanoseconds and therefore cover only a reduced range |
| (roughly years 1677-2262 for timestamps, +/-292 years of microseconds for intervals); that |
| limitation applies to Arrow interchange paths such as `toPandas`, not to this cache. |
| |
| ### String and Binary |
| - StringType (including collated strings) |
| - BinaryType |
| |
| ### Complex Types |
| - ArrayType |
| - StructType |
| - MapType |
| - Nested combinations of the above |
| |
| ### Other Types |
| - VariantType |
| - GeometryType, GeographyType |
| - User-defined types (UDTs) whose underlying representation is itself supported |
| |
| ### Unsupported Types |
| |
| Arrow cache covers every type the default cache serializer supports, plus some it |
| does not (for example geometry and geography). Types that Arrow cannot represent |
| (such as `ObjectType`) are not silently dropped or routed to a different cache |
| serializer: there is no per-type fallback, because the cache serializer is chosen |
| once via the static `spark.sql.cache.serializer` configuration and then handles |
| every cached relation. Attempting to cache an unsupported type fails with an |
| `UNSUPPORTED_DATATYPE` error when the cache is materialized. |
| |
| ## Statistics and Filter Pushdown |
| |
| Arrow cache automatically collects min/max statistics for the following types: |
| - Boolean |
| - Numeric types (Byte, Short, Int, Long, Float, Double) |
| - Decimal |
| - Date, Timestamp, and Timestamp without time zone (TIMESTAMP_NTZ) |
| - Nanosecond-precision timestamps |
| - Time |
| - Year-month and day-time intervals |
| - String (using collation-aware comparison for collated strings) |
| |
| Other types (Binary, Variant, calendar intervals, and complex types such as |
| Array/Struct/Map) are cached but do not contribute min/max bounds, so they only |
| record null counts and sizes. |
| |
| These statistics enable partition pruning when filtering: |
| |
| ```scala |
| val df = spark.range(10000000).cache() |
| |
| // This filter can skip batches using min/max statistics |
| df.filter("id > 5000000").count() |
| ``` |
| |
| ## Memory Management |
| |
| The cached data itself lives on the JVM heap, not off-heap. Each cached batch is stored as a |
| serialized Arrow IPC byte array (`Array[Byte]`), and the default `Dataset.cache()` storage level is |
| the deserialized `MEMORY_AND_DISK`, so those bytes are retained as ordinary heap objects (and spill |
| to disk under memory pressure). Arrow's off-heap allocators are used only for the transient |
| `VectorSchemaRoot`s created while encoding a batch for caching and while decoding a batch on read; |
| these are released as soon as the encode/decode completes. |
| |
| **Sizing implications**: |
| - Size the JVM heap (executor memory) for the cached data, since that is where it resides. This is |
| the main knob for cache capacity. |
| - `spark.executor.memoryOverhead` covers the transient off-heap encode/decode buffers, which are |
| proportional to a single batch, not to the total cached size. It generally does not need to grow |
| with the size of the cache. |
| - Arrow cache is often **more memory-efficient** than the default cache for the heap-resident bytes: |
| efficient zstd compression, a compact columnar layout without per-value Java object overhead, |
| and better compression ratios for strings and complex types. |
| |
| **Memory Cleanup**: |
| - The transient off-heap encode/decode roots are released when each task completes. |
| - The heap-resident cached bytes are released when the DataFrame is unpersisted or evicted, like any |
| other cached block. |
| |
| You can monitor cache block sizes through the Storage tab in the Spark UI. |
| |
| ## Limitations and Considerations |
| |
| 1. **Static Configuration**: Cache serializer must be set before SparkSession creation |
| 2. **Memory Overhead**: Arrow format has a small per-batch overhead |
| 3. **Compatibility**: Cannot mix cache formats - recache needed when switching |
| 4. **Compression Trade-off**: Higher compression = lower memory but slower reads |
| |
| ## Migration from Default Cache |
| |
| The cache serializer is resolved from `spark.sql.cache.serializer` only on first use and is then |
| held in a process-wide field that is not reset when a SparkSession stops. As a result, **switching |
| cache formats requires a fresh JVM** once any cache has been materialized -- stopping and |
| rebuilding the SparkSession in the same process keeps using the originally resolved serializer. |
| |
| To migrate from the default cache to Arrow cache: |
| |
| 1. **Start a new JVM / driver process** (a brand-new Spark application). |
| 2. **Build the SparkSession with the Arrow serializer**: |
| ```scala |
| val spark = SparkSession.builder() |
| .config("spark.sql.cache.serializer", |
| "org.apache.spark.sql.execution.columnar.ArrowCachedBatchSerializer") |
| .getOrCreate() |
| ``` |
| 3. **Cache your DataFrames** as usual. |
| |
| **Note**: Cache data is never shared across formats; each application caches in whichever format |
| its serializer produces. |
| |
| ## Troubleshooting |
| |
| ### Out of Memory Errors |
| |
| If you encounter OOM errors with Arrow cache: |
| |
| 1. Reduce batch size: |
| ```scala |
| spark.conf.set("spark.sql.execution.arrow.maxRecordsPerBatch", "5000") // Default: 10000 |
| ``` |
| |
| 2. Enable compression: |
| ```scala |
| spark.conf.set("spark.sql.execution.arrow.compression.codec", "zstd") |
| ``` |
| |
| 3. Reduce compression level: |
| ```scala |
| spark.conf.set("spark.sql.execution.arrow.compression.zstd.level", "1") |
| ``` |
| |
| ### Slow Performance |
| |
| If Arrow cache is slower than expected: |
| |
| 1. Enable vectorized reader: |
| ```scala |
| spark.conf.set("spark.sql.inMemoryColumnarStorage.enableVectorizedReader", "true") |
| ``` |
| |
| 2. Reduce or disable compression (decompression is part of every read): |
| ```scala |
| spark.conf.set("spark.sql.execution.arrow.compression.zstd.level", "1") // faster, less ratio |
| // or, for read-heavy workloads where memory is not the constraint: |
| spark.conf.set("spark.sql.execution.arrow.compression.codec", "none") |
| ``` |
| Note: `lz4` is not recommended. Arrow's Java LZ4 codec always uses the pure-Java Commons |
| Compress framed LZ4 implementation, which is far slower than zstd. |
| |
| 3. Increase batch size (if memory allows): |
| ```scala |
| spark.conf.set("spark.sql.execution.arrow.maxRecordsPerBatch", "20000") |
| ``` |
| |
| ## Configuration Reference |
| |
| | Configuration | Default | Description | |
| |---------------|---------|-------------| |
| | `spark.sql.cache.serializer` | DefaultCachedBatchSerializer | Cache format serializer class | |
| | `spark.sql.execution.arrow.compression.codec` | `none` | Compression codec (none, lz4, zstd) | |
| | `spark.sql.execution.arrow.compression.zstd.level` | `3` | Zstd compression level (negative = faster, up to 22) | |
| | `spark.sql.execution.arrow.maxRecordsPerBatch` | `10000` | Maximum rows per Arrow batch | |
| | `spark.sql.execution.arrow.maxBytesPerBatch` | `64MB` | Maximum bytes per Arrow batch (whichever limit is hit first applies) | |
| | `spark.sql.execution.arrow.cache.prefetch.enabled` | `false` | Prefetch the next batch in the background while the current one is consumed | |
| | `spark.sql.inMemoryColumnarStorage.enableVectorizedReader` | `true` | Enable vectorized cache reading | |
| |
| ## Example: Complete Application |
| |
| ```scala |
| import org.apache.spark.sql.SparkSession |
| |
| object ArrowCacheExample { |
| def main(args: Array[String]): Unit = { |
| // Create SparkSession with Arrow cache |
| val spark = SparkSession.builder() |
| .appName("ArrowCacheExample") |
| .master("local[*]") |
| .config("spark.sql.cache.serializer", |
| "org.apache.spark.sql.execution.columnar.ArrowCachedBatchSerializer") |
| .config("spark.sql.execution.arrow.compression.codec", "zstd") |
| .config("spark.sql.inMemoryColumnarStorage.enableVectorizedReader", "true") |
| .getOrCreate() |
| |
| try { |
| // Read columnar data source |
| val df = spark.read.parquet("large_dataset.parquet") |
| |
| // Cache with Arrow format |
| df.cache() |
| |
| // Queries benefit from zero-copy reads and statistics |
| val result1 = df.filter("age > 30").select("name", "age").count() |
| println(s"Filtered count: $result1") |
| |
| val result2 = df.groupBy("country").agg(sum("sales")).collect() |
| println(s"Aggregation result: ${result2.mkString(", ")}") |
| |
| // Uncache when done |
| df.unpersist() |
| |
| } finally { |
| spark.stop() |
| } |
| } |
| } |
| ``` |
| |
| ## Best Practices |
| |
| 1. **Use with Columnar Sources**: Maximum benefit with Parquet/ORC |
| 2. **Enable Statistics**: Let Arrow cache collect min/max for filter pushdown |
| 3. **Monitor Memory**: Watch cache block sizes in the Spark UI Storage tab; the cached data lives on the JVM heap |
| 4. **Test First**: Benchmark your workload before production deployment |
| 5. **Compression**: Start with `none` for read-heavy workloads, or `zstd` (default level 3) when memory matters; avoid `lz4` (Arrow's Java LZ4 codec is pure-Java Commons Compress and much slower than zstd) |
| 6. **Vectorization**: Enable vectorized reader for primitive-heavy workloads |
| |
| ## Further Reading |
| |
| - [Apache Arrow Project](https://arrow.apache.org/) |
| - [Spark Caching Documentation](https://spark.apache.org/docs/latest/sql-performance-tuning.html#caching-data-in-memory) |
| - [Arrow IPC Format](https://arrow.apache.org/docs/format/Columnar.html#ipc-file-format) |