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++

optimization

Проблема межпроцессного взаимодействия в 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 строк:

  1. JVM executor берёт строку из UnsafeRow-буфера
  2. Конвертирует её в Java-объекты (Integer, String, Double)
  3. Сериализует в байты через Pickle (Python-протокол сериализации объектов)
  4. Записывает в UNIX socket (локальный сетевой интерфейс)
  5. Python-воркер читает из сокета
  6. Десериализует из Pickle в Python-объекты (int, str, float)
  7. Вызывает Python UDF с этими объектами
  8. Сериализует результат обратно в Pickle
  9. Записывает в сокет → 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 векторного регистра.

Выровненный буфер гарантирует:

  1. VMOVDQA32 (aligned load) вместо VMOVDQU32 (unaligned load) - aligned быстрее на 10-15%
  2. Кэш-строка никогда не разрезает векторный 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 состоит из:

  1. FlatBuffer metadata: структурированное описание схемы батча, количество строк, смещения и длины каждого буфера - всего несколько килобайт
  2. 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, которая:

  1. Внутри функции конвертирует pd.Series в pa.Array через pa.Array.from_pandas()
  2. Логирует: тип Arrow array, адрес и размер Values Buffer, количество NULL-значений, выровнен ли буфер по 64 байтам
  3. Выполняет вычисление и возвращает результат

Запустите и проверьте логи - убедитесь, что видите Arrow-буферы.

Задача 3: Оптимальный размер батча

Реализуйте benchmark из практики (разные maxRecordsPerBatch). Запустите для схемы с 5 колонками и 20 колонками. Постройте таблицу «batch_size → время → использование памяти».

На основе результатов обоснуйте: какой maxRecordsPerBatch оптимален для вашего кейса и почему?

Задача 4: Interoperability со сторонней библиотекой

Реализуйте pipeline:

  1. Spark читает Parquet (или сгенерированный датасет)
  2. Конвертирует в Arrow через .toArrow() (Spark 3.3+) или через pandas
  3. DuckDB или Polars выполняет дополнительный аналитический запрос (например, расчёт квантилей)
  4. Результат возвращается в Spark через spark.createDataFrame(result_arrow)

Замерьте overhead на каждом переходе. Объясните, почему это быстрее, чем писать всё только в Spark или только в DuckDB/Polars.