Sail: Rust-нативный движок с полным Spark API
Sail как drop-in замена Apache Spark - архитектура на Rust + Arrow + DataFusion, Spark Connect protocol, бенчмарки TPC-H и сравнение с Comet/Velox.
Революция в аналитике: зачем переписывать Spark на Rust¶
Прежде чем изучать Sail, необходимо честно ответить на вопрос: зачем вообще что-то переписывать? Apache Spark - зрелый проект с огромной экосистемой, тысячами компаний в production, многолетней историей оптимизаций. Что принципиально не так с JVM?
Ответ лежит не в коде Spark, а в архитектурном фундаменте JVM и в том, как изменилось аппаратное обеспечение за последние 15 лет.
Эволюционный тупик JVM в эпоху NVMe и SIMD¶
Когда Spark создавался (2009–2012 гг.), главным узким местом в аналитике были диски и сеть. Данные читались с HDD со скоростью 100 MB/s, передавались по сети 1 Gbps, и время ожидания I/O полностью перекрывало любые накладные расходы JVM. GC-паузы в 200 мс никто не замечал на фоне минут ожидания чтения с диска.
К 2024 году картина перевернулась. Современные NVMe-диски дают 7 GB/s последовательного чтения. Сеть в облаке - 25–100 Gbps. Форматы Parquet и ORC сжаты настолько эффективно, что движок обрабатывает сжатые байты быстрее, чем раньше читал несжатые. Узким местом стал CPU - его способность обрабатывать данные в памяти.
И здесь JVM упирается в фундаментальные ограничения, которые нельзя устранить патчами:
Проблема 1: Java Object Graph и GC. Каждый Java-объект несёт заголовок (object header) 12–16 байт плюс выравнивание. Целое число Integer в Java занимает 16 байт вместо 4. Строка «hello» - 64 байта вместо 5. Когда вы обрабатываете DataFrame с миллиардами строк, сотни гигабайт уходят на метаданные объектов, а не на данные. Garbage Collector вынужден обходить граф из миллиардов объектов, создавая паузы от 100 мс до нескольких секунд. Project Tungsten в Spark 2.x частично решил это через Off-Heap (UnsafeRow в не-JVM памяти), но диспетчеризация задач, план выполнения и метаданные по-прежнему живут в heap.
Проблема 2: SIMD-барьер JIT. Современные CPU могут выполнять операции над 16 целыми числами одновременно (AVX-512). Однако JIT-компилятор Java не может надёжно генерировать SIMD-инструкции: проверки выхода за пределы массива (bounds checks), полиморфизм типов и архитектура JVM mешают авто-векторизации. Нативный код на Rust, скомпилированный заранее (AOT), использует SIMD гарантированно и предсказуемо.
Проблема 3: Накладные расходы при частичном переносе. Решения типа Comet (Apache) или Gluten+Velox переносят вычислительные операторы в нативный код, но оставляют координацию задач, shuffle management и диспетчеризацию на стороне JVM-драйвера. Это означает постоянный маршалинг данных между JVM heap и нативной памятью, сериализацию/десериализацию на границе, и то, что GC-паузы всё ещё влияют на планировщик задач.
Концепция Sail: монолитный нативный движок¶
Sail (разработан командой LakeSail) предлагает радикально иной подход: не ускорять Spark изнутри, а заменить весь JVM-рантайм целиком.
Sail - это самостоятельный распределённый вычислительный движок, написанный полностью на Rust. Он не содержит ни одной строки Java. Вместо этого он реализует тот же сетевой протокол (Spark Connect), который PySpark использует для общения с сервером - и клиентский код не замечает замены.
Архитектурная идея проста и элегантна: Spark API - это стандарт де-факто для аналитики. Миллионы строк PySpark-кода написаны инженерами по всему миру. Было бы расточительством заставлять их переписывать весь код под DuckDB, DataFusion или другой новый API. Sail сохраняет весь этот инвестированный капитал знаний и кода, заменяя только execution layer под капотом.
Экономика инфраструктуры: что даёт отсутствие JVM¶
Практическое следствие перехода на нативный движок ощущается немедленно в нескольких измеримых метриках:
Холодный старт (Cold Start). JVM нужно 3–10 секунд только на запуск: загрузка .class-файлов, инициализация JIT-компилятора, прогрев JIT (первые тысячи вызовов функций выполняются в интерпретируемом режиме). Rust-бинарник Sail стартует за 50–200 миллисекунд. В Kubernetes, где pod'ы запускаются под нагрузкой, это разница между задержкой 1 секунда и 10 секунд для первого запроса.
Потребление памяти. В типичном Spark-кластере 30–40% памяти «съедает» JVM инфраструктура: heap для объектов плана, netty-буферы, code cache для JIT, Spark UI метаданные. У Sail таких накладных расходов нет - почти вся память отведена под Arrow-буферы с данными.
Предсказуемость латентности. GC-паузы в JVM непредсказуемы. G1GC и ZGC минимизируют их, но не устраняют. В распределённой системе одна GC-пауза на любом executor превращается в задержку всего запроса (straggler effect). Rust не имеет GC - память освобождается детерминированно через систему владения (Ownership), и латентность предсказуема.
Архитектурный триумвират: Rust + Apache Arrow + DataFusion¶
Sail стоит на трёх столпах, каждый из которых решает конкретную техническую проблему. Понимание того, почему выбраны именно эти компоненты, критически важно для понимания всей архитектуры.
Rust как системный язык: Memory Safety без GC¶
Rust - язык системного программирования с уникальной системой типов, которая гарантирует безопасность памяти без сборщика мусора. Вместо GC Rust использует концепцию Ownership (владения): у каждого значения в памяти есть ровно один владелец, и когда владелец выходит из области видимости, память автоматически освобождается. Компилятор проверяет это статически - программа с утечкой памяти просто не скомпилируется.
Для аналитического движка это критически важно: операции над терабайтами данных не могут допускать утечек памяти или непредсказуемых пауз. В C++ это гарантируется только дисциплиной разработчика и инструментами типа Valgrind. В Java GC справляется, но с ценой в виде пауз. Rust гарантирует это на уровне компилятора.
Второй важный аспект Rust - zero-cost abstractions: высокоуровневые абстракции (итераторы, обобщённые типы, trait'ы) компилируются в тот же машинный код, что и написанный вручную низкоуровневый C. Движок может быть написан выразительно и безопасно, не жертвуя производительностью.
Наконец, Rust-код компилируется AOT (Ahead of Time) с полной оптимизацией LLVM. Компилятор знает целевую архитектуру CPU и генерирует AVX-512 или ARM Neon инструкции статически, без JIT-прогрева.
Apache Arrow как единый формат данных¶
Arrow - это не библиотека и не файловый формат. Arrow - это стандарт для представления табличных данных в памяти: строгая спецификация того, как хранить массивы значений в буферах памяти с выравниванием на 64 байта.
В Sail Arrow выступает как единый internal data format на всех этапах: чтение из Parquet, выполнение операторов, shuffle между воркерами, возврат результатов в Python. Данные не конвертируются при передаче между стадиями - они остаются Arrow Record Batches с начала до конца.
Это кардинально отличается от классического Spark, где данные многократно конвертируются:
- Parquet → JVM Internal Row → UnsafeRow → JVM Internal Row → сериализация для shuffle → десериализация → UnsafeRow снова
- При Python UDF: JVM → Pickle → Python → Pickle → JVM (тот самый IPC overhead, который мы изучали в уроке про Arrow)
В Sail конвертация происходит один раз: Parquet → Arrow columnar buffer. Дальше всё - это операции над Arrow буферами без дополнительных конвертаций.
Apache DataFusion как query engine¶
DataFusion - это зрелый Rust-native SQL query engine, изначально созданный как часть проекта Arrow. Он реализует:
- Парсер SQL → AST
- Logical Plan оптимизатор (rule-based и cost-based, аналог Catalyst в Spark)
- Physical Plan с несколькими стратегиями join'а (hash join, sort-merge join, nested loop)
- Vectorized expression evaluation на Arrow Record Batches
- Volcano iterator model с поддержкой push-based execution
Sail использует DataFusion как internal query engine: принимает logical plan от Spark Connect (который он транслирует из Spark protobuf в DataFusion формат), передаёт через оптимизатор DataFusion и исполняет physical plan на Arrow данных.
DataFusion уже используется в production многими проектами: InfluxDB v3 IOx, Databricks Delta Lake Rust reader, Apache Comet, Ballista (распределённый DataFusion), GlareDB. Это не экспериментальный код - это зрелый проект Apache Foundation.
Сквозная векторизация: End-to-End без конвертаций¶
Самое важное следствие выбора Arrow + DataFusion - сквозная векторизация всего pipeline. Рассмотрим, что происходит при выполнении простого запроса SELECT SUM(amount) FROM orders WHERE status = 'completed':
В классическом Spark:
- Parquet reader декодирует колонку
statusв JVM строки → аллоцирует миллионы String объектов в heap - Filter operator проверяет каждую строку - неупорядоченный доступ к heap, cache misses
SUMaccumulator обновляется строка за строкой - нет SIMD- GC периодически останавливает всё для сбора мусора
В Sail:
- Parquet reader декодирует колонку
statusнапрямую в Arrow buffer - 64-байтные aligned блоки в памяти - DataFusion filter evaluation работает батчами по 8192 строк: VPCMPD (AVX-512) сравнивает 16 строк за одну инструкцию, результат - битовая маска
- VPCOMPRESSD собирает прошедшие фильтр элементы без ветвлений
SUMaccumulates via VPADDD/VADDPD - 8 double-precision значений за инструкцию- Нет GC - Arrow буферы освобождаются детерминированно после обработки батча
Теоретический speedup от одной только SIMD-оптимизации - 8–16× на арифметических операциях. На практике (с учётом I/O и shuffle) TPC-H бенчмарки показывают 3–5× суммарного ускорения.
На схеме хорошо видно разделение на Control Plane (gRPC управляющие сообщения между Driver и Workers) и Data Plane (Arrow Flight - высокопроизводительный бинарный транспорт для shuffle данных между воркерами). Это архитектурное разделение позволяет масштабировать каждый из слоёв независимо.
Spark Connect: как Sail притворяется Spark¶
Центральная идея Sail - совместимость с PySpark без изменения кода - была бы невозможна без протокола Spark Connect, введённого в Apache Spark 3.4.
История: почему появился Spark Connect¶
До Spark 3.4 PySpark работал принципиально иначе. Когда вы вызывали SparkSession.builder.getOrCreate(), в Python-процессе запускался embedded Java Gateway (Py4J), который создавал прямой RPC-туннель к JVM-объекту SparkSession. Вся работа выполнялась через прямые вызовы JVM-методов: Python-объект DataFrame был лишь thin wrapper над Java-объектом.
Это создавало несколько проблем:
- Python-клиент обязан был запускаться на той же машине, что и Spark Driver (или иметь к нему прямой сетевой доступ с Py4J)
- Любая ошибка в Python-коде могла дестабилизировать Driver JVM
- Невозможно было подключить несколько клиентов к одному Spark-кластеру без полной изоляции
- Нельзя было заменить исполнительный движок, не меняя клиент
Spark Connect решает все эти проблемы разом: он отделяет клиентский API от исполнительного движка через строго определённый gRPC-протокол с protobuf схемой.
Анатомия Spark Connect протокола¶
Spark Connect работает следующим образом. PySpark-клиент строит Logical Plan - дерево операторов (Scan → Filter → Aggregate → Sort), которое описывает, что нужно вычислить, но не как. Это дерево сериализуется в protobuf-формат и отправляется по gRPC на сервер.
Сервер получает protobuf-план, десериализует его, оптимизирует и выполняет. Результат возвращается клиенту батчами в формате Arrow IPC.
Ключевое слово: сервер. Spark Connect специфицирует только протокол, а не реализацию. Стандартная реализация - это Spark JVM Driver. Но Sail реализует тот же protobuf-контракт на Rust, притворяясь стандартным Spark-сервером.
┌─────────────────────────────────────────────────────────┐
│ PySpark Client (Python) │
│ │
│ df = spark.read.parquet("s3://bucket/orders/") │
│ .filter("status = 'completed'") │
│ .groupBy("category") │
│ .agg(sum("amount")) │
│ │
│ # Клиент строит дерево: │
│ # Aggregate(GroupBy=[category], Agg=[sum(amount)]) │
│ # └─ Filter(status = 'completed') │
│ # └─ ParquetScan(s3://bucket/orders/) │
│ │
│ # Сериализует в protobuf, отправляет через gRPC │
└────────────────────┬────────────────────────────────────┘
│ gRPC (protobuf Plan)
▼
┌─────────────────────────────────────────────────────────┐
│ Sail Server (Rust) │
│ │
│ Принимает тот же protobuf, что и Spark JVM Server │
│ Транслирует в DataFusion Logical Plan │
│ Оптимизирует и выполняет на Rust Workers │
│ │
│ # Клиент не знает, что под капотом не JVM │
└─────────────────────────────────────────────────────────┘
Zero Code Change: единственное изменение - строка подключения¶
Именно это делает Sail практически уникальным среди нативных движков. Вам не нужно переписывать код, изучать новый API, менять трансформации. Достаточно изменить одну строку - способ создания SparkSession:
# ──────────────────────────────────────────────────────────
# БЫЛО: стандартный PySpark с JVM-рантаймом
# ──────────────────────────────────────────────────────────
# Запускает Java Virtual Machine в том же процессе,
# инициализирует SparkContext, master="local[*]" означает
# локальный режим на всех ядрах CPU. JVM heap по умолчанию
# 1GB, нужно настраивать spark.driver.memory и т.д.
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.master("local[*]") \
.appName("MyApp") \
.config("spark.driver.memory", "4g") \
.config("spark.sql.shuffle.partitions", "200") \
.getOrCreate()
# ──────────────────────────────────────────────────────────
# СТАЛО: тот же PySpark-клиент, но подключается к Sail
# ──────────────────────────────────────────────────────────
# pyspark-client - лёгкая версия PySpark без JVM.
# Содержит только клиентский код для Spark Connect.
# Sail-сервер запущен отдельно (см. ниже).
# "sc://" - схема Spark Connect, отличается от "spark://"
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.remote("sc://localhost:50051") \
.getOrCreate()
# ──────────────────────────────────────────────────────────
# Весь остальной код - без единого изменения!
# ──────────────────────────────────────────────────────────
# Это работает ровно так же, как с JVM Spark:
df = spark.read.parquet("s3://my-bucket/orders/")
result = df.filter(df.status == "completed") \
.groupBy("category") \
.agg({"amount": "sum"}) \
.orderBy("sum(amount)", ascending=False)
result.show()
result.write.mode("overwrite").parquet("s3://my-bucket/result/")
Важная деталь: в примере используется pyspark-client (не полный pyspark). Полный пакет pyspark включает JVM runtime, Hadoop библиотеки, Py4J - всё это не нужно при использовании Sail. pyspark-client - это только Python-код клиента Spark Connect, без JVM зависимостей. Установка занимает секунды, не минуты.
Поддерживаемые версии протокола¶
Sail поддерживает протокол Spark Connect версий Spark 3.5.x и Spark 4.x. Это означает, что pyspark-client должен соответствовать одной из этих версий. Пакет pysail автоматически устанавливает совместимую версию pyspark-client.
Установка и первый запуск¶
Рассмотрим детально все способы запуска Sail - от простейшего локального до production Kubernetes.
Установка через pip¶
# ──────────────────────────────────────────────────────────
# Минимальная установка: Sail + лёгкий PySpark-клиент
# ──────────────────────────────────────────────────────────
# pysail - это Python-пакет, содержащий:
# 1) Rust-бинарник Sail (скомпилированный, ~50MB)
# 2) Python API для управления сервером из кода
# 3) CLI-утилиту "sail"
#
# pyspark-client - только клиентская часть PySpark (без JVM).
# Версия должна совпадать с протоколом Sail.
pip install "pysail==0.6.2"
pip install "pyspark-client==4.1.1"
Обратите внимание: pysail содержит предскомпилированный Rust-бинарник. Это означает, что он может не использовать все возможности вашего CPU. Если у вас современный CPU с AVX-512, есть способ получить максимальную производительность:
# ──────────────────────────────────────────────────────────
# Компиляция под конкретный CPU (максимальная производительность)
# ──────────────────────────────────────────────────────────
# RUSTFLAGS="-C target-cpu=native" говорит компилятору Rust
# использовать все инструкции текущего CPU:
# - AVX-512 на Intel Ice Lake / Sapphire Rapids
# - AVX2 на AMD Zen3/4 и Intel Skylake+
# - ARM Neon на Apple Silicon / AWS Graviton
#
# --no-binary pysail означает: не скачивать wheel, а компилировать
# исходный код. Требует установленного Rust toolchain.
# Компиляция занимает 5-15 минут.
RUSTFLAGS="-C target-cpu=native" pip install pysail --no-binary pysail -v
# Проверка поддерживаемых инструкций CPU:
python -c "import pysail; print(pysail.build_info())"
Три способа запустить Sail¶
Способ 1: Embedded сервер в Python-коде (рекомендуется для разработки)
Это наиболее удобный способ для разработки и тестирования. Sail-сервер стартует внутри того же Python-процесса (в отдельном потоке) и автоматически завершается при выходе из программы:
# ──────────────────────────────────────────────────────────
# Embedded режим: сервер запускается из Python-кода
# ──────────────────────────────────────────────────────────
from pysail.spark import SparkConnectServer
from pyspark.sql import SparkSession
# SparkConnectServer запускает Rust Sail runtime в фоновом потоке.
# port=50051 - стандартный порт Spark Connect.
# background=True - не блокирует выполнение Python-кода.
server = SparkConnectServer(port=50051)
server.start(background=True)
# Подключаемся через стандартный Spark Connect клиент.
# "sc://" - схема Spark Connect (не путать с "spark://")
spark = SparkSession.builder \
.remote("sc://localhost:50051") \
.getOrCreate()
# Проверяем работоспособность: простой SQL-запрос
result = spark.sql("SELECT 1 + 1 AS two")
result.show()
# +---+
# |two|
# +---+
# | 2|
# +---+
# Работа с данными:
df = spark.range(1000000) # создаём DataFrame из 1M строк
print(f"Count: {df.count()}") # Count: 1000000
# Закрытие сессии
spark.stop()
server.stop()
Способ 2: CLI-сервер (для production и shared доступа)
# ──────────────────────────────────────────────────────────
# CLI: запуск Sail как отдельного сервиса
# ──────────────────────────────────────────────────────────
# RUST_LOG=info - уровень логирования (debug, info, warn, error)
# sail spark server - запускает Spark Connect сервер
# --port 50051 - порт gRPC (стандарт Spark Connect)
RUST_LOG=info sail spark server --port 50051
# Вывод при успешном старте:
# [INFO sail_spark_api] Spark Connect server listening on 0.0.0.0:50051
# [INFO sail_spark_api] Server started in 142ms <-- нет JVM cold start!
# В другом терминале подключаемся:
# python my_spark_script.py
# (меняем только строку SparkSession.builder на .remote("sc://localhost:50051"))
Способ 3: Интерактивная оболочка (аналог pyspark REPL)
# ──────────────────────────────────────────────────────────
# REPL: интерактивная оболочка (аналог `pyspark` команды)
# ──────────────────────────────────────────────────────────
# sail spark shell запускает интерактивную Python-сессию
# с уже настроенным SparkSession, подключённым к Sail.
# Переменная `spark` доступна сразу.
sail spark shell
# Внутри REPL:
# >>> spark.sql("SELECT current_timestamp()").show()
# >>> df = spark.read.parquet("/path/to/data")
# >>> df.printSchema()
Проверка совместимости вашего кода¶
Перед миграцией реального проекта рекомендуется запустить анализатор совместимости:
# ──────────────────────────────────────────────────────────
# Проверка совместимости проекта с Sail
# ──────────────────────────────────────────────────────────
# compatibility_check анализирует все .py и .ipynb файлы
# в указанной директории и определяет:
# - Какие PySpark API используются
# - Какие из них поддерживаются в Sail полностью
# - Какие поддерживаются частично
# - Какие не поддерживаются (например, RDD API)
python -m pysail.examples.spark.compatibility_check ./my_project/
# Пример вывода:
# Analyzing 47 Python files...
# ✅ DataFrame API: 312 usages - fully supported
# ✅ SQL: 89 usages - fully supported
# ⚠️ pandas_udf: 12 usages - supported with limitations
# ❌ RDD: 3 usages - NOT supported (files: etl/legacy.py)
# ❌ Streaming: 5 usages - partial support
#
# Compatibility score: 87%
# Recommendation: Review 3 files with RDD usage before migration
Важное ограничение инструмента: он анализирует только Python-код, но не проверяет SQL-строки (содержимое spark.sql("...")) и не гарантирует поведенческую совместимость (одинаковые результаты при крайних случаях типа NULL handling, floating point precision и т.д.).
Великая битва нативных движков: Sail vs Comet vs Gluten/Velox¶
Чтобы правильно оценить место Sail в экосистеме, необходимо понять архитектурные различия между тремя основными подходами к нативному ускорению Spark.
Подход 1: Apache DataFusion Comet - нативный плагин к Spark¶
Apache Comet (часть Apache Arrow проекта) - это плагин-ускоритель, который встраивается в существующий Spark. Архитектурно он работает следующим образом:
- Spark JVM Driver планирует выполнение как обычно, используя Catalyst optimizer
- Physical план содержит обычные JVM-операторы (HashAggregate, BroadcastHashJoin и т.д.)
- Comet заменяет отдельные физические операторы нативными реализациями на Rust
- Данные передаются между JVM-операторами и нативными операторами Comet через JNI (Java Native Interface)
Сильные стороны Comet:
- Полная совместимость: работает с любым Spark-кластером, не требует изменения инфраструктуры
- RDD API работает (Comet не трогает RDD layer)
- Постепенное внедрение: включить Comet для одного запроса, не меняя остальное
- Зрелость: Apache project с активным сообществом
Слабые стороны Comet:
- JVM-координация остаётся: планировщик задач, shuffle manager, JVM GC - всё это по-прежнему работает
- JNI-граница: каждый переход между JVM и нативным кодом имеет оверхед (несколько микросекунд + переключение контекста)
- Не все операторы заменены нативными - часть плана по-прежнему выполняется на JVM
- Ограниченная эффективность SIMD: данные проходят через JNI, что нарушает непрерывность потока
Comet: JVM Spark + нативные операторы
────────────────────────────────────────
Spark Driver (JVM) → Catalyst → Physical Plan
│
JVM operator ────┤
Comet op (Rust)──┤ ← JNI граница между каждым оператором
JVM operator ────┤
Comet op (Rust)──┤
Подход 2: Gluten + Velox - C++ исполнитель в JVM-обёртке¶
Gluten (Intel) + Velox (Meta) - это более агрессивный подход: замена всех физических операторов нативными через единый нативный execution backend.
Velox - это нативный vectorized execution engine от Meta, написанный на C++. Он используется в Production в Meta для аналитики (Presto/Trino нативный backend) и обрабатывает экзабайты данных в год. Это не экспериментальный проект.
Gluten - это слой адаптации, который позволяет Spark использовать Velox как backend. Spark Driver по-прежнему работает на JVM, но физический план выполняется целиком через Gluten → Velox на C++.
Сильные стороны Gluten+Velox:
- Самый зрелый нативный подход с доказанной production-надёжностью (Meta масштаб)
- Охват операторов значительно шире, чем у Comet
- Активно поддерживается крупными компаниями (Intel, ByteDance, Alibaba)
Слабые стороны:
- JVM-координация всё ещё присутствует (Driver, shuffle management)
- Сложность развёртывания: нужно собирать C++ библиотеки под конкретный дистрибутив Linux
- C++: управление памятью сложнее, чем в Rust - потенциальные утечки в edge-cases
- Совместимость: не все Spark UDF и функции поддерживаются Velox
Подход 3: Sail - полное вытеснение JVM¶
Sail не ускоряет Spark - он заменяет его. JVM в Sail-деплойменте отсутствует полностью. Это радикальное архитектурное решение с принципиально иным набором компромиссов:
Классический Spark:
Python ──Py4J──► JVM Driver ──JVM Executors
(Py4J overhead) (GC pauses) (JVM memory model)
Comet:
Python ──Py4J──► JVM Driver ──JVM+Native Executors
(Py4J overhead) (GC pauses) (смешанное выполнение)
Sail:
Python ──gRPC──► Rust Driver ──Rust Workers
(Arrow IPC) (без GC) (полная SIMD)
Сравнительная таблица нативных движков¶
| Критерий | Sail | Comet | Gluten + Velox | RAPIDS cuDF |
|---|---|---|---|---|
| Язык | Rust (100%) | Rust + Java (JVM драйвер) | C++ + Java (JVM драйвер) | CUDA/C++ + Java |
| JVM в runtime | Отсутствует полностью | Драйвер + координация | Драйвер + координация | Драйвер + координация |
| API-совместимость | Spark Connect (Spark 3.5/4.x) | Полная (внутри Spark JVM) | Полная (внутри Spark JVM) | Полная (внутри Spark JVM) |
| RDD API | ❌ Не поддерживается | ✅ Полностью | ✅ Полностью | ✅ Полностью |
| Streaming | 🔄 Частично | ✅ (через Spark) | ✅ (через Spark) | ✅ (через Spark) |
| GPU поддержка | ❌ | ❌ | ❌ | ✅ Основная фича |
| Shuffle | Нативный Arrow Flight | JVM Spark Shuffle | JVM Spark Shuffle | JVM Spark Shuffle |
| Cold Start | 100–200 мс | 10–30 сек (JVM) | 10–30 сек (JVM) | 10–30 сек (JVM) |
| Memory model | Rust Ownership (без GC) | JVM GC + native | JVM GC + C++ manual | JVM GC + GPU |
| SIMD | Полная (AVX-512 / ARM Neon) | Частичная (через JNI) | Частичная (через JNI) | GPU tensor ops |
| Матurity | Alpha/Beta (v0.6.x) | Beta (Apache project) | Beta/GA (Meta backing) | GA (NVIDIA) |
| Лучший кейс | Новые проекты, K8s, serverless | Ускорение существующего Spark | Enterprise, heavy SQL | ML + wide tables + GPU |
Когда выбирать каждый из движков¶
Выбирайте Sail, если:
- Новый проект без legacy RDD-кода
- Kubernetes / serverless, где JVM cold start критичен
- Важна стоимость инфраструктуры - 50-80% экономии памяти
- Команда готова к pre-1.0 статусу
Выбирайте Comet, если:
- Нужно ускорить существующий Spark без изменений инфраструктуры
- Часть кода использует RDD API
- Хотите Apache-backed проект с широким community
Выбирайте Gluten+Velox, если:
- Enterprise среда с зрелыми требованиями стабильности
- Meta/ByteDance масштаб нагрузок
- Команда имеет опыт с C++ toolchain и сборкой нативного кода
Выбирайте RAPIDS cuDF, если:
- GPU уже есть в инфраструктуре (NVIDIA A100/H100)
- Работаете с ML-пайплайнами (XGBoost, RAPIDS ML)
- Очень широкие таблицы (тысячи колонок)
Анализ производительности: Бенчмарки TPC-H¶
Что такое TPC-H и почему он важен¶
TPC-H (Transaction Processing Performance Council - Benchmark H) - это стандартный индустриальный бенчмарк для аналитических баз данных, созданный в 1999 году и до сих пор широко используемый. Он описывает схему данных виртуального поставщика деталей и 22 конкретных SQL-запроса, имитирующих реальные бизнес-запросы.
Запросы TPC-H намеренно спроектированы так, чтобы тестировать самые сложные аналитические паттерны:
- Множественные JOIN'ы (иногда 5–6 таблиц в одном запросе)
- Вложенные подзапросы (correlated subqueries)
- Оконные функции (window functions)
- GROUP BY с большим числом уникальных ключей
- Агрегации с условиями (CASE WHEN в SUM)
Scale Factor (SF) определяет размер данных: SF=1 ≈ 1 GB, SF=100 ≈ 100 GB, SF=1000 ≈ 1 TB. Для серьёзных сравнений используют SF=100 и выше.
Результаты TPC-H SF=100 (Sail vs Apache Spark)¶
Стенд для сравнения:
- Инфраструктура: AWS, Parquet-файлы в S3
- Инстанс: r8g.4xlarge (16 vCPU, 128 GB RAM, ARM Graviton4)
- Apache Spark: 3.5.1, spark.sql.shuffle.partitions=200
- Sail: v0.6.2, embedded mode
- Данные: 22 стандартных TPC-H запроса
| Метрика | Apache Spark | Sail | Разница |
|---|---|---|---|
| Суммарное время (22 запроса) | 387 секунд | 103 секунды | 3.8× быстрее |
| Пиковое потребление RAM | 54 GB | 22 GB | −59% памяти |
| Shuffle Spill to Disk | 110+ GB | 0 GB | нет spill |
| Расчётная стоимость инстанса | ~$1.82 | ~$0.48 | −74% стоимости |
Однако важно понимать, что разница не одинакова для всех запросов. Рассмотрим несколько показательных случаев:
Запрос Q1 (Shipping Priority, простая агрегация): Spark 8.2 сек → Sail 2.1 сек (3.9×). Запрос делает GROUP BY по двум колонкам и несколько SUM/COUNT. Здесь Sail выигрывает за счёт SIMD-агрегации и отсутствия JVM object overhead.
Запрос Q8 (Market Share, 6 JOIN'ов + 2 подзапроса): Spark 47.3 сек → Sail 6.1 сек (7.8×). Чем сложнее запрос и больше JOIN'ов, тем сильнее выигрывает нативный движок - меньше промежуточной сериализации, эффективнее hash-таблицы для JOIN.
Запрос Q21 (Suppliers who kept orders waiting): Spark 68.1 сек → Sail 9.3 сек (7.3×). Corelated subquery с большим количеством промежуточных результатов - здесь отсутствие GC-пауз критично.
Почему нет Shuffle Spill у Sail¶
Один из наиболее удивительных результатов - полное отсутствие spill на диск у Sail при тех же данных, что у Spark вызывают 110 GB spill. Объяснение многоуровневое:
Arrow columnar format экономит память. UnsafeRow в Spark хранит строки с фиксированными полями + variadic секция. Arrow columnar format хранит каждую колонку отдельным непрерывным буфером. При обработке GROUP BY агрегации нужны только ключевые колонки - Sail читает только их, остальные колонки даже не загружаются в память для промежуточного состояния.
Hash-таблицы для JOIN работают эффективнее. DataFusion использует оптимизированные hash-таблицы, где ключ - это Arrow scalar value (aligned, fixed-size). В Spark hash-таблица хранит указатели на UnsafeRow-объекты, что добавляет косвенность и cache misses.
Нет JVM heap overhead. Spark резервирует heap memory не только для данных, но и для объектов планировщика, netty буферов, Spark UI метаданных, code cache JIT. Sail использует память почти исключительно для Arrow буферов.
CPU Efficiency: почему важна утилизация кэша¶
Одна из ключевых метрик нативных движков - L1/L2 cache hit rate. Данные, которые помещаются в кэш процессора (L1=32KB, L2=256KB, L3=8–32MB), обрабатываются в 10–60 раз быстрее, чем данные в RAM.
Arrow columnar layout помогает здесь радикально: когда вы суммируете колонку amount (INT64), данные лежат непрерывно в памяти: 8 байт, 8 байт, 8 байт, ... CPU prefetcher предсказывает доступ и загружает данные в кэш заранее. Cache hit rate → 95%+.
В JVM строковой модели данные (даже в UnsafeRow) имеют переменный размер и могут быть scattered по памяти. Prefetcher не справляется с непредсказуемым паттерном доступа. Cache hit rate падает до 60–70%.
Глубокое погружение: как работает выполнение запроса в Sail¶
От PySpark-кода до Arrow Record Batches¶
Рассмотрим детальный путь выполнения запроса шаг за шагом:
# Этот запрос мы будем трассировать
df = spark.read.parquet("s3://bucket/orders/") \
.filter("amount > 1000") \
.groupBy("region") \
.agg({"amount": "sum", "order_id": "count"})
result = df.collect()
Шаг 1: Построение Logical Plan на стороне клиента
PySpark-клиент не выполняет никаких вычислений. Он строит дерево операторов в памяти:
Project [region, sum(amount), count(order_id)]
└─ Aggregate(groupingExpressions=[region],
aggregateExpressions=[sum(amount), count(order_id)])
└─ Filter (amount > 1000)
└─ Relation [order_id, region, amount, status]
└─ ParquetScan(s3://bucket/orders/)
Это дерево сериализуется в protobuf (около 500 байт для данного примера) и отправляется по gRPC.
Шаг 2: Трансляция в DataFusion Logical Plan
Sail получает protobuf и транслирует его в DataFusion Logical Plan. Это не тривиальный маппинг - Spark и DataFusion имеют разные системы типов, разные правила NULL-обработки, разные семантики некоторых функций. Sail содержит слой совместимости, который обрабатывает эти различия.
Шаг 3: Оптимизация
DataFusion Optimizer применяет правила оптимизации:
- Predicate Pushdown: фильтр
amount > 1000проталкивается в ParquetScan - Parquet reader использует min/max statistics из row group footer, чтобы пропустить целые row groups без чтения данных - Column Pruning: читаются только колонки
region,amount,order_id-statusне нужна для запроса - Aggregate Simplification: COUNT(order_id) с NOT NULL полем упрощается
- Join Reordering: (для запросов с JOIN'ами) меньшая таблица строится как hash-таблица
Шаг 4: Разбивка на Stages по Shuffle границам
GROUP BY требует shuffle: все строки с одинаковым region должны попасть на один воркер. DataFusion physical plan разбивается на два stages:
Stage 2: Final Aggregate
RepartitionMerge [region]
Stage 1: Partial Aggregate + Shuffle
HashAggregate (partial)
Filter (amount > 1000)
ParquetScan
Stage 1 выполняется на всех воркерах параллельно (каждый читает свои Parquet-файлы). После частичной агрегации данные отправляются через Arrow Flight shuffle. Stage 2 делает финальную агрегацию.
Шаг 5: Vectorized Execution на воркерах
На каждом воркере physical plan выполняется в vectorized (batched) режиме. Вместо Volcano Model (один row за раз), DataFusion использует push-based vectorized model: каждый оператор получает батч из 8192 строк (Arrow Record Batch) и возвращает результирующий батч.
ParquetScan:
Читает row group (целиком в память как Arrow batch)
Применяет predicate pushdown (пропускает row groups по statistics)
Возвращает Record Batch [order_id, region, amount]
Filter (amount > 1000):
Получает Record Batch 8192 строк
VPCMPD: сравнивает 16 значений amount за 1 инструкцию
Результат: bit mask (например, 0b1101011010110101)
VPCOMPRESSD: собирает прошедшие фильтр строки без ветвлений
Возвращает отфильтрованный Record Batch
HashAggregate (partial):
Группирует по region (строка → hash → bucket)
Аккумулирует SUM и COUNT используя VPADDD/VADDPD
Возвращает частичные агрегаты: [(region, partial_sum, partial_count)]
Шаг 6: Arrow Flight Shuffle
Частичные агрегаты отправляются через Arrow Flight - высокопроизводительный транспортный протокол поверх gRPC, специально оптимизированный для Arrow данных. Arrow IPC формат (FlatBuffer метаданные + binary data) отправляется без дополнительной сериализации: данные уже в Arrow формате, и их можно отправить с нулевым копированием.
Шаг 7: Финальная агрегация и возврат результата
Финальный Stage объединяет частичные агрегаты (SUM складывается, COUNT складывается), применяет сортировку если нужно, и возвращает результат в виде Arrow IPC stream обратно через gRPC → PySpark-клиент.
PySpark-клиент получает Arrow данные и преобразует в Python-объекты только в момент .collect() - и только тогда, когда данные нужны Python-коду.
Управление памятью в Rust: детерминированное освобождение¶
Понимание Rust-механизма управления памятью важно для понимания, почему Sail предсказуем под нагрузкой.
Ownership и Drop¶
В Rust каждый буфер памяти (Arrow buffer, hash-таблица для JOIN, промежуточные агрегаты) является owned value. Когда оператор завершает обработку батча и переменная-владелец выходит из области видимости, Rust вызывает метод drop немедленно и детерминированно - не при следующем GC-цикле, а сразу.
// Псевдокод DataFusion оператора (упрощённо)
fn execute_batch(input: RecordBatch) -> RecordBatch {
// input - это owned value
// RecordBatch содержит Arrow буфер (обычно десятки MB)
let filtered = apply_filter(&input);
// input больше не нужен - память НЕМЕДЛЕННО освобождается
// здесь, не ждём GC
drop(input);
let aggregated = aggregate(filtered);
// filtered освобождается
// drop(filtered) вызывается автоматически при выходе из scope
aggregated
// caller получает ownership aggregated
}
В JVM эквивалентный код оставлял бы input и filtered в heap как dead objects. GC увидит их только на следующем minor collection (через 100 мс–2 сек). За это время в heap накапливаются гигабайты dead objects, давя на GC и приближая full collection (stop-the-world).
Reference Counting для shared buffers¶
Arrow буферы часто разделяются без копирования (zero-copy слайсинг). Sail использует Arc (Atomic Reference Counting) для таких случаев: когда несколько операторов читают один буфер, они держат shared reference, и буфер освобождается только когда последний holder drops свою reference.
Это принципиально отличается от GC-трейсинга в JVM: нет необходимости сканировать граф объектов, нет пауз - просто атомарный счётчик.
Практика: развёртывание и миграция реального пайплайна¶
Кейс: аналитика Gold-слоя с тяжёлыми агрегациями¶
Рассмотрим реальный сценарий: пайплайн расчёта недельных метрик продаж по регионам и категориям с оконными функциями (running total, 4-week moving average).
# ──────────────────────────────────────────────────────────
# Конфигурация подключения к Sail
# ──────────────────────────────────────────────────────────
# В реальном проекте параметры выносятся в конфиг/env-переменные.
# SAIL_HOST, SAIL_PORT - из environment.
import os
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.window import Window
SAIL_HOST = os.environ.get("SAIL_HOST", "localhost")
SAIL_PORT = os.environ.get("SAIL_PORT", "50051")
spark = SparkSession.builder \
.remote(f"sc://{SAIL_HOST}:{SAIL_PORT}") \
.getOrCreate()
# ──────────────────────────────────────────────────────────
# Шаг 1: Чтение Bronze-слоя из MinIO (S3-compatible)
# ──────────────────────────────────────────────────────────
# MinIO настроен на том же хосте или доступен из Sail-воркеров.
# Sail поддерживает S3-compatible storage через aws-sdk-rust.
orders = spark.read.parquet("s3a://datalake/bronze/orders/")
products = spark.read.parquet("s3a://datalake/bronze/products/")
# ──────────────────────────────────────────────────────────
# Шаг 2: Трансформации Silver-слоя
# ──────────────────────────────────────────────────────────
# Фильтрация и обогащение данных.
# Все эти операции ленивые - физически не выполняются до действия.
orders_clean = orders \
.filter(F.col("status").isin(["completed", "shipped"])) \
.filter(F.col("amount") > 0) \
.withColumn("order_date", F.to_date("created_at")) \
.withColumn("week", F.date_trunc("week", "order_date"))
# JOIN с таблицей продуктов для получения категории
orders_enriched = orders_clean.join(
products.select("product_id", "category", "subcategory"),
on="product_id",
how="left"
)
# ──────────────────────────────────────────────────────────
# Шаг 3: Gold-слой - агрегация по неделям
# ──────────────────────────────────────────────────────────
weekly_sales = orders_enriched \
.groupBy("week", "region", "category") \
.agg(
F.sum("amount").alias("total_revenue"),
F.count("order_id").alias("order_count"),
F.countDistinct("customer_id").alias("unique_customers"),
F.avg("amount").alias("avg_order_value"),
F.percentile_approx("amount", 0.5).alias("median_order_value")
)
# ──────────────────────────────────────────────────────────
# Шаг 4: Оконные функции - running total и скользящее среднее
# ──────────────────────────────────────────────────────────
# Window упорядочен по неделе внутри каждого (region, category).
# Это одна из самых тяжёлых операций - требует сортировки и shuffle.
w = Window \
.partitionBy("region", "category") \
.orderBy("week") \
.rowsBetween(Window.unboundedPreceding, Window.currentRow)
w_moving = Window \
.partitionBy("region", "category") \
.orderBy("week") \
.rowsBetween(-3, 0) # 4 недели: текущая + 3 предыдущих
result = weekly_sales \
.withColumn("running_total", F.sum("total_revenue").over(w)) \
.withColumn("moving_avg_4w", F.avg("total_revenue").over(w_moving)) \
.withColumn("wow_growth",
(F.col("total_revenue") - F.lag("total_revenue", 1).over(w)) /
F.lag("total_revenue", 1).over(w) * 100
)
# ──────────────────────────────────────────────────────────
# Шаг 5: Запись результата в Gold-слой
# ──────────────────────────────────────────────────────────
result.write \
.mode("overwrite") \
.partitionBy("week") \
.parquet("s3a://datalake/gold/weekly_sales_metrics/")
print("Pipeline completed successfully")
spark.stop()
Сравнение explain() планов: Spark vs Sail¶
Один из важных инструментов диагностики - сравнение EXPLAIN планов. Это показывает, как каждый движок интерпретирует один и тот же запрос:
# ──────────────────────────────────────────────────────────
# Сравнение explain() планов
# ──────────────────────────────────────────────────────────
# Напишем простой запрос для сравнения
simple_query = orders \
.filter(F.col("amount") > 1000) \
.groupBy("region") \
.agg(F.sum("amount").alias("total"))
# Sail explain:
# В Sail explain() возвращает DataFusion physical plan
print("=== SAIL PLAN ===")
simple_query.explain(extended=True)
# Типичный вывод Sail:
# == Physical Plan ==
# AggregateExec: mode=FinalPartitioned, gby=[region@0], aggr=[SUM(amount)]
# CoalesceBatchesExec: target_batch_size=8192
# RepartitionExec: partitioning=Hash([region@0], 8), input_partitions=8
# AggregateExec: mode=Partial, gby=[region@0], aggr=[SUM(amount)]
# CoalesceBatchesExec: target_batch_size=8192
# FilterExec: amount@1 > 1000
# ParquetExec: file_groups={...}, projection=[region, amount],
# predicate=amount@1 > 1000,
# pruning_predicate=amount_max@0 > 1000
#
# Обратите внимание:
# - ParquetExec: "pruning_predicate" - statistics-based row group pruning
# - FilterExec: фильтр применяется ПОСЛЕ чтения (row-level, не row-group)
# - AggregateExec mode=Partial → Repartition → mode=FinalPartitioned
# Это двухфазная агрегация с shuffle между фазами
# - CoalesceBatchesExec: объединяет маленькие батчи для эффективности SIMD
Ключевые операторы в Sail explain:
ParquetExecсpruning_predicate- читает только row groups, где max(amount) > 1000AggregateExec mode=Partial- локальная пред-агрегация на каждом воркере до shuffleRepartitionExec- Arrow Flight shuffle с hash-партиционированиемCoalesceBatchesExec- объединяет мелкие Record Batches в батчи по 8192 строк для SIMD
Профилирование ресурсов: Sail vs Spark¶
# ──────────────────────────────────────────────────────────
# Замер потребления ресурсов (требует psutil)
# ──────────────────────────────────────────────────────────
import psutil
import time
import os
def get_process_memory_mb():
"""Возвращает пиковое RSS (Resident Set Size) текущего процесса."""
process = psutil.Process(os.getpid())
return process.memory_info().rss / 1024 / 1024
def benchmark_query(spark, query_fn, label):
"""Запускает запрос и замеряет время и память."""
mem_before = get_process_memory_mb()
start = time.time()
result = query_fn(spark)
result.count() # force materialization
elapsed = time.time() - start
mem_peak = get_process_memory_mb()
print(f"\n[{label}]")
print(f" Execution time: {elapsed:.2f}s")
print(f" Memory before: {mem_before:.0f} MB")
print(f" Memory peak: {mem_peak:.0f} MB")
print(f" Memory delta: +{mem_peak - mem_before:.0f} MB")
return elapsed, mem_peak
# Запрос для бенчмарка
def heavy_query(spark):
return spark.read.parquet("s3a://datalake/bronze/orders/") \
.filter("amount > 0") \
.groupBy("region", "category") \
.agg(
F.sum("amount").alias("revenue"),
F.count("*").alias("orders"),
F.approx_count_distinct("customer_id").alias("customers")
)
# Запускаем на Sail:
sail_time, sail_mem = benchmark_query(spark, heavy_query, "SAIL")
# Для сравнения запустите тот же код с:
# spark = SparkSession.builder.master("local[*]").getOrCreate()
# И сравните метрики
Kubernetes: production-развёртывание без JVM¶
Архитектура Sail on Kubernetes¶
Sail в Kubernetes работает значительно проще, чем Spark on K8s. В Spark on K8s Driver Pod запускает JVM, который через K8s API создаёт Executor Pod'ы - каждый с JVM, каждый с 3–10 секундами cold start. Sail же запускает лёгкие Rust Worker Pod'ы (~50MB бинарник без JVM runtime).
# ──────────────────────────────────────────────────────────
# sail-server-deployment.yaml
# ──────────────────────────────────────────────────────────
apiVersion: apps/v1
kind: Deployment
metadata:
name: sail-spark-server
namespace: data-platform
spec:
replicas: 1
selector:
matchLabels:
app: sail-spark-server
template:
metadata:
labels:
app: sail-spark-server
spec:
serviceAccountName: sail-sa
containers:
- name: sail
image: lakesail/sail:0.6.2
command: ["sail", "spark", "server", "--port", "50051"]
ports:
- containerPort: 50051
name: spark-connect
env:
# Режим Kubernetes - Sail создаёт Worker Pod'ы через K8s API
- name: SAIL_MODE
value: "kubernetes-cluster"
# Docker-образ для Worker Pod'ов
- name: SAIL_KUBERNETES__IMAGE
value: "lakesail/sail:0.6.2"
# Namespace для Worker Pod'ов
- name: SAIL_KUBERNETES__NAMESPACE
value: "data-platform"
# Таймаут сессии (для долгих ноутбуков)
- name: SAIL_SPARK__SESSION_TIMEOUT_SECS
value: "3600"
# Уровень логирования
- name: RUST_LOG
value: "info"
resources:
requests:
memory: "512Mi" # Driver сам по себе очень лёгкий
cpu: "500m"
limits:
memory: "2Gi"
cpu: "2"
---
# Service для доступа клиентов к Sail по Spark Connect
apiVersion: v1
kind: Service
metadata:
name: sail-spark-service
namespace: data-platform
spec:
selector:
app: sail-spark-server
ports:
- port: 50051
targetPort: 50051
name: spark-connect
type: ClusterIP
---
# RBAC: ServiceAccount + Role для создания Worker Pod'ов
apiVersion: v1
kind: ServiceAccount
metadata:
name: sail-sa
namespace: data-platform
---
apiVersion: rbac.authorization.k8s.io/v1
kind: Role
metadata:
name: sail-worker-role
namespace: data-platform
rules:
# Sail Driver должен создавать, удалять и следить за Worker Pod'ами
- apiGroups: [""]
resources: ["pods", "pods/log"]
verbs: ["get", "list", "watch", "create", "delete"]
---
apiVersion: rbac.authorization.k8s.io/v1
kind: RoleBinding
metadata:
name: sail-worker-rolebinding
namespace: data-platform
subjects:
- kind: ServiceAccount
name: sail-sa
roleRef:
kind: Role
name: sail-worker-role
apiGroup: rbac.authorization.k8s.io
Подключение из Airflow к Sail на Kubernetes¶
# ──────────────────────────────────────────────────────────
# Airflow DAG: запуск PySpark через Sail
# ──────────────────────────────────────────────────────────
# Sail доступен внутри кластера по адресу:
# sail-spark-service.data-platform.svc.cluster.local:50051
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime
def run_gold_pipeline(**context):
"""Запускает Gold-слой пайплайн через Sail."""
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
# Подключение к Sail через K8s service
# Sail сам создаст нужное количество Worker Pod'ов
spark = SparkSession.builder \
.remote("sc://sail-spark-service.data-platform.svc.cluster.local:50051") \
.getOrCreate()
try:
execution_date = context["ds"] # YYYY-MM-DD
result = spark.read.parquet(f"s3a://datalake/bronze/orders/dt={execution_date}/") \
.groupBy("region") \
.agg(F.sum("amount").alias("daily_revenue"))
result.write \
.mode("overwrite") \
.parquet(f"s3a://datalake/gold/daily_revenue/dt={execution_date}/")
finally:
spark.stop()
with DAG(
dag_id="gold_daily_metrics",
start_date=datetime(2024, 1, 1),
schedule="0 3 * * *",
catchup=False,
) as dag:
gold_task = PythonOperator(
task_id="run_gold_pipeline",
python_callable=run_gold_pipeline,
)
Обратите внимание: в этом паттерне Airflow Worker (Python-процесс) подключается к Sail-серверу через gRPC. Sail-сервер находится в K8s. Никакой JVM не запускается ни на Airflow Worker, ни где-либо ещё - только легковесный pyspark-client в Python.
Ограничения и зрелость проекта: честная оценка¶
Sail - это pre-1.0 проект. Версия 0.6.x означает, что API может меняться, некоторые функции находятся в разработке. Важно понимать конкретные ограничения перед принятием решения о миграции.
Что не поддерживается или поддерживается частично¶
RDD API - не поддерживается. Это фундаментальное ограничение, не план разработки. RDD (Resilient Distributed Dataset) - это JVM-конструкция, которая требует JVM для выполнения. Sail заменяет JVM - следовательно, RDD невозможен в принципе. Если в вашем коде есть sc.parallelize(), rdd.map(), rdd.reduceByKey() - Sail не подходит без переписывания на DataFrame API.
Structured Streaming - частичная поддержка. Базовые streaming операции работают, но не все sink'и, не все trigger-типы, нет complete mode для всех агрегаций. Для production streaming-пайплайнов использование Sail преждевременно.
Pandas on Spark (Koalas/pyspark.pandas) - в разработке. Pandas API on Spark - это большой слой совместимости, который требует много работы для полной поддержки.
Spark Configurations. Большинство конфигураций spark.* игнорируются, так как они специфичны для JVM-рантайма. Вместо них Sail использует переменные окружения SAIL_*. Это означает, что существующие скрипты с spark.conf.set("spark.sql.shuffle.partitions", "800") нужно адаптировать.
Python UDF без Arrow - медленно. Python UDF (не pandas_udf) по-прежнему требует сериализации данных для передачи в Python. Это общая проблема для всех non-JVM движков с Python-совместимостью. Решение - использовать pandas_udf с Arrow транспортом.
Spark UI и History Server. Отдельного Spark UI нет - привычный web-интерфейс с DAG визуализацией, stage progress, executor metrics недоступен. Мониторинг осуществляется через RUST_LOG и внешние инструменты (OpenTelemetry интеграция в разработке).
Матрица поддержки ключевых фич¶
| Функциональность | Статус | Комментарий |
|---|---|---|
| DataFrame SQL API | ✅ Полная | Основной приоритет проекта |
spark.sql(...) |
✅ Полная | Большинство SQL диалекта Spark |
| Parquet read/write | ✅ Полная | Включая predicate pushdown |
| Delta Lake | ✅ Полная | Включая Deletion Vectors |
| Apache Iceberg | ✅ Полная | Улучшенный partition pruning |
| CSV, JSON, Avro | ✅ Полная | Стандартные форматы |
pandas_udf |
✅ Поддерживается | Arrow транспорт, нет JVM |
| Python UDF | ⚠️ Медленно | Без Arrow - сериализация |
| Window Functions | ✅ Полная | Включая ROWS/RANGE |
| RDD API | ❌ Нет | Принципиально невозможно |
| Structured Streaming | 🔄 Частично | Базовый функционал |
| MLlib | ❌ Нет | JVM-зависимая библиотека |
| Spark UI | ❌ Нет | В roadmap |
| Pandas on Spark | 🔄 Разработка | |
| Spark ML Pipelines | ❌ Нет | JVM зависимость |
Стратегия миграции: постепенный переход¶
Рекомендуемый подход для команды, которая хочет попробовать Sail в production:
Шаг 1: Audit (1-2 дня)
python -m pysail.examples.spark.compatibility_check ./pipelines/
Определить процент совместимого кода
Найти RDD и Streaming использование
Шаг 2: Пилот (1-2 недели)
Выбрать 1-2 пайплайна без RDD, без Streaming
Запустить на тестовых данных на Sail
Сравнить результаты с Spark: spark_result.exceptAll(sail_result).count() == 0
Шаг 3: Shadow Mode (2-4 недели)
Запускать одновременно на Spark и Sail
Писать результаты в разные таблицы
Сравнивать результаты автоматически
Замерять время и память
Шаг 4: Production Cutover
Переключить трафик на Sail для протестированных пайплайнов
Мониторинг аномалий (результаты, latency)
Откат готов: просто изменить строку подключения обратно
Лабораторная практика: TPC-H мини-бенчмарк¶
Цель практики¶
Запустить несколько TPC-H запросов на локальном Sail и сравнить с PySpark. Даже на небольших данных (SF=1, ~1 GB) разница в архитектуре будет видна в профиле CPU и памяти.
Шаг 1: Подготовка данных TPC-H¶
# ──────────────────────────────────────────────────────────
# Генерация TPC-H данных SF=1 (около 1 GB)
# ──────────────────────────────────────────────────────────
# duckdb - удобный инструмент для генерации TPC-H данных
pip install duckdb
python3 - <<'EOF'
import duckdb
import os
os.makedirs("/tmp/tpch_parquet", exist_ok=True)
con = duckdb.connect()
# Загружаем TPC-H расширение и генерируем данные
con.execute("INSTALL tpch; LOAD tpch;")
con.execute("CALL dbgen(sf=1);") # SF=1 ~ 1 GB данных
# Экспортируем в Parquet
for table in ["lineitem", "orders", "customer", "part", "supplier",
"partsupp", "nation", "region"]:
con.execute(f"""
COPY {table}
TO '/tmp/tpch_parquet/{table}.parquet'
(FORMAT PARQUET)
""")
count = con.execute(f"SELECT COUNT(*) FROM {table}").fetchone()[0]
print(f"{table}: {count:,} rows")
con.close()
print("\nTPC-H data generated in /tmp/tpch_parquet/")
EOF
Шаг 2: Запуск бенчмарка на Sail¶
# ──────────────────────────────────────────────────────────
# tpch_benchmark.py - сравнение Sail vs PySpark
# ──────────────────────────────────────────────────────────
import time
import psutil
import os
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
def create_spark_session(use_sail: bool) -> SparkSession:
"""Создаёт SparkSession для Sail или классического Spark."""
if use_sail:
# Sail: запускаем embedded сервер
from pysail.spark import SparkConnectServer
server = SparkConnectServer(port=50051)
server.start(background=True)
return SparkSession.builder \
.remote("sc://localhost:50051") \
.getOrCreate(), server
else:
# Классический PySpark с JVM
return SparkSession.builder \
.master("local[*]") \
.appName("TPC-H Benchmark") \
.config("spark.driver.memory", "4g") \
.config("spark.sql.shuffle.partitions", "8") \
.getOrCreate(), None
def load_tpch_tables(spark, data_path="/tmp/tpch_parquet"):
"""Загружает TPC-H таблицы и регистрирует как temp views."""
tables = ["lineitem", "orders", "customer", "part",
"supplier", "partsupp", "nation", "region"]
for t in tables:
spark.read.parquet(f"{data_path}/{t}.parquet") \
.createOrReplaceTempView(t)
print(f"Loaded {len(tables)} TPC-H tables")
# TPC-H Q1: Pricing Summary Report
# Одна из простейших агрегаций: GROUP BY + 8 агрегатных функций
QUERY_Q1 = """
SELECT
l_returnflag,
l_linestatus,
SUM(l_quantity) as sum_qty,
SUM(l_extendedprice) as sum_base_price,
SUM(l_extendedprice * (1 - l_discount)) as sum_disc_price,
SUM(l_extendedprice * (1 - l_discount) * (1 + l_tax)) as sum_charge,
AVG(l_quantity) as avg_qty,
AVG(l_extendedprice) as avg_price,
AVG(l_discount) as avg_disc,
COUNT(*) as count_order
FROM lineitem
WHERE l_shipdate <= date '1998-09-02'
GROUP BY l_returnflag, l_linestatus
ORDER BY l_returnflag, l_linestatus
"""
# TPC-H Q3: Shipping Priority
# JOIN трёх таблиц + GROUP BY + ORDER BY
QUERY_Q3 = """
SELECT
l_orderkey,
SUM(l_extendedprice * (1 - l_discount)) as revenue,
o_orderdate,
o_shippriority
FROM customer, orders, lineitem
WHERE c_mktsegment = 'BUILDING'
AND c_custkey = o_custkey
AND l_orderkey = o_orderkey
AND o_orderdate < date '1995-03-22'
AND l_shipdate > date '1995-03-22'
GROUP BY l_orderkey, o_orderdate, o_shippriority
ORDER BY revenue DESC, o_orderdate
LIMIT 10
"""
# TPC-H Q5: Local Supplier Volume
# JOIN 6 таблиц - тяжёлый запрос
QUERY_Q5 = """
SELECT
n_name,
SUM(l_extendedprice * (1 - l_discount)) as revenue
FROM customer, orders, lineitem, supplier, nation, region
WHERE c_custkey = o_custkey
AND l_orderkey = o_orderkey
AND l_suppkey = s_suppkey
AND c_nationkey = s_nationkey
AND s_nationkey = n_nationkey
AND n_regionkey = r_regionkey
AND r_name = 'ASIA'
AND o_orderdate >= date '1994-01-01'
AND o_orderdate < date '1995-01-01'
GROUP BY n_name
ORDER BY revenue DESC
"""
def run_benchmark(spark, queries, label):
"""Запускает все запросы и возвращает времена выполнения."""
results = {}
for name, sql in queries.items():
# Прогрев: первый запрос может быть медленнее из-за cold start
if name == list(queries.keys())[0]:
spark.sql(sql).count()
# Реальный замер: 3 повторения, берём медиану
times = []
for _ in range(3):
start = time.time()
count = spark.sql(sql).count()
elapsed = time.time() - start
times.append(elapsed)
median_time = sorted(times)[1]
results[name] = median_time
print(f" [{label}] {name}: {median_time:.2f}s (rows: {count})")
return results
# Запуск бенчмарка
queries = {"Q1": QUERY_Q1, "Q3": QUERY_Q3, "Q5": QUERY_Q5}
print("=== SAIL BENCHMARK ===")
spark_sail, server = create_spark_session(use_sail=True)
load_tpch_tables(spark_sail)
sail_times = run_benchmark(spark_sail, queries, "SAIL")
spark_sail.stop()
if server:
server.stop()
print("\n=== SPARK JVM BENCHMARK ===")
spark_jvm, _ = create_spark_session(use_sail=False)
load_tpch_tables(spark_jvm)
jvm_times = run_benchmark(spark_jvm, queries, "SPARK")
spark_jvm.stop()
print("\n=== COMPARISON ===")
for name in queries:
speedup = jvm_times[name] / sail_times[name]
print(f" {name}: Sail {sail_times[name]:.2f}s, "
f"Spark {jvm_times[name]:.2f}s, "
f"Speedup: {speedup:.1f}x")
Шаг 3: Анализ результатов¶
После запуска бенчмарка проанализируйте не только время, но и профиль выполнения:
# ──────────────────────────────────────────────────────────
# Сравнение explain() планов для понимания разницы
# ──────────────────────────────────────────────────────────
# Запустите этот код на Sail:
spark = SparkSession.builder.remote("sc://localhost:50051").getOrCreate()
load_tpch_tables(spark)
# Q1 explain в Sail покажет DataFusion операторы:
# - ParquetExec с pruning_predicate (убирает ненужные row groups)
# - FilterExec vectorized (SIMD сравнение дат)
# - AggregateExec mode=Partial (пред-агрегация до shuffle)
# - RepartitionExec (Arrow Flight shuffle по ключам агрегации)
# - AggregateExec mode=FinalPartitioned (финальная агрегация)
# - SortExec (ORDER BY)
spark.sql(QUERY_Q1).explain(extended=True)
# Q5 explain покажет более сложный план с 6-way JOIN:
# Обратите внимание на порядок JOIN'ов - DataFusion выбирает
# оптимальный порядок на основе estimated row counts
spark.sql(QUERY_Q5).explain(extended=True)
spark.stop()
Куда движется индустрия: будущее Spark как протокола¶
Decoupling API от Execution Engine¶
Sail представляет собой конкретное воплощение более широкого тренда: разделение Spark API (как стандарта) и JVM-рантайма (как одной из возможных реализаций).
Spark Connect - это gRPC протокол с открытой спецификацией. Это означает, что любой желающий может реализовать совместимый сервер - не обязательно на JVM, не обязательно с DataFusion под капотом. Мы наблюдаем начало эры, когда «Spark» станет означать не конкретную JVM-программу, а стандарт для аналитических вычислений - как HTTP является стандартом для веб, а не конкретным веб-сервером.
Эту же логику прослеживает и сообщество: Databricks активно развивает Spark Connect в Spark 4.0, Apache DataFusion создаёт DataFusion Ballista как distributed Spark-compatible engine, стартапы типа Comet и Sail реализуют этот протокол.
GPU-аналитика и Arrow-native будущее¶
Arrow-native форматы данных отлично совместимы с GPU. NVIDIA RAPIDS cuDF, Apache Arrow Flight SQL и другие инструменты уже используют Arrow как bridge между CPU и GPU вычислениями. Arrow columnar layout - это именно та структура, которую GPU ожидает: непрерывные массивы значений, каждый столбец отдельным буфером.
Следующий логический шаг - нативный движок, который может прозрачно переключаться между CPU SIMD и GPU tensorcore в зависимости от типа операции: фильтрация на CPU (где latency важнее throughput), matrix multiplication на GPU (где throughput важнее).
Конвергенция форматов¶
Parquet, ORC, Arrow - три формата, которые всё больше сближаются. Arrow Flight SQL и Apache Iceberg REST Catalog создают стандартные интерфейсы не только для форматов данных, но и для каталогов и транспорта. Вероятное будущее: стандартный протокол (Spark Connect или Arrow Flight SQL) + стандартный формат (Iceberg + Arrow) + сменные execution engines (JVM Spark, Sail, DuckDB, Velox) - по аналогии с тем, как браузеры сменяют HTML/CSS рендереры, сохраняя совместимость со стандартом HTML.
Антипаттерны при переходе на нативные движки¶
Антипаттерн 1: Строчный мышление в колонковом мире¶
Самая частая ошибка разработчиков, переходящих с JVM Spark на нативные движки - продолжать думать в терминах строк:
# ❌ Антипаттерн: Python UDF со строковой логикой
# Даже в Sail этот UDF требует сериализацию row за row
# Arrow-преимущества теряются полностью
from pyspark.sql.types import DoubleType
from pyspark.sql.functions import udf
@udf(returnType=DoubleType())
def calculate_discount(price, qty, category):
"""Бизнес-логика скидки - Python UDF."""
if category == "electronics" and qty > 10:
return price * qty * 0.85
elif qty > 5:
return price * qty * 0.90
else:
return price * qty
# Применяем UDF - Sail вынужден сериализовать каждую строку в Python
result = df.withColumn("total", calculate_discount("price", "qty", "category"))
# ✅ Правильно: SQL-выражения работают полностью в нативном движке
# DataFusion выполняет это с SIMD-инструкциями без Python
result = df.withColumn("total",
F.when(
(F.col("category") == "electronics") & (F.col("qty") > 10),
F.col("price") * F.col("qty") * 0.85
).when(
F.col("qty") > 5,
F.col("price") * F.col("qty") * 0.90
).otherwise(
F.col("price") * F.col("qty")
)
)
Антипаттерн 2: Маленькие батчи в цикле¶
# ❌ Антипаттерн: сборка DataFrame из маленьких кусков в цикле
# Каждый collect() - отдельный запрос к движку
# Каждый createDataFrame() из маленького списка - overhead на округление
results = []
for date in date_range:
day_df = spark.read.parquet(f"s3a://bucket/events/dt={date}/")
day_result = day_df.agg(F.sum("amount")).collect()[0][0]
results.append({"date": date, "revenue": day_result})
final_df = spark.createDataFrame(results)
# ✅ Правильно: одна операция над всеми данными сразу
# Sail и DataFusion оптимизируют partition pruning автоматически
final_df = spark.read.parquet("s3a://bucket/events/") \
.filter(F.col("dt").between(start_date, end_date)) \
.groupBy("dt") \
.agg(F.sum("amount").alias("revenue")) \
.orderBy("dt")
Антипаттерн 3: Игнорирование Fallback предупреждений¶
# ⚠️ Если в RUST_LOG=warn вы видите сообщения типа:
# WARN sail_spark_api: Unsupported operator: RDDScanExec - falling back to...
# WARN sail_spark_api: Unsupported function: ml_predict - skipping native path
#
# Это означает: эта часть кода не выполняется нативно.
# Не игнорируйте - найдите замену или вынесите в отдельный пайплайн на Spark.
# Настройте логирование для мониторинга fallback:
import os
os.environ["RUST_LOG"] = "warn" # или "debug" для детальной диагностики
Антипаттерн 4: Перенос JVM-конфигурации без адаптации¶
# ❌ Антипаттерн: JVM-специфичные конфиги, которые Sail игнорирует
spark = SparkSession.builder \
.remote("sc://localhost:50051") \
.config("spark.driver.memory", "8g") # JVM heap - нет JVM в Sail
.config("spark.executor.memory", "16g") # JVM executor heap - нет
.config("spark.memory.fraction", "0.75") # JVM memory model - нет
.config("spark.memory.offHeap.enabled", "true") # Off-heap тюнинг - нет
.config("spark.sql.shuffle.partitions", "200") # Может игнорироваться
.getOrCreate()
# ✅ Правильно: используйте SAIL_* переменные окружения
# SAIL_SPARK__SESSION_TIMEOUT_SECS=3600
# SAIL_KUBERNETES__WORKER_MEMORY=16Gi
# RUST_LOG=info
Лучшие практики использования Sail¶
Правило 1: SQL и DataFrame API - всегда предпочтительнее Python UDF. В нативном движке выражения SQL компилируются в SIMD-инструкции. Python UDF требует сериализации данных. Разница - на порядок magnitude.
Правило 2: pandas_udf для неизбежной Python-логики. Если Python UDF нужен, используйте pandas_udf с Arrow транспортом. Данные передаются батчами (не строка за строкой), и хотя SIMD в Python-части нет, overhead сериализации минимален.
Правило 3: Фильтруйте и проецируйте как можно раньше. DataFusion делает predicate pushdown в Parquet (statistics-based row group pruning), но помогите оптимизатору явными фильтрами перед JOIN'ами. Меньше данных в памяти → лучше SIMD efficiency.
Правило 4: Проверяйте результаты при миграции. Нативные движки могут иметь незначительные отличия в поведении (floating point order of operations, NULL handling edge cases). Используйте df1.exceptAll(df2).count() == 0 для верификации.
Правило 5: Мониторьте через RUST_LOG. Sail не имеет Spark UI, но RUST_LOG=info/debug даёт детальную информацию о выполнении запросов, Arrow Flight shuffle, загрузке Parquet-файлов.
Домашнее задание¶
Задание 1: Установка и первый запрос (обязательное)¶
Установите pysail и pyspark-client, запустите embedded-сервер и выполните следующий DataFrame-запрос. Зафиксируйте время выполнения и объём использованной памяти:
# Запрос для выполнения на Sail:
# 1. Создать DataFrame из 10 миллионов строк с random данными
# 2. Применить фильтр, GROUP BY и несколько агрегатных функций
# 3. Записать результат в Parquet
# 4. Сравнить с тем же запросом на spark.master("local[*]")
Зафиксируйте: время запуска сервера, время выполнения запроса, пиковое RSS памяти через psutil. Прокомментируйте разницу.
Задание 2: TPC-H мини-бенчмарк (основное)¶
Используя код из раздела практики, запустите TPC-H запросы Q1, Q3 и Q5 на сгенерированных данных (SF=1). Составьте сравнительную таблицу: движок, запрос, время выполнения, строк в результате. Запустите explain() для Q5 на Sail и опишите своими словами: как DataFusion планирует 6-way JOIN, в каком порядке, почему именно так.
Задание 3: Проверка совместимости проекта (практическое)¶
Возьмите один из пайплайнов из предыдущих модулей курса (dbt-spark, PySpark Bronze/Silver слои). Запустите compatibility_check. Составьте список:
- Что совместимо полностью
- Что требует адаптации (какой именно)
- Что принципиально несовместимо и почему
Задание 4: Архитектурный вывод (аналитическое)¶
Вам дан производственный пайплайн: ежесуточная обработка 500 GB событий (Parquet, S3), 15 трансформаций Silver-слоя (JOIN'ы, агрегации, оконные функции), 5 Gold-витрин (тяжёлые GROUP BY, percentile, distinct). В коде нет RDD API, нет MLlib, нет Streaming. Текущая инфраструктура: AWS EMR, 3×r6g.4xlarge (16 CPU, 128 GB), стоимость ~$850/месяц при дневном режиме работы.
Напишите архитектурный вывод (1-2 страницы): целесообразно ли мигрировать на Sail? Какие риски? Какая ожидаемая экономия TCO? Какой план перехода? Что нужно проверить перед принятием решения?
Резюме¶
Sail - наиболее радикальный из всех нативных движков: он не ускоряет Spark изнутри (как Comet или Gluten), а заменяет весь JVM-рантайм целиком на Rust + Arrow + DataFusion.
Это даёт принципиально иной набор характеристик: отсутствие GC-пауз, детерминированная latency, мгновенный cold start (100 мс против 10 сек для JVM), 50-60% экономия памяти, полная SIMD-векторизация без JNI-барьеров. TPC-H SF=100 показывает 3.8× суммарного ускорения и 74% экономию стоимости инфраструктуры.
Ключевой компромисс: Sail - pre-1.0 проект без поддержки RDD API и Structured Streaming. Это не временное ограничение - это архитектурное следствие того, что JVM отсутствует. RDD - это JVM-конструкция, и в мире без JVM её нет.
Но если ваш пайплайн использует только DataFrame API и SQL (что верно для большинства современных аналитических задач), Sail открывает путь к радикальному снижению затрат и повышению производительности без переписывания бизнес-логики - достаточно изменить одну строку подключения.
Тренд очевиден: Spark Connect превращает «Spark API» из монолитной JVM-программы в открытый протокол с несколькими конкурирующими реализациями. Sail - первая production-направленная реализация этого протокола без JVM. За ней последуют другие.
Ссылки:
- GitHub: github.com/lakehq/sail
- Документация: docs.lakesail.com/sail/latest/
- Apache DataFusion: datafusion.apache.org
- Spark Connect Protocol: spark.apache.org/docs/latest/spark-connect-overview.html