Apache Arrow: in-memory columnar format, zero-copy передача и Record Batch как граница Python ↔ JVM ↔ C++
Apache Arrow: in-memory columnar format, zero-copy передача и Record Batch как граница Python ↔ JVM ↔ C++
Проблема межпроцессного взаимодействия в PySpark¶
Архитектурный раскол: два мира PySpark¶
Apache Spark написан на Scala и работает на JVM (Java Virtual Machine). Когда появился PySpark - Python-интерфейс к Spark - инженеры столкнулись с фундаментальной проблемой: нужно позволить Python-коду взаимодействовать с JVM-кодом. Это два принципиально разных рантайма, два разных процесса операционной системы, две разные модели памяти.
В типичной PySpark-архитектуре:
- Driver: Python-процесс (ваш
spark-submit my_script.py) + JVM-процесс SparkContext, общающиеся через Py4J (RPC-мост) - Executor: JVM-процесс, выполняющий задачи Spark. При наличии Python UDF - дополнительно запускается Python-воркер (отдельный процесс для выполнения Python-кода)
[Python Driver] ←── Py4J (RPC) ──→ [JVM SparkContext]
│
▼
[JVM Executor 1] [JVM Executor 2]
│ │
[Python Worker 1] [Python Worker 2]
(только при Python UDF)
Ключевое: Python Worker - это отдельный процесс. Он не разделяет память с JVM Executor. Чтобы передать данные из JVM в Python - их нужно физически скопировать через границу процессов.
Py4J и построчная сериализация: «проклятие»¶
До появления Arrow PySpark передавал данные между JVM и Python через Py4J и Pickle-сериализацию, по одной строке за раз.
Вот что происходило при выполнении Python UDF на датасете из N строк:
- JVM executor берёт строку из UnsafeRow-буфера
- Конвертирует её в Java-объекты (
Integer,String,Double) - Сериализует в байты через Pickle (Python-протокол сериализации объектов)
- Записывает в UNIX socket (локальный сетевой интерфейс)
- Python-воркер читает из сокета
- Десериализует из Pickle в Python-объекты (
int,str,float) - Вызывает Python UDF с этими объектами
- Сериализует результат обратно в Pickle
- Записывает в сокет → JVM читает, десериализует
Для каждой из N строк - 9 шагов. При N = 100 000 000 строк - 900 000 000 операций только на коммуникацию.
Накладные расходы огромны:
- Pickle-сериализация/десериализация: CPU-bound операция
- Создание миллиарда Python-объектов: давление на GC Python
- Чтение/запись в сокет: даже localhost имеет overhead
- Создание JVM-объектов для Pickle: давление на JVM GC
На практике Python UDF работал в 10-100 раз медленнее эквивалентной Scala-функции. Это делало Python UDF практически неприменимыми для Big Data - накладные расходы на коммуникацию съедали 80-90% времени.
Почему именно такая архитектура¶
Py4J был разумным решением для своего времени (2014): он позволил использовать Spark API из Python без переписывания всего на Python. Цена - коммуникационный overhead. Для интерактивной работы (небольшие датасеты, быстрый прототип) это приемлемо. Для production-нагрузок на терабайты - катастрофа.
Apache Arrow решил эту проблему не путём оптимизации Py4J, а через принципиально иную модель обмена данными: вместо строчной сериализации - передача готовых колоночных буферов через разделяемую память.
Анатомия Apache Arrow: единый стандарт памяти¶
Что такое Apache Arrow¶
Apache Arrow - открытый стандарт, определяющий точный физический layout колоночных данных в оперативной памяти. Появился в 2016 году как проект Apache, созданный при участии Wes McKinney (автор Pandas) и разработчиков из Dremio, Two Sigma, InflightData.
Ключевое слово - стандарт. Arrow не просто описывает концепцию, он задаёт до байта точную организацию данных: где хранится значение 42 для колонки int32, как закодирован NULL, как определяется длина строки. Благодаря этому:
- Код на C++ читает Arrow-буфер, созданный Java-кодом - без сериализации
- Python-код читает Arrow-буфер, созданный Rust-кодом - без копирования
- Spark передаёт данные в DuckDB - zero-copy
До Arrow каждая система имела свой формат: Pandas - NumPy arrays с object-колонками для строк; JVM Spark - UnsafeRow blob; R - SEXP-структуры; ClickHouse - свой внутренний формат. Передача данных между любыми двумя системами требовала промежуточной сериализации.
Arrow стал lingua franca аналитических систем - универсальным языком обмена данными.
Columnar in-memory: почему колонки, а не строки¶
Arrow хранит данные по колонкам, а не по строкам. Для аналитических запросов это принципиально важно по трём причинам:
Кэш-локальность при агрегации. Запрос SUM(amount) читает только колонку amount. При колоночном хранении все значения amount лежат рядом - одна непрерывная область памяти. CPU-prefetcher загружает их в L1-кэш быстро и без промахов. При строчном хранении значения amount разделены другими полями каждой строки.
SIMD-дружественность. SIMD-инструкция AVX-512 загружает 64 байта подряд (16 × int32). При колоночном хранении это ровно 16 последовательных значений amount. При строчном - бессмысленная смесь разных полей разных строк.
Эффективное сжатие. Данные одного типа одного домена сжимаются лучше, чем перемежающиеся данные разных типов. Dictionary encoding для колонки status с 5 значениями даёт 10-50× сжатие.
Разница Arrow vs Parquet¶
Частое заблуждение: Arrow и Parquet - это «одно и то же». На самом деле это разные форматы для разных целей:
| Свойство | Apache Arrow | Apache Parquet |
|---|---|---|
| Назначение | In-memory processing | On-disk storage |
| Формат | Не сжатый (или лёгкое сжатие) | Глубокое сжатие |
| Оптимизация | CPU computation, IPC | Storage I/O, disk space |
| Поиск внутри | O(1) по индексу | O(log N) по row group |
| Поддержка NULL | Validity bitmap | Definition levels |
| Применение | Обмен между процессами, in-memory DB | Долгосрочное хранение |
Типичный пайплайн: Parquet на диске → читаем → конвертируем в Arrow in-memory → обрабатываем → пишем обратно в Parquet. Arrow - это «рабочий формат», Parquet - «архивный».
Разница Arrow vs Pandas¶
Pandas (до версии 2.0) хранит данные в NumPy arrays - это тоже колоночный формат, но с отличиями:
- Object dtype: строки, смешанные типы, Python-объекты хранятся как массив указателей на Python-объекты. Это разрушает cache-locality и не совместимо с SIMD
- Нет стандартизированного layout для NULL: Pandas использует
NaNдля float (семантика IEEE 754), отдельныйNone/pd.NAдля других типов - непоследовательно - Нет cross-language стандарта: Pandas layout специфичен для Python
Arrow решает все три проблемы: строки хранятся в contiguous буферах (не указатели), NULL - через validity bitmap, layout стандартизирован для C++/Java/Rust.
Pandas 2.0 ввёл ArrowDtype - бэкенд на Apache Arrow вместо NumPy. Это делает современный Pandas нативно совместимым с Arrow-экосистемой.
Физический layout Arrow: буферы, массивы, Record Batch¶
Структура Arrow Array¶
Каждая колонка в Arrow - это Array: типизированная структура из нескольких буферов. Для разных типов данных структура разная.
Fixed-size primitive array (int32, int64, float64):
Validity Bitmap Buffer:
Бит i = 1 → значение не NULL
Бит i = 0 → значение NULL
Биты упакованы: 8 значений в 1 байт
Выровнен по 64 байта (кэш-строка = AVX-512 регистр)
Пример для [1, NULL, 3, 4, NULL]:
Bytes: 0b00001101 (биты 0,2,3 = 1; биты 1,4 = 0)
↑↑↑↑↑
43210 (порядок: бит 0 = первый элемент)
Values Buffer:
Непрерывный массив примитивов:
[1][?][3][4][?] ← для NULL сохраняется 0 (или любое значение, игнорируется)
Каждый элемент: фиксированный размер (int32 = 4 байта)
Выровнен по 64 байта
Важно: NULL-значение не отсутствует физически. Оно присутствует в Value Buffer (обычно как 0), а Validity Bitmap указывает, что его нужно игнорировать. Это позволяет SIMD-инструкциям работать с буфером без conditional checks - NULL обрабатывается через маску.
Variable-size binary array (string, bytes):
Строки имеют переменную длину - нельзя хранить их как fixed-size array. Arrow использует три буфера:
Validity Bitmap Buffer:
[1][1][0][1] ← NULL маска
Offsets Buffer (int32 или int64):
[0][5][12][12][18] ← N+1 значений offset
offset[i] = начало строки i в Data Buffer
offset[i+1] = конец строки i (= начало строки i+1)
Длина строки i = offset[i+1] - offset[i]
Data Buffer (bytes):
[h][e][l][l][o][w][o][r][l][d][ ][f][o][o]
0 1 2 3 4 5 6 7 8 9 10 11 12 13 ...
Для строки 0: offset[0]=0, offset[1]=5 → bytes [0..4] = "hello" Для строки 1: offset[1]=5, offset[2]=12 → bytes [5..11] = "world "
NULL-строка (индекс 2): offset[2]=12, offset[3]=12 → длина 0, данных нет
Такая структура позволяет:
- O(1) доступ к строке i по индексу (без сканирования)
- Contiguous хранение всех данных (хорошо для кэша и SIMD)
- Эффективный range scan строк
Dictionary-encoded array:
Dictionary Array:
Indices Buffer: [0][2][1][0][2][1] ← int8/int16/int32 индексы
Validity Bitmap: [1][1][1][1][1][1]
Dictionary Values Array (string):
["pending", "completed", "cancelled"] ← словарь уникальных значений
Итоговые данные: ["pending", "cancelled", "completed", "pending", "cancelled", "completed"]
Dictionary encoding экономит память при повторяющихся значениях. Для колонки status с 5 уникальными значениями и 100M строк: вместо 100M строк - только индексы (1-4 байта на строку) + словарь из 5 строк. Фильтрация WHERE status = 'completed' по индексам (int) без декодирования строк - быстрее.
Выравнивание памяти: 64 байта как ключевое число¶
Arrow требует, чтобы все буферы начинались с адреса, кратного 64 байтам. Не случайно: 64 байта - это размер кэш-строки процессора и размер AVX-512 векторного регистра.
Выровненный буфер гарантирует:
- VMOVDQA32 (aligned load) вместо VMOVDQU32 (unaligned load) - aligned быстрее на 10-15%
- Кэш-строка никогда не разрезает векторный load: если буфер выровнен по 64 байтам, загрузка 64 байт никогда не пересечёт границу кэш-строки
Это кажется деталью реализации, но в нагруженных аналитических пайплайнах эти 10-15% складываются в ощутимый выигрыш.
Schema: типовая информация для всего batch¶
Arrow Schema - метаданные, описывающие типы всех колонок:
import pyarrow as pa
schema = pa.schema([
pa.field("user_id", pa.int64()),
pa.field("amount", pa.float64()),
pa.field("status", pa.dictionary(pa.int8(), pa.utf8())),
pa.field("created_at", pa.timestamp("us", tz="UTC")),
pa.field("tags", pa.list_(pa.utf8())), # Variable-length list
pa.field("metadata", pa.struct([ # Nested struct
pa.field("source", pa.utf8()),
pa.field("version", pa.int32()),
])),
])
print(schema)
Schema - это самодостаточные метаданные: получатель данных знает типы без дополнительного запроса. При IPC-передаче Schema отправляется один раз в начале потока, затем только данные Record Batch'ей.
Record Batch: основная единица передачи¶
Record Batch - это фиксированное количество строк, представленных как набор Arrow Array, все одинаковой длины:
import pyarrow as pa
import numpy as np
# Создаём Record Batch из 5 строк
batch = pa.record_batch({
"user_id": pa.array([1, 2, 3, 4, 5], type=pa.int64()),
"amount": pa.array([100.0, None, 250.0, 75.0, 500.0], type=pa.float64()),
"status": pa.array(["completed", "pending", "completed", "cancelled", "completed"]),
})
print(f"Схема: {batch.schema}")
print(f"Строк: {batch.num_rows}") # 5
print(f"Колонок: {batch.num_columns}") # 3
print(f"Размер (bytes): {batch.nbytes}") # Размер всех буферов
# Инспекция конкретной колонки
amount_array = batch.column("amount")
print(f"Validity bitmap: {amount_array.is_valid}")
# [True, False, True, True, True] ← строка 1 (None) = False
Record Batch - это граница между системами. Именно Record Batch передаётся:
- Из Spark JVM в Python Worker (pandas UDF)
- Из Python Worker в Spark JVM (результат pandas UDF)
- Из Spark в DuckDB (zero-copy)
- Из PySpark в ML-фреймворк (TensorFlow, PyTorch)
- По сети через Arrow Flight Protocol
Ключевое свойство: Record Batch самодостаточен - содержит и данные, и схему (типы), и NULL-маски. Получатель не нуждается в дополнительном контексте.
Zero-Copy: передача без копирования¶
Что такое zero-copy¶
Zero-copy - это техника передачи данных между процессами или компонентами системы без физического копирования байт из одной области памяти в другую. Вместо копирования передаётся указатель (адрес начала буфера) или файловый дескриптор разделяемой памяти.
Обычная передача данных (с копированием):
Process A:
[Data Buffer at addr 0x7f00]
│
▼ memcpy() или socket write()
Копируем N байт в новый буфер
│
▼ socket read() и memcpy()
Process B:
[Data Buffer at addr 0x3e00] ← НОВАЯ область памяти, КОПИЯ
Zero-copy через разделяемую память:
OS: Создаём shared memory segment
Process A:
[Data Buffer at addr 0x7f00] ← Пишем данные в shared memory
│
▼ Передаём только:
- Адрес/handle разделяемой памяти (8 байт)
- Смещение (8 байт)
- Размер (8 байт)
│
Process B:
[Маппируем тот же shared memory segment]
[Data Buffer at addr 0x7f00] ← ТОТ ЖЕ физический адрес!
Никакого копирования!
В контексте Apache Arrow: JVM создаёт Arrow-буфер в off-heap памяти (за пределами Java heap). Python Worker маппирует эту же память через IPC handle. Нет Pickle, нет сокетов, нет копирования - Python просто «видит» ту же область памяти, что и JVM.
Механизм Arrow IPC¶
Arrow IPC (Inter-Process Communication) - это протокол для передачи Record Batch'ей между процессами. Существует два режима:
File Format: Record Batch'и сохраняются в файл (или memory-mapped file) для последующего чтения.
Stream Format: Record Batch'и передаются последовательно как поток (через сокет, pipe, shared memory).
Структура Arrow IPC Stream:
[Schema Message] ← Типы колонок (один раз в начале)
[Record Batch Message] ← Данные (meta: буферы offsets+lengths)
[Record Batch Message]
[Record Batch Message]
...
[EOS (End of Stream)]
Каждый Record Batch Message состоит из:
- FlatBuffer metadata: структурированное описание схемы батча, количество строк, смещения и длины каждого буфера - всего несколько килобайт
- Raw binary data: сами данные колонок - основной объём
При zero-copy передаче через shared memory:
- Process A записывает raw binary data в shared memory
- Передаёт Process B только FlatBuffer metadata (несколько KB через обычный сокет)
- Process B маппирует shared memory и интерпретирует данные используя metadata
Overhead: только metadata (KB) вместо полного объёма данных (GB).
Zero-copy в практике PySpark¶
В PySpark с включённым Arrow передача данных происходит через Arrow IPC over socket (не полный zero-copy через shared memory, но значительно эффективнее Pickle):
JVM Executor Python Worker
│ │
├─ Собираем ColumnarBatch ───────► │
│ (данные уже в off-heap) │
├─ Конвертируем в Arrow ────────► │
│ RecordBatch (нет копирования, │
│ только metadata creation) │
├─ Записываем в socket ─────────► │
│ ├─ Читаем из socket
│ ├─ Создаём pa.RecordBatch
│ │ (wrap вокруг буфера)
│ ├─ Конвертируем в pd.DataFrame
│ │ (Arrow → NumPy arrays)
│ ├─ Вызываем Pandas UDF
│ ├─ pd.Series → pa.Array
│◄── Читаем результат ──────────────── ├─ Записываем в socket
Для 100M строк × 5 колонок × 8 байт = 4 GB данных:
- Pickle (без Arrow): 4 GB × 100M iterations × serialize/deserialize = часы
- Arrow IPC (с Arrow): 4 GB × 1 batch transfer × write/read = секунды
Как Arrow меняет PySpark¶
Включение Arrow-оптимизации¶
# Включить Arrow для всех PySpark ↔ Python операций
spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "true")
# Размер Record Batch при передаче (строк на батч)
# Большие батчи = лучше для SIMD, хуже для памяти
spark.conf.set("spark.sql.execution.arrow.maxRecordsPerBatch", "10000")
# Безопасный режим: если Arrow не работает для данного типа - fallback на Pickle
spark.conf.set("spark.sql.execution.arrow.pyspark.fallback.enabled", "true")
В Spark 3.x Arrow включён по умолчанию (true). Но понимание этой конфигурации важно при диагностике проблем производительности.
toPandas() и createDataFrame(): до и после Arrow¶
Самое наглядное применение Arrow - конвертация между Spark DataFrame и Pandas DataFrame:
import pandas as pd
import time
# Генерируем 10M строк в Spark
df_spark = spark.range(10_000_000) \
.withColumn("amount", (F.rand() * 1000).cast("double")) \
.withColumn("category", (F.rand() * 100).cast("int").cast("string")) \
.withColumn("status", F.when(F.rand() > 0.5, "completed").otherwise("pending"))
# Вариант 1: WITHOUT Arrow
spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "false")
start = time.time()
df_pandas_slow = df_spark.toPandas()
t_without = time.time() - start
print(f"Без Arrow (Pickle row-by-row): {t_without:.1f}s")
# Вариант 2: WITH Arrow
spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "true")
start = time.time()
df_pandas_fast = df_spark.toPandas()
t_with = time.time() - start
print(f"С Arrow (batch columnar): {t_with:.1f}s")
print(f"Ускорение: {t_without / t_with:.1f}x")
Типичный результат для 10M строк:
- Без Arrow: ~120 секунд (row-by-row Pickle)
- С Arrow: ~3-5 секунд (batch Arrow IPC)
- Ускорение: 25-40×
Что происходит под капотом с Arrow: Spark executor собирает все строки в Arrow RecordBatch'и → сериализует через Arrow IPC → Driver получает Arrow stream → pa.Table → Pandas DataFrame. Вместо 10M Python-объектов - несколько больших буферов.
createDataFrame() с Pandas: обратная передача¶
# Pandas → Spark с Arrow
import pandas as pd
import numpy as np
df_large_pandas = pd.DataFrame({
"user_id": np.random.randint(0, 1_000_000, size=5_000_000),
"revenue": np.random.uniform(0, 10_000, size=5_000_000),
"category": np.random.choice(["A", "B", "C", "D"], size=5_000_000),
})
# Без Arrow: ~60s (построчное создание Row объектов)
spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "false")
start = time.time()
df_spark_slow = spark.createDataFrame(df_large_pandas)
df_spark_slow.count() # trigger action
print(f"Без Arrow: {time.time() - start:.1f}s")
# С Arrow: ~2-4s (columnar Arrow IPC → Spark ColumnarBatch)
spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "true")
start = time.time()
df_spark_fast = spark.createDataFrame(df_large_pandas)
df_spark_fast.count()
print(f"С Arrow: {time.time() - start:.1f}s")
pandas_udf: батчевая обработка через Arrow¶
pandas_udf - это декоратор, который говорит Spark: «передавай данные в мою Python-функцию не строка за строкой, а батчами как Pandas Series/DataFrame, и используй Arrow для передачи».
from pyspark.sql.functions import pandas_udf
import pandas as pd
# Скалярный pandas_udf: принимает Series, возвращает Series
@pandas_udf("double")
def calculate_tax(amount: pd.Series, rate: pd.Series) -> pd.Series:
"""
Вызывается не N раз (по строке), а K раз (по батчу),
где K = N / batch_size.
Внутри - NumPy-векторизованные операции (SIMD!).
"""
return (amount * rate).where(amount > 0, 0.0)
# Arrow lifecycle:
# 1. Spark собирает batc h строк (default: 10000)
# 2. Конвертирует в Arrow RecordBatch
# 3. Передаёт в Python via IPC
# 4. Python видит pa.RecordBatch → pd.DataFrame колонки → pd.Series
# 5. UDF выполняется (NumPy vectorized)
# 6. Результат pd.Series → pa.Array → Arrow IPC → JVM
result = df_spark.withColumn(
"tax",
calculate_tax(F.col("amount"), F.lit(0.2))
)
В Spark UI физический план покажет разницу:
# Без Arrow (обычный udf):
*(1) Project [user_id#1, BatchEvalPython [udf(amount#2)], [amount#2]]
↑
Построчная сериализация!
# С Arrow (pandas_udf):
*(1) Project [user_id#1, ArrowEvalPython [calculate_tax(amount#2, 0.2)], [amount#2]]
↑
Батчевая Arrow-передача!
Оператор ArrowEvalPython - явное подтверждение в EXPLAIN, что используется Arrow.
Grouped Map pandas_udf: обработка на уровне партиции¶
Более мощный вариант - применение Pandas-функции ко всей группе данных:
from pyspark.sql.functions import pandas_udf, PandasUDFType
schema = "user_id long, amount double, percentile double"
@pandas_udf(schema, PandasUDFType.GROUPED_MAP)
def add_percentile(group: pd.DataFrame) -> pd.DataFrame:
"""
group - это pd.DataFrame со всеми строками одного user_id.
Всё это через один Arrow RecordBatch!
"""
group["percentile"] = group["amount"].rank(pct=True)
return group
result = df_spark.groupBy("user_id").apply(add_percentile)
Spark передаёт все строки одного user_id как единый Arrow RecordBatch в Python. Внутри функции - полноценный Pandas DataFrame с возможностью любых Pandas-операций.
Инспекция Arrow в Python: практический разбор¶
Работа с pyarrow напрямую¶
import pyarrow as pa
import pyarrow.compute as pc
import numpy as np
# Создаём Arrow Array вручную
amounts = pa.array([100, None, 250, 75, 500, None, 180], type=pa.int32())
statuses = pa.array(["completed", "pending", "completed", "cancelled",
"completed", "pending", "cancelled"])
# Инспекция NULL-маски
print(f"Validity bitmap: {amounts.is_valid.to_pylist()}")
# [True, False, True, True, True, False, True]
print(f"NULL count: {amounts.null_count}") # 2
# Все буферы колонки
for i, buf in enumerate(amounts.buffers()):
if buf is not None:
print(f"Buffer {i}: {buf.size} bytes at address {buf.address:#x}")
# Buffer 0 (validity): 1 bytes at 0x7f1234560000 ← NULL bitmap
# Buffer 1 (values): 28 bytes at 0x7f1234570040 ← 7 × int32
# Создаём Record Batch
batch = pa.record_batch({
"amount": amounts,
"status": statuses,
})
print(f"Batch: {batch.num_rows} rows, {batch.nbytes} bytes")
# Операции через Arrow Compute (без Python loops)
# pc = pyarrow.compute - все операции векторизованы (SIMD!)
filtered = pc.filter(batch, pc.greater(amounts, pa.scalar(100)))
total = pc.sum(filtered.column("amount"))
print(f"Сумма amount > 100: {total}")
# Конвертация в Pandas и обратно
df = batch.to_pandas()
print(df)
batch_back = pa.RecordBatch.from_pandas(df)
print(f"Восстановлено: {batch_back.num_rows} rows")
Инспекция Arrow при PySpark pandas_udf¶
import pyarrow as pa
@pandas_udf("long")
def inspect_and_compute(series: pd.Series) -> pd.Series:
"""
Внутри pandas_udf мы можем получить Arrow-буфер напрямую.
Это показывает, что Pandas Series backed by Arrow buffer.
"""
# Pandas Series → Apache Arrow Array
arrow_array = pa.Array.from_pandas(series)
print(f"Arrow type: {arrow_array.type}")
print(f"Buffer address: {arrow_array.buffers()[1].address:#x}") # Values buffer
print(f"Buffer size: {arrow_array.buffers()[1].size} bytes")
print(f"NULL count: {arrow_array.null_count}")
# SIMD-векторизованная операция через NumPy (backed by Arrow buffer)
return series ** 2 # NumPy square - AVX-optimized
result = df_spark.select(inspect_and_compute("amount"))
Измерение накладных расходов IPC¶
import time
import pyarrow as pa
import pyarrow.ipc as ipc
import io
# Создаём большой Record Batch (симуляция)
n = 1_000_000
batch = pa.record_batch({
"user_id": pa.array(np.random.randint(0, 1_000_000, n), type=pa.int64()),
"amount": pa.array(np.random.uniform(0, 10_000, n), type=pa.float64()),
"category": pa.array(np.random.choice(["A","B","C","D","E"], n)),
})
print(f"Batch size: {batch.nbytes / 1_000_000:.1f} MB")
# Бенчмарк: Arrow IPC serialization/deserialization
buffer = io.BytesIO()
start = time.perf_counter()
writer = ipc.new_stream(buffer, batch.schema)
writer.write_batch(batch)
writer.close()
t_write = time.perf_counter() - start
buffer.seek(0)
start = time.perf_counter()
reader = ipc.open_stream(buffer)
batch_read = reader.read_next_batch()
t_read = time.perf_counter() - start
print(f"IPC Write: {t_write*1000:.1f}ms ({batch.nbytes/t_write/1e9:.1f} GB/s)")
print(f"IPC Read: {t_read*1000:.1f}ms ({batch.nbytes/t_read/1e9:.1f} GB/s)")
print(f"Итого overhead IPC: {(t_write+t_read)*1000:.1f}ms для {batch.nbytes/1e6:.0f} MB")
Типичный вывод:
Batch size: 48.0 MB
IPC Write: 42ms (1.1 GB/s)
IPC Read: 18ms (2.7 GB/s)
Итого overhead IPC: 60ms для 48 MB
Для сравнения: Pickle serialization тех же данных - 5-10 секунд вместо 60 мс.
Arrow как ядро Modern Data Stack¶
Диаграмма: Arrow как lingua franca аналитических систем¶
Схема показывает Apache Arrow как центральный узел современного аналитического стека. Каждая система понимает Arrow - это значит, что данные между любыми двумя из них передаются без сериализации. Spark отдаёт Arrow buffer в Velox - нет конвертации. Python получает Arrow buffer от Spark - нет Pickle. DuckDB читает Arrow RecordBatch от Polars - нет копирования. Arrow - это не просто формат, это инфраструктура нулевых накладных расходов между системами.
Spark + DuckDB: zero-copy аналитический пайплайн¶
Один из самых мощных применений Arrow - совместное использование Spark и DuckDB. Spark обрабатывает петабайты распределённо, DuckDB выполняет аналитические запросы в памяти. Передача через Arrow - практически без overhead:
import duckdb
import pyarrow as pa
# Spark вычисляет и возвращает как Arrow Table
df_spark = spark.sql("""
SELECT user_id, SUM(amount) AS revenue, COUNT(*) AS orders
FROM transactions
WHERE event_date >= '2024-01-01'
GROUP BY user_id
""")
# Конвертируем в Arrow (Arrow → DuckDB через IPC)
arrow_table = df_spark.toArrow() # Spark 3.3+
# или: arrow_table = pa.Table.from_batches(df_spark._collect_as_arrow())
# DuckDB читает Arrow Table без копирования
con = duckdb.connect()
con.register("revenue_data", arrow_table) # Регистрируем как virtual table
# DuckDB выполняет запрос с AVX-512 нативной обработкой
result = con.execute("""
SELECT
PERCENTILE_CONT(0.5) WITHIN GROUP (ORDER BY revenue) AS median_revenue,
PERCENTILE_CONT(0.95) WITHIN GROUP (ORDER BY revenue) AS p95_revenue,
AVG(revenue) AS avg_revenue,
STDDEV(revenue) AS stddev_revenue
FROM revenue_data
WHERE orders >= 3
""").fetchdf()
print(result)
DuckDB читает Arrow Table через zero-copy (общий буфер памяти). Всё сложное агрегирование (PERCENTILE - квантили - дорогая операция) выполняется DuckDB с полной SIMD-векторизацией. Результат - обычный Pandas DataFrame.
Spark + Polars: аналитика на Rust¶
Polars - DataFrame-библиотека на Rust с нативной Arrow-поддержкой. Отлично подходит для обработки средних объёмов (GB) с максимальной скоростью:
import polars as pl
import pyarrow as pa
# Получаем Arrow RecordBatch'и из Spark
arrow_batches = df_spark.toArrow() # pa.Table
# Polars читает Arrow Table напрямую (zero-copy)
df_polars = pl.from_arrow(arrow_batches)
# Polars выполняет сложную трансформацию с SIMD-векторизацией
result_polars = (
df_polars
.filter(pl.col("status") == "completed")
.with_columns([
(pl.col("amount") * 0.2).alias("tax"),
pl.col("created_at").dt.truncate("1d").alias("date"),
])
.group_by("date", "category")
.agg([
pl.col("amount").sum().alias("total_revenue"),
pl.col("amount").mean().alias("avg_order"),
pl.len().alias("order_count"),
])
.sort("date", "category")
)
# Возвращаем в Spark через Arrow (если нужно)
result_spark = spark.createDataFrame(result_polars.to_arrow())
Apache DataFusion + Comet: Arrow-native Spark execution¶
Apache DataFusion - это query engine на Rust с нативной Arrow-поддержкой. Apache DataFusion Comet - плагин для Spark (аналог Gluten/Velox), который заменяет Tungsten операторы DataFusion'овскими.
# Подключение Comet плагина к Spark
spark = SparkSession.builder \
.config("spark.plugins", "org.apache.datafusion.comet.CometPlugin") \
.config("spark.comet.enabled", "true") \
.config("spark.comet.exec.enabled", "true") \
.config("spark.comet.exec.all.operator.enabled", "true") \
.getOrCreate()
# Тот же Spark SQL API - но выполняется на Rust + Arrow + AVX-512
df = spark.read.parquet("s3a://bucket/data/")
result = df.filter("amount > 100").groupBy("user_id").sum("amount")
result.show()
# В логах: CometExec operators вместо Tungsten
Arrow Flight: высокопроизводительный RPC для аналитики¶
Arrow Flight - это высокопроизводительный RPC-протокол для передачи Arrow данных по сети. По сути - gRPC, но вместо Protobuf использует Arrow IPC для передачи батчей.
Преимущество над JDBC/ODBC: JDBC/ODBC передают данные в строчном формате с Protobuf-сериализацией. Flight передаёт Arrow RecordBatch'и - колоночный формат без конвертации.
Сравнение скоростей (1 TB query result):
| Протокол | Скорость передачи | Overhead |
|---|---|---|
| JDBC (строчный) | ~50 MB/s | Высокий |
| Arrow Flight (колоночный) | ~3-5 GB/s | Минимальный |
Arrow Flight используется в: Dremio (полётная передача результатов запросов), Google BigQuery Storage API (быстрая выгрузка данных), Spark 4.x (Flight SQL в экспериментальном режиме).
Arrow limitations и fallback-поведение¶
Типы данных, которые Arrow не поддерживает в PySpark¶
Не все типы данных Spark конвертируются в Arrow. При несовместимом типе Spark автоматически переходит на Pickle (если включён fallback.enabled):
from pyspark.sql.types import *
# Эти типы Arrow НЕ поддерживает (в PySpark):
unsupported_schema = StructType([
StructField("map_col", MapType(StringType(), IntegerType())), # MapType
StructField("nested_list", ArrayType(ArrayType(StringType()))), # Вложенный Array
StructField("timestamp_ntz", TimestampNTZType()), # Зависит от версии
])
# Проверить, включён ли fallback:
print(spark.conf.get("spark.sql.execution.arrow.pyspark.fallback.enabled"))
# true - Spark молча переключится на Pickle для несовместимых колонок
При fallback в логах появится предупреждение:
WARN ArrowUtils: Failed to use Arrow optimization. ...
Falling back to non-Arrow conversion: MapType not supported
Если fallback выключен (false) - будет исключение. Это полезно для диагностики: убедиться, что Arrow работает для всей схемы.
Проблема timestamp и timezone¶
Timestamp - источник частых проблем при Arrow-передаче между Python и JVM:
# JVM Spark: TimestampType = микросекунды + UTC
# Python Pandas: datetime64[ns] = наносекунды + timezone-aware или naive
# Проблема: разные точности и разные timezone semantics
@pandas_udf("timestamp")
def round_to_hour(ts: pd.Series) -> pd.Series:
# pd.Series с dtype datetime64[us, UTC] (Arrow-backed)
return ts.dt.floor("h")
# Убедиться, что timezone корректна:
from pyspark.sql.types import TimestampType
spark.conf.set("spark.sql.session.timeZone", "UTC")
Рекомендация: всегда используйте UTC в spark.sql.session.timeZone при работе с pandas_udf и timestamp-колонками.
Размер батча: баланс между памятью и производительностью¶
# По умолчанию: 10000 строк на батч
spark.conf.get("spark.sql.execution.arrow.maxRecordsPerBatch") # "10000"
# Слишком маленький батч (1000):
# + Мало памяти на батч
# - Больше IPC-вызовов (overhead на каждый вызов UDF)
# - Хуже SIMD (меньше данных за раз)
spark.conf.set("spark.sql.execution.arrow.maxRecordsPerBatch", "1000")
# Слишком большой батч (100000):
# + Меньше IPC-вызовов
# + Лучше SIMD-векторизация
# - Больше памяти на Python Worker
# - Риск OOM при широких строках (много колонок × много строк)
spark.conf.set("spark.sql.execution.arrow.maxRecordsPerBatch", "100000")
# Практическое правило: batch_size × row_size_bytes < 100-200 MB
# Для row_size = 100 bytes → batch_size = 1_000_000-2_000_000
# Для row_size = 1000 bytes (широкие строки) → batch_size = 100_000-200_000
Диагностика Arrow-проблем¶
# Включить детальное логирование Arrow
import logging
logging.getLogger("py4j.java_gateway").setLevel(logging.INFO)
# Проверить в EXPLAIN, что Arrow используется:
df.withColumn("result", my_pandas_udf("col")).explain("formatted")
# Ищем: "ArrowEvalPython" - хорошо
# Ищем: "BatchEvalPython" - Pickle, не Arrow
# Проверить метрики в Spark UI:
# Stage → Task Metrics → Python execution time
# Если Python execution time >> Task duration - UDF overhead
Лабораторная практика: Arrow benchmark¶
Сценарий: выгрузка 20 GB в Pandas для ML-модели¶
Типичный сценарий в Data Science: Spark обрабатывает большой датасет, результат нужно передать в Python для обучения ML-модели.
import time
import psutil
import os
def get_memory_mb():
"""Получить использование памяти текущего процесса в MB."""
process = psutil.Process(os.getpid())
return process.memory_info().rss / 1024 / 1024
# Генерируем 20 GB датасет в Spark (приблизительно)
# 50M строк × ~400 байт = ~20 GB
df_large = (
spark.range(50_000_000)
.withColumn("user_id", (F.rand() * 10_000_000).cast("long"))
.withColumn("amount", (F.rand() * 10_000).cast("double"))
.withColumn("quantity", (F.rand() * 100).cast("int"))
.withColumn("price", (F.rand() * 1000).cast("double"))
.withColumn("category", (F.rand() * 50).cast("int").cast("string"))
.withColumn("status", F.when(F.rand() > 0.3, "completed")
.when(F.rand() > 0.6, "pending")
.otherwise("cancelled"))
.withColumn("created_at", F.current_timestamp())
.cache() # Кэшируем чтобы исключить I/O из benchmark
)
df_large.count() # Trigger cache
# Шаг 1: БЕЗ Arrow
spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "false")
mem_before = get_memory_mb()
start = time.time()
try:
df_pandas_slow = df_large.limit(1_000_000).toPandas() # Берём 1M для теста
t_slow = time.time() - start
mem_after = get_memory_mb()
print(f"Без Arrow: {t_slow:.1f}s, память: +{mem_after-mem_before:.0f} MB")
except MemoryError:
print("OOM! Без Arrow слишком много памяти")
# Шаг 2: С Arrow
spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "true")
mem_before = get_memory_mb()
start = time.time()
df_pandas_fast = df_large.limit(1_000_000).toPandas()
t_fast = time.time() - start
mem_after = get_memory_mb()
print(f"С Arrow: {t_fast:.1f}s, память: +{mem_after-mem_before:.0f} MB")
print(f"Ускорение: {t_slow/t_fast:.1f}x")
# Шаг 3: Инспекция dtypes - Arrow-backed Pandas
print("\nDtypes с Arrow:")
print(df_pandas_fast.dtypes)
# Заметим: Arrow-backed Pandas может показывать ArrowDtype для некоторых колонок
# Pandas 2.x: StringDtype ("string[python]") → ArrowDtype("large_string[pyarrow]")
Программный разбор Record Batch¶
import pyarrow as pa
import pyarrow.ipc as ipc
import numpy as np
# Создаём реалистичный Record Batch
n = 100_000
batch = pa.record_batch({
"user_id": pa.array(np.random.randint(0, 1_000_000, n), type=pa.int64()),
"amount": pa.array(np.where(
np.random.random(n) < 0.05, None,
np.random.uniform(0, 10_000, n)
), type=pa.float64()),
"status": pa.array(np.random.choice(
["completed", "pending", "cancelled"], n
)),
"created_at": pa.array(
pd.date_range("2024-01-01", periods=n, freq="1s"),
type=pa.timestamp("us", tz="UTC")
),
})
print("=== Record Batch Info ===")
print(f"Rows: {batch.num_rows:,}")
print(f"Columns: {batch.num_columns}")
print(f"Size: {batch.nbytes / 1024:.1f} KB")
print(f"Schema:\n{batch.schema}")
print("\n=== Column Analysis ===")
for col_name in batch.schema.names:
col = batch.column(col_name)
bufs = col.buffers()
print(f"\nColumn '{col_name}' ({col.type}):")
print(f" NULL count: {col.null_count}")
if bufs[0] is not None:
print(f" Validity bitmap: {bufs[0].size} bytes "
f"({col.null_count} NULLs)")
if bufs[1] is not None:
print(f" Values buffer: {bufs[1].size} bytes "
f"at addr {bufs[1].address:#x}")
# Проверяем выравнивание
if bufs[1] is not None:
aligned = bufs[1].address % 64 == 0
print(f" 64-byte aligned: {aligned}")
# Arrow Compute: векторизованные операции без Python loops
import pyarrow.compute as pc
completed_mask = pc.equal(batch.column("status"), "completed")
completed_amounts = pc.filter(batch.column("amount"), completed_mask)
total = pc.sum(completed_amounts)
mean = pc.mean(completed_amounts)
print(f"\n=== Compute Results ===")
print(f"Completed transactions: {pc.sum(completed_mask).as_py():,}")
print(f"Total revenue: {total.as_py():.2f}")
print(f"Avg order: {mean.as_py():.2f}")
Бенчмарк: оптимальный размер батча¶
import time
results = []
for batch_size in [100, 500, 1000, 5000, 10000, 50000, 100000]:
spark.conf.set("spark.sql.execution.arrow.maxRecordsPerBatch", str(batch_size))
@pandas_udf("double")
def dummy_computation(s: pd.Series) -> pd.Series:
return (s * 1.2).clip(lower=0)
df_test = spark.range(1_000_000).withColumn("amount", (F.rand() * 1000))
start = time.time()
df_test.withColumn("result", dummy_computation("amount")) \
.write.format("noop").mode("overwrite").save()
elapsed = time.time() - start
results.append({"batch_size": batch_size, "time_s": elapsed})
print(f"batch_size={batch_size:6d}: {elapsed:.2f}s")
# Типичный вывод:
# batch_size= 100: 45.2s ← слишком много вызовов UDF
# batch_size= 500: 18.7s
# batch_size= 1000: 11.2s
# batch_size= 5000: 5.8s
# batch_size= 10000: 4.1s ← sweet spot для большинства случаев
# batch_size= 50000: 3.9s ← почти то же, но больше памяти
# batch_size=100000: 3.8s ← риск OOM на широких схемах
Best Practices: максимизация Arrow-эффективности¶
Правило 1: всегда используйте pandas_udf вместо udf¶
# ПЛОХО: строчный udf → Pickle row-by-row
@F.udf("double")
def slow_udf(x):
return x * 1.2
# ХОРОШО: pandas_udf → Arrow batch
@pandas_udf("double")
def fast_udf(x: pd.Series) -> pd.Series:
return x * 1.2 # NumPy vectorized, SIMD-friendly
Правило 2: предпочитайте SQL-функции pandas_udf, а pandas_udf - обычному udf¶
# Лучшее: SQL built-in (JVM-native, SIMD в Tungsten, нет Python)
df.withColumn("result", F.col("amount") * 1.2)
# Хорошее: pandas_udf (Arrow batch, NumPy SIMD в Python)
@pandas_udf("double")
def multiply(s: pd.Series) -> pd.Series:
return s * 1.2
df.withColumn("result", multiply("amount"))
# Плохое: обычный udf (Pickle row-by-row)
@F.udf("double")
def multiply_slow(x):
return x * 1.2
df.withColumn("result", multiply_slow("amount"))
Правило 3: фильтруйте до UDF, не после¶
# ПЛОХО: UDF выполняется на всех строках
df.withColumn("score", expensive_ml_udf("features")) \
.filter("score > 0.8")
# ХОРОШО: сначала фильтр (SIMD SQL), потом UDF на меньшем объёме
df.filter("pre_score > 0.5 AND status = 'active'") \
.withColumn("score", expensive_ml_udf("features")) \
.filter("score > 0.8")
Правило 4: управляйте размером батча под вашу схему¶
# Формула: batch_size ≈ target_memory_mb × 1024 × 1024 / row_size_bytes
# row_size_bytes ≈ sum(col_sizes) ≈ num_columns × avg_col_size
# Для широких схем (100+ колонок × 8 байт = 800 байт/строка):
# batch_size = 100 MB / 800 B = 125_000 → установить 100_000
# Для узких схем (5 колонок × 8 байт = 40 байт/строка):
# batch_size = 100 MB / 40 B = 2_500_000 → установить 1_000_000
# Практически: начните с 10_000, увеличивайте пока нет OOM на воркере
Правило 5: используйте ArrowDtype в Pandas 2.x¶
import pandas as pd
import pyarrow as pa
# Pandas 2.x: явно просить Arrow backend
df_pandas = df_spark.toPandas()
# Конвертировать в Arrow-backed dtypes для следующего шага:
df_arrow_backed = df_pandas.convert_dtypes(dtype_backend="pyarrow")
# Теперь все колонки - ArrowDtype, нет object-колонок
print(df_arrow_backed.dtypes)
# user_id int64[pyarrow]
# amount double[pyarrow]
# status large_string[pyarrow]
# Передача в DuckDB: zero-copy (ArrowDtype → Arrow buffer → DuckDB)
import duckdb
con = duckdb.connect()
con.register("data", df_arrow_backed)
result = con.execute("SELECT AVG(amount) FROM data WHERE amount > 100").df()
Антипаттерны Arrow¶
Антипаттерн 1: маленькие батчи в цикле¶
# ПЛОХО: создаём много маленьких DataFrame и конвертируем по одному
for batch in get_data_chunks(chunk_size=100): # 100 строк за раз!
df_chunk = spark.createDataFrame(batch) # Overhead на каждое создание
process(df_chunk)
# ХОРОШО: создаём один большой DataFrame
all_data = list(get_data_chunks(chunk_size=100))
df_all = spark.createDataFrame(pd.DataFrame(all_data))
process(df_all)
Антипаттерн 2: игнорировать предупреждения о fallback¶
WARN ArrowUtils: Failed to use Arrow optimization for toPandas. ...
Это предупреждение означает, что ваш запрос работает без Arrow (через Pickle). Если видите это - исправьте схему или тип данных. Распространённые причины: MapType, ArrayType(ArrayType(...)), DecimalType при старых версиях Arrow.
Антипаттерн 3: Python-циклы внутри pandas_udf¶
# ПЛОХО: Python цикл внутри pandas_udf убивает смысл батчевой обработки
@pandas_udf("double")
def slow_inner(s: pd.Series) -> pd.Series:
result = []
for value in s: # ЦИКЛ! Не векторизованно!
result.append(value * 1.2 if value > 0 else 0.0)
return pd.Series(result)
# ХОРОШО: NumPy/Pandas векторизованные операции
@pandas_udf("double")
def fast_inner(s: pd.Series) -> pd.Series:
return (s * 1.2).clip(lower=0.0) # Одна numpy операция над массивом
Домашнее задание¶
Задача 1: Бенчмарк Arrow vs Pickle¶
Создайте Spark DataFrame с 5 миллионами строк и следующей схемой: user_id (long), revenue (double), category (string), status (string), created_at (timestamp).
Измерьте время toPandas() с Arrow включённым и выключенным. Дополнительно:
- Используйте
psutilдля замера потребления памяти процесса Python Driver до и после операции - Запустите с
maxRecordsPerBatch = 1000,10000,100000и зафиксируйте разницу
Задача 2: Инспекция Record Batch¶
Напишите PySpark pandas_udf, которая:
- Внутри функции конвертирует
pd.Seriesвpa.Arrayчерезpa.Array.from_pandas() - Логирует: тип Arrow array, адрес и размер Values Buffer, количество NULL-значений, выровнен ли буфер по 64 байтам
- Выполняет вычисление и возвращает результат
Запустите и проверьте логи - убедитесь, что видите Arrow-буферы.
Задача 3: Оптимальный размер батча¶
Реализуйте benchmark из практики (разные maxRecordsPerBatch). Запустите для схемы с 5 колонками и 20 колонками. Постройте таблицу «batch_size → время → использование памяти».
На основе результатов обоснуйте: какой maxRecordsPerBatch оптимален для вашего кейса и почему?
Задача 4: Interoperability со сторонней библиотекой¶
Реализуйте pipeline:
- Spark читает Parquet (или сгенерированный датасет)
- Конвертирует в Arrow через
.toArrow()(Spark 3.3+) или через pandas - DuckDB или Polars выполняет дополнительный аналитический запрос (например, расчёт квантилей)
- Результат возвращается в Spark через
spark.createDataFrame(result_arrow)
Замерьте overhead на каждом переходе. Объясните, почему это быстрее, чем писать всё только в Spark или только в DuckDB/Polars.