Columnar Data in 2026: Apache Arrow Memory Layout, Parquet Storage & DuckDB Vectorization

Table of Contents(13 sections)
Columnar data formats are foundational to modern analytical systems, enabling sub-second query performance over petabyte-scale datasets. This guide dissects the interplay between Apache Parquet for persistent storage, Apache Arrow for in-memory processing, and DuckDB's vectorized execution engine, illustrating how these technologies collectively form a high-performance analytical data stack.
Apache Parquet: Disk-Optimized Columnar Storage
Apache Parquet is the de facto standard for columnar data storage in big data ecosystems. Its design prioritizes efficient disk I/O and predicate pushdown, crucial for analytical workloads.
Parquet File Structure
A Parquet file is composed of row groups, which are horizontal partitions of the data. Within each row group, data is stored column-wise. This structure allows for reading only the necessary columns for a query, significantly reducing I/O.
Parquet File
├── File Metadata (footer)
│ └── Schema, Row Group Metadata, Column Metadata
├── Row Group 1
│ ├── Column 1 Chunk
│ │ ├── Page 1 (Data, Dictionary, Index)
│ │ ├── Page 2
│ │ └── ...
│ ├── Column 2 Chunk
│ └── ...
├── Row Group 2
│ ├── Column 1 Chunk
│ └── ...
└── ...
Encoding and Compression
Parquet employs various encoding schemes and compression algorithms to minimize storage footprint and improve read performance.
- Dictionary Encoding: For columns with low cardinality, unique values are stored in a dictionary, and the column data stores indices into this dictionary. This is highly effective for string or categorical data.
- Run Length Encoding (RLE): Efficient for sequences of repeated values.
- Bit-Packed Encoding: For integers, packs values into the minimum number of bits required.
- Compression: Snappy, Gzip, Zstd, Brotli. Zstd offers an excellent balance of compression ratio and speed.
Predicate Pushdown and Column Statistics
Parquet's file format includes statistics (min/max values, null counts) for each column chunk within a row group. Query engines leverage these statistics to prune entire row groups or pages that cannot possibly satisfy a query's predicates, avoiding unnecessary I/O.
Consider a query SELECT sum(value) FROM table WHERE timestamp > '2023-01-01'. If a row group's timestamp column statistics indicate its maximum timestamp is 2022-12-31, that entire row group is skipped.
import pandas as pd
import pyarrow as pa
import pyarrow.parquet as pq
import numpy as np
import os
# Generate sample data
data = {
'id': np.arange(100000),
'timestamp': pd.to_datetime(pd.date_range('2022-01-01', periods=100000, freq='H')),
'category': np.random.choice(['A', 'B', 'C', 'D'], 100000),
'value': np.random.rand(100000) * 100
}
df = pd.DataFrame(data)
# Convert to PyArrow Table
table = pa.Table.from_pandas(df)
# Define Parquet file path
parquet_file = 'sample_data.parquet'
# Write to Parquet with specific options
# row_group_size: Controls the number of rows in each row group.
# Smaller row groups can improve pruning but increase metadata overhead.
# Larger row groups reduce metadata but might read more data than needed.
# compression: Zstd offers good balance.
# write_statistics: Essential for predicate pushdown.
pq.write_table(
table,
parquet_file,
row_group_size=10000, # 10,000 rows per row group
compression='zstd',
write_statistics=True,
data_page_size=8 * 1024 * 1024 # 8MB data page size
)
print(f"Parquet file '{parquet_file}' created.")
# Inspect Parquet file metadata
parquet_schema = pq.read_schema(parquet_file)
print("\nParquet Schema:")
print(parquet_schema)
parquet_metadata = pq.read_metadata(parquet_file)
print("\nParquet Metadata:")
print(parquet_metadata)
# Example of reading with predicate pushdown
# This query will only read row groups where 'timestamp' is after '2023-01-01'
# and only the 'timestamp' and 'value' columns.
start_time = pd.Timestamp('2023-01-01')
filtered_table = pq.read_table(
parquet_file,
columns=['timestamp', 'value'],
filters=[('timestamp', '>', start_time)]
)
print(f"\nFiltered table (first 5 rows) after predicate pushdown and column projection:")
print(filtered_table.to_pandas().head())
# Clean up
os.remove(parquet_file)
Apache Arrow: In-Memory Columnar Format
Apache Arrow defines a language-agnostic, contiguous memory layout for columnar data. Its primary goal is to eliminate serialization/deserialization overhead when data moves between different systems or processes.
Arrow Memory Layout
Arrow's strength lies in its standardized, flat, and contiguous memory buffers. For a given column, all values of the same type are stored together in a single buffer. This layout is highly cache-friendly and enables efficient SIMD (Single Instruction, Multiple Data) operations.
For a simple integer column:
[value1, value2, value3, ..., valueN]
For a nullable integer column:
[validity_bitmap_buffer], [value1, value2, value3, ..., valueN]
The validity bitmap indicates which values are null.
For variable-length types (e.g., strings):
[offset_buffer], [data_buffer], [validity_bitmap_buffer]
The offset buffer stores the starting position and length of each string within the data buffer.
import pyarrow as pa
import numpy as np
# Example: Integer Array
int_array = pa.array([1, 2, None, 4, 5], type=pa.int32())
print("Integer Array:")
print(int_array)
print(f" Type: {int_array.type}")
print(f" Buffers: {[b.to_pybytes() for b in int_array.buffers()]}") # Validity bitmap, Data buffer
print(f" Is Null: {int_array.is_null()}")
print(f" Values: {int_array.to_numpy()}")
# Example: String Array
string_array = pa.array(["apple", "banana", None, "cherry"], type=pa.string())
print("\nString Array:")
print(string_array)
print(f" Type: {string_array.type}")
print(f" Buffers: {[b.to_pybytes() for b in string_array.buffers()]}") # Validity bitmap, Offset buffer, Data buffer
print(f" Is Null: {string_array.is_null()}")
print(f" Values: {string_array.to_numpy()}")
# Example: RecordBatch (collection of arrays/columns)
schema = pa.schema([
pa.field('id', pa.int64()),
pa.field('name', pa.string()),
pa.field('score', pa.float32())
])
record_batch = pa.RecordBatch.from_arrays(
[
pa.array([1, 2, 3, 4]),
pa.array(["Alice", "Bob", "Charlie", "David"]),
pa.array([90.5, 88.0, 92.1, 78.9])
],
schema=schema
)
print("\nRecordBatch:")
print(record_batch)
print(f" Number of columns: {record_batch.num_columns}")
print(f" Number of rows: {record_batch.num_rows}")
print(f" Column 'name' data buffer (first 10 bytes): {record_batch.column('name').buffers()[2].to_pybytes()[:10]}")
Zero-Copy Reads
The most significant advantage of Arrow is its zero-copy capability. When data is already in Arrow format, it can be passed between processes, languages, or even different machines (via Arrow Flight) without serialization or deserialization. This drastically reduces CPU cycles and latency.
For instance, a Parquet file can be read into an Arrow Table or RecordBatch directly. The Parquet reader maps the on-disk columnar data to Arrow's in-memory layout, often with minimal data transformation.
DuckDB: Vectorized Query Engine
DuckDB is an in-process SQL OLAP database designed for analytical workloads. Its core strength lies in its vectorized query execution engine, which operates directly on columnar data, often in Arrow format, leveraging SIMD instructions for extreme performance.
Vectorized Execution
Instead of processing one row at a time (tuple-at-a-time processing), DuckDB processes data in batches (vectors) of typically 1024-4096 rows. Each operation (e.g., filter, join, aggregate) works on entire vectors, allowing for:
- Reduced Function Call Overhead: A single function call processes many values.
- CPU Cache Efficiency: Data is accessed sequentially, leading to high cache hit rates.
- SIMD Instruction Utilization: Modern CPUs can perform the same operation on multiple data points simultaneously using SIMD instructions (e.g., AVX2, AVX-512). DuckDB's columnar layout is perfectly suited for this.
Integration with Parquet and Arrow
DuckDB can directly query Parquet files without requiring an explicit CREATE TABLE statement or data ingestion. It reads Parquet files, converts them into its internal columnar representation (which is Arrow-compatible), and then executes queries using its vectorized engine. This eliminates the ETL step for many analytical tasks.
import duckdb
import pyarrow as pa
import pyarrow.parquet as pq
import pandas as pd
import numpy as np
import os
import time
# 1. Create a large Parquet file for demonstration
num_rows = 10_000_000
data = {
'id': np.arange(num_rows),
'timestamp': pd.to_datetime(pd.date_range('2020-01-01', periods=num_rows, freq='10s')),
'category': np.random.choice(['A', 'B', 'C', 'D', 'E'], num_rows),
'value': np.random.rand(num_rows) * 1000,
'is_active': np.random.choice([True, False], num_rows, p=[0.7, 0.3])
}
df_large = pd.DataFrame(data)
table_large = pa.Table.from_pandas(df_large)
large_parquet_file = 'large_sample_data.parquet'
pq.write_table(table_large, large_parquet_file, row_group_size=100000, compression='zstd', write_statistics=True)
print(f"Created large Parquet file: {large_parquet_file} with {num_rows} rows.")
# 2. Initialize DuckDB
con = duckdb.connect(database=':memory:', read_only=False)
# 3. Query Parquet directly
print("\n--- Querying Parquet file directly with DuckDB ---")
query = f"""
SELECT
category,
COUNT(*) AS count,
AVG(value) AS avg_value,
MAX(timestamp) AS max_timestamp
FROM '{large_parquet_file}'
WHERE timestamp >= '2023-01-01' AND is_active = TRUE
GROUP BY category
ORDER BY count DESC;
"""
start_time = time.time()
result_df = con.execute(query).fetchdf()
end_time = time.time()
print(f"Query executed in {end_time - start_time:.4f} seconds.")
print(result_df)
# 4. Demonstrate Arrow integration: Create an Arrow Table and query it
print("\n--- Querying an in-memory PyArrow Table with DuckDB ---")
# Create a smaller Arrow table for in-memory demo
arrow_data = pa.table({
'city': pa.array(['NYC', 'LA', 'Chicago', 'NYC', 'LA']),
'population': pa.array([8.4, 3.9, 2.7, 8.4, 3.9]),
'area_sq_mi': pa.array([302, 469, 234, 302, 469])
})
# Register the Arrow table as a view in DuckDB
con.register('my_arrow_table', arrow_data)
arrow_query = """
SELECT
city,
SUM(population) AS total_population
FROM my_arrow_table
WHERE area_sq_mi > 300
GROUP BY city
ORDER BY total_population DESC;
"""
start_time = time.time()
arrow_result_df = con.execute(arrow_query).fetchdf()
end_time = time.time()
print(f"Arrow Table query executed in {end_time - start_time:.4f} seconds.")
print(arrow_result_df)
# 5. Clean up
con.close()
os.remove(large_parquet_file)
Architectural Comparison
| Feature | Apache Parquet | Apache Arrow | DuckDB |
|---|---|---|---|
| Purpose | Disk storage format | In-memory data format | In-process OLAP database engine |
| Data Location | Persistent storage (HDFS, S3, local disk) | RAM | RAM (can spill to disk) |
| Primary Use | Long-term archival, data lakes, data warehousing | Inter-process communication, zero-copy data exchange | Analytical queries, ETL, data exploration |
| Memory Layout | Columnar, compressed, page-oriented | Columnar, contiguous, uncompressed (mostly) | Columnar, vectorized, internal representation |
| Serialization | High (disk I/O, compression/decompression) | Zero-copy (if already in Arrow format) | Minimal (operates directly on columnar data) |
| Query Engine | None (storage only) | None (format only) | Full SQL engine, vectorized execution |
| Predicate Pushdown | Yes (via column statistics) | N/A (in-memory, direct access) | Yes (leveraging Parquet stats, internal optimizations) |
| SIMD Optimization | Indirect (decompression, then processing) | Direct (memory layout enables SIMD) | Direct (core of vectorized engine) |
| Typical Latency | Seconds to minutes (for large scans) | Milliseconds (data transfer) | Milliseconds to seconds (for complex queries) |
Production Gotchas & Troubleshooting
- Parquet Row Group Size Misconfiguration:
- Symptom: Extremely slow queries on Parquet files, especially with
WHEREclauses. - Problem: Too many small row groups lead to excessive metadata reads and high overhead. Too few large row groups reduce pruning effectiveness, forcing the engine to read more data than necessary.
- Fix: Aim for row groups between 128MB and 512MB uncompressed size. This typically translates to 100,000 to 1,000,000 rows per row group, depending on schema width. Use
pq.write_table(..., row_group_size=...)or equivalent in other languages.
- Symptom: Extremely slow queries on Parquet files, especially with
- Missing or Stale Parquet Statistics:
- Symptom: Query engines (like DuckDB, Spark) perform full scans even with highly selective predicates.
- Problem: Parquet files were written without statistics (
write_statistics=False) or the statistics are outdated due to data corruption/manual modification. - Fix: Always ensure
write_statistics=Truewhen writing Parquet. If files are already written, consider rewriting them or using tools that can generate/update statistics (e.g.,ANALYZE TABLEin some data warehouses, or re-ingestion).
- Arrow Memory Leaks in Long-Running Processes:
- Symptom: Python processes using
pyarrowconsume increasing amounts of RAM over time, eventually crashing. - Problem:
pyarrowobjects, especiallyTableandRecordBatch, manage off-heap memory. If references are held inadvertently, or ifclose()methods are not called on IPC readers/writers, memory might not be released promptly. - Fix: Explicitly delete large Arrow objects (
del my_arrow_table) when no longer needed. For IPC streams, ensurereader.close()andwriter.close()are called. Usegc.collect()for aggressive garbage collection in Python, though this is often a band-aid for deeper reference management issues. Profile memory usage with tools likememrayortracemalloc.
- Symptom: Python processes using
- DuckDB Out-of-Memory (OOM) Errors:
- Symptom: DuckDB crashes or reports OOM errors during complex queries, especially joins or aggregations on large datasets.
- Problem: DuckDB, while efficient, still needs RAM for intermediate results. If the working set exceeds available memory, it will spill to disk, but if the spill itself is too large or the system runs out of disk space, OOM can occur.
- Fix:
- Increase available RAM for the process.
- Optimize queries: filter early, project only necessary columns.
- Tune DuckDB memory limits:
SET memory_limit='8GB';. - Ensure sufficient temporary disk space for spills.
- For very large joins, consider pre-aggregating data or using a different join strategy if possible.
- Performance Degradation with Small DuckDB Batches:
- Symptom: DuckDB queries are slower than expected, especially on smaller datasets or when processing data from non-optimized sources.
- Problem: While DuckDB is vectorized, if the input data source provides very small batches (e.g., 1 row at a time from a custom iterator), the overhead of vectorized processing can outweigh the benefits.
- Fix: Ensure data is fed to DuckDB in reasonably sized batches (e.g., 1024-4096 rows) if using custom data sources or iterators. When reading Parquet or Arrow, DuckDB handles batching internally.
Frequently Asked Questions
-
When should I use Parquet versus Arrow? Use Parquet for long-term, persistent storage on disk, especially in data lakes or warehouses, where compression and predicate pushdown are critical for cost and I/O efficiency. Use Arrow for in-memory data transfer and processing between different systems or components within the same system, where zero-copy semantics and high-performance vectorized operations are paramount. They are complementary, not mutually exclusive.
-
Can DuckDB directly read data from S3 or other cloud storage? Yes, DuckDB has built-in support for reading Parquet and CSV files directly from S3, GCS, and Azure Blob Storage. You typically need to configure credentials (e.g., AWS access keys) via
SETcommands or environment variables. For example:INSTALL httpfs; LOAD httpfs; SET s3_region='us-east-1'; SELECT * FROM 's3://my-bucket/data.parquet';. -
How does DuckDB handle data types and schema evolution from Parquet? DuckDB infers the schema directly from the Parquet file metadata. It supports a wide range of data types, including nested types. For schema evolution, if new columns are added to a Parquet file, DuckDB will typically ignore them unless explicitly queried. If column types change incompatibly, it will raise an error. It's generally robust to additive schema changes.
-
What are the performance implications of using
pyarrow.Table.to_pandas()? Converting a largepyarrow.Tableto apandas.DataFrameinvolves copying data from Arrow's contiguous memory buffers into pandas' internal, often fragmented, memory layout. This is a copy operation, not zero-copy, and can be a performance bottleneck and memory hog for very large datasets. For analytical tasks, it's often more efficient to perform operations directly onpyarrow.Tableusing libraries likepolarsorduckdbwhich can consume Arrow tables directly. -
Is DuckDB suitable for production OLTP workloads? No. DuckDB is an OLAP (Online Analytical Processing) database, optimized for complex analytical queries over large datasets. It is not designed for high-concurrency, low-latency transactional workloads (OLTP) that require frequent small writes, updates, and deletes. For OLTP, traditional relational databases like PostgreSQL or MySQL are more appropriate. DuckDB excels as an embedded analytical engine, often replacing Spark or Pandas for local data processing.
Free In-Browser Developer Tools
Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.
Related Articles

Open Table Formats in 2026: Apache Iceberg v2 vs Delta Lake 3.0 on Cloud Object Storage
Comprehensive guide covering open table formats in 2026: apache iceberg v2 vs delta lake 3.0 on cloud object storage with production-grade architecture and code examples.
Read more
Fast Data Science: DuckDB and Polars for High-Performance Analytics
Ditch Pandas memory bloat: accelerate data pipelines with Polars (Rust multi-threading, LazyFrames, Arrow) and DuckDB (in-process vectorized SQL, out-of-core Parquet streaming).
Read more
ClickHouse vs DuckDB in 2026: In-Memory Embedded OLAP vs Distributed Vectorized Warehouses
Comprehensive guide covering clickhouse vs duckdb in 2026: in-memory embedded olap vs distributed vectorized warehouses with production-grade architecture and code examples.
Read more