Photon (Databricks): SIMD-движок на C++, что заменяет в Tungsten, 3–10× на Join и агрегациях
Архитектура Photon: почему Databricks переписали Spark на C++, как работают векторизованные операторы с AVX-512, fallback в JVM, и когда ускорение достигает 10×.
Рождение Photon: как создатели Spark отказались от своего движка¶
Paradox современной дата-инженерии: в 2020 году компания Databricks - та самая компания, которая создала Apache Spark - объявила, что написала новый execution engine на C++, заменяющий выполнение запросов в их управляемой Databricks Runtime. Движок называется Photon, и его появление стало признанием фундаментального ограничения JVM-архитектуры.
Почему это произошло именно сейчас, а не раньше? Ответ лежит в изменении аппаратной реальности облака.
Исторический контекст: когда I/O перестал быть узким местом¶
В 2012–2016 годах, когда Spark вытеснял MapReduce, главным узким местом аналитических задач было дисковое I/O. Данные лежали на HDD-дисках (100–200 MB/s последовательного чтения). Spark был быстрее MapReduce именно потому, что держал промежуточные результаты в памяти и избегал лишних записей на диск. На фоне I/O любые накладные расходы JVM - миллисекундные GC-паузы, object overhead - были незаметны.
К 2018–2020 годам ситуация изменилась радикально. NVMe-диски в облаке дают 3–7 GB/s. Сети в AWS/Azure/GCP - 25–100 Gbps. Форматы Parquet и Delta Lake сжаты настолько эффективно, что один NVMe-диск успевает прочитать и декомпрессировать гигабайты в секунду. Данные стали поступать в движок быстрее, чем JVM успевал их обрабатывать.
Узким местом стал CPU - его способность выполнять операции над данными в памяти. И здесь JVM упирается в фундаментальный потолок.
Project Tungsten: блестящая оптимизация с ограниченным потолком¶
Project Tungsten (Spark 2.0, 2016) был революционным шагом. Разработчики Spark осознали, что Java object model катастрофически неэффективна для аналитических данных:
- Java
Integerзанимает 16 байт вместо 4 (object header + padding + pointer) - Миллиарды Java-объектов создают давление на GC
- Случайный доступ к объектам в heap = cache misses
Tungsten решил эти проблемы через:
UnsafeRow - бинарный формат для строк данных в Off-Heap памяти. Вместо Java-объектов данные хранятся как raw bytes с фиксированными смещениями. Целое число INT32 занимает ровно 4 байта, как и должно.
Whole-Stage Code Generation (WSCG) - технология, при которой Catalyst генерирует специализированный Java-байткод для каждого конкретного запроса. Вместо виртуальных вызовов интерфейсных методов (next() в Volcano model) - один монолитный цикл без накладных расходов на индирекцию.
Off-Heap Memory - управление памятью вне JVM heap через sun.misc.Unsafe, устраняющее паузы GC.
Tungsten дал Spark огромный прирост производительности. Но к 2020 году достиг своего потолка по трём причинам:
Причина 1: JIT не может гарантировать SIMD. WSCG генерирует Java-байткод. JIT компилирует его в нативный код. Теоретически JIT может авто-векторизовать циклы и генерировать AVX инструкции. На практике JIT делает это ненадёжно: bounds checks в каждом обращении к массиву, полиморфизм типов, ограничения JVM runtime - всё это мешает авто-векторизации. WSCG порождает огромные монолитные методы, которые JIT компилятор вынужден анализировать целиком - и зачастую отказывается от агрессивной оптимизации из-за сложности.
Причина 2: Строковая (row-based) парадигма. UnsafeRow - это всё ещё строковый формат. Данные для одной строки хранятся рядом: [id, amount, region, category, timestamp, ...]. Когда оператор SUM(amount) обходит миллион строк, он пропускает через кэш CPU гигабайты данных, из которых только amount ему нужен. Остальные колонки загружаются в кэш, вытесняют нужные данные и выбрасываются.
Причина 3: Overhead WSCG компиляции. Генерация, компиляция и прогрев JIT-кода добавляет latency к каждому запросу (100–500 мс). Для коротких запросов (Serverless, интерактивный SQL) это заметная часть времени выполнения.
Концепция Photon: написать всё заново¶
В 2019–2020 годах Databricks создали новую команду с задачей: написать vectorized execution engine на C++, который заменит JVM-выполнение в Databricks Runtime. Главные принципы:
- Columnar от рождения: данные хранятся и обрабатываются по колонкам, а не строкам
- SIMD как первый класс: операторы написаны с явным использованием intrinsics AVX-2/AVX-512, не полагаясь на авто-векторизацию
- Нет GC: вся память управляется вручную через C++ allocators
- AOT компиляция: код скомпилирован заранее под конкретную микроархитектуру CPU (Intel Xeon, AMD EPYC, ARM Graviton)
- Совместимость с Catalyst: Photon не заменяет оптимизатор - только execution layer
Первый production-релиз Photon вышел в Databricks Runtime 9.1 (2021). С тех пор он стал основным execution engine для SQL-аналитики на Databricks.
Photon против Tungsten: битва двух философий¶
Рассмотрим детально, что именно Photon заменяет в архитектуре Spark, а что остаётся неизменным.
Что Photon НЕ заменяет¶
Важно понимать границы Photon. Он не является новым Spark. Следующие компоненты остаются без изменений:
- Spark API (PySpark, Scala DataFrame, SQL) - ваш код работает без изменений
- Catalyst optimizer - весь конвейер оптимизации (Logical Plan → Physical Plan) остаётся на Scala/JVM
- Spark Scheduler - распределение задач по executors, retry-логика, стадии
- Shuffle - данные по-прежнему пишутся на диск/сеть между стадиями (хотя Photon оптимизирует сериализацию)
- Cluster Manager (YARN/Kubernetes/Databricks Cluster Manager) - управление ресурсами не меняется
Что Photon ЗАМЕНЯЕТ¶
- Физические операторы:
HashAggregateExec→PhotonHashAggregateExec,SortMergeJoinExec→PhotonSortMergeJoinExecи т.д. - Parquet reader: стандартный VectorizedParquetReader заменяется нативным C++ ридером
- Hash tables: для JOIN и GROUP BY - нативные C++ hash maps, оптимизированные под размеры кэша CPU
- Expression evaluation: вычисление выражений (
amount * 1.2 + tax,CASE WHEN,LIKE) выполняется в C++ с SIMD - Memory management для данных: Arrow-подобные columnar буферы в Off-Heap вместо UnsafeRow
На схеме видно: Catalyst и Scheduler - общий фундамент для обоих путей. Справа от развилки - принципиально разные реализации. Photon-путь полностью нативный, без JVM в execution loop.
Whole-Stage Code Generation vs Векторизованное выполнение¶
Это самое принципиальное различие между Tungsten и Photon. Рассмотрим на конкретном примере запроса SELECT region, SUM(amount) FROM orders WHERE amount > 1000 GROUP BY region.
Tungsten / WSCG подход:
Catalyst генерирует Java-байткод примерно такого вида (псевдокод):
// WSCG генерирует монолитный метод для целого pipeline
void processPartition(Iterator<UnsafeRow> input) {
// Инициализация hash map для GROUP BY
HashMap<byte[], long[]> aggregateMap = new HashMap<>();
while (input.hasNext()) { // Цикл по СТРОКАМ (row-at-a-time)
UnsafeRow row = input.next();
// Читаем поля из бинарного формата UnsafeRow
double amount = row.getDouble(1); // offset вычислен Catalyst
String region = row.getString(2);
// Фильтр amount > 1000
if (amount <= 1000.0) continue;
// Группировка: hash key
byte[] key = region.getBytes();
long[] agg = aggregateMap.getOrDefault(key, new long[]{0L, 0L});
// Аккумулируем SUM(amount)
agg[0] += Double.doubleToLongBits(amount);
agg[1]++; // count
aggregateMap.put(key, agg);
}
// Финализация - возврат результатов строками
for (Map.Entry<byte[], long[]> entry : aggregateMap.entrySet()) {
// Создаём UnsafeRow для каждого результата
emitRow(entry.getKey(), entry.getValue());
}
}
Проблемы этого кода:
- Цикл
while (input.hasNext())- одна строка за итерацию, нет возможности применить SIMD row.getDouble(1)- чтение с переменным offset из бинарного буфера, bounds check JVM при каждом вызовеHashMap<byte[], long[]>- JVM-объекты в heap, давление на GC- JIT должен компилировать ВЕСЬ этот монолитный метод целиком - трудно оптимизировать
Photon / Vectorized подход:
Photon обрабатывает те же данные принципиально иначе. Каждый оператор работает с батчем значений (обычно 1024–4096 строк), представленным как contiguous C++ массивы:
// Упрощённый псевдокод Photon-оператора Filter
// В реальности использует SIMD intrinsics напрямую
void PhotonFilterExec::process(ColumnarBatch* input, ColumnarBatch* output) {
// input->amount - указатель на contiguous double массив, 4096 значений
// Все 4096 значений лежат РЯДОМ в памяти - идеально для кэша CPU
const double* amounts = input->column<double>("amount");
int n = input->num_rows; // 4096
// Validity mask: какие строки прошли фильтр
// 64-битный integer = битовая маска для 64 строк
std::vector<uint64_t> filter_mask(n / 64 + 1);
// Vectorized comparison: amount > 1000.0
// Компилятор C++ с -O3 -march=native генерирует AVX-512:
// VCMPGTPD zmm0, zmm1, [memory] (сравнивает 8 double за такт)
for (int i = 0; i < n; i += 8) {
// Автовекторизация компилятором ИЛИ явные intrinsics
__m512d batch = _mm512_loadu_pd(amounts + i);
__m512d threshold = _mm512_set1_pd(1000.0);
__mmask8 cmp_result = _mm512_cmp_pd_mask(batch, threshold, _CMP_GT_OQ);
// cmp_result - 8-битная маска: 1 = прошёл фильтр, 0 = нет
filter_mask[i/64] |= ((uint64_t)cmp_result << (i % 64));
}
// Применяем маску: VPCOMPRESSQ собирает прошедшие элементы
int out_count = _mm512_compressstore_filter(
output->amount_col, amounts, filter_mask, n
);
output->set_valid_count(out_count);
}
Ключевые различия:
- Цикл работает над массивом из 4096 значений, а не строками по одной
- AVX-512 обрабатывает 8
double(64-битных) за одну инструкцию VCMPGTPD - Данные лежат contiguous в памяти - CPU prefetcher предсказывает доступ и загружает данные заранее
- Нет JVM bounds checks, нет GC, нет object allocation
Этот подход - vectorized execution model - пришёл из базы данных MonetDB и DuckDB. Он противоположен как строковому Volcano model, так и WSCG: операторы конвейерные, но каждый вызов обрабатывает векторы, а не строки.
Магия под капотом: SIMD-ускорение операций¶
AVX-512 в аналитике: чем больше регистр, тем быстрее¶
Современные серверные CPU в облаке (Intel Xeon Ice Lake/Sapphire Rapids, AMD EPYC Zen3/4) имеют AVX-512 инструкции. ZMM-регистры (512 бит) вмещают:
- 16 ×
float(32 бит) - 8 ×
double(64 бит) - 16 ×
int32 - 8 ×
int64 - 64 ×
int8(для bitmap операций)
Это означает, что одна инструкция VCMPGTPD сравнивает 8 пар double-значений с одним порогом. На одной итерации цикла фильтрации Photon проверяет 8 значений amount > 1000.0 одновременно.
Для колонки из 100 миллионов double-значений:
- JVM цикл: 100M итераций, каждая ≈ 3–5 ns → ~300–500 ms
- AVX-512 цикл: 100M / 8 = 12.5M итераций, каждая ≈ 1 ns → ~12.5 ms
- Теоретический speedup: 24–40× только на фильтрации
На практике speedup меньше из-за I/O и shuffle накладных расходов, но на CPU-bound запросах (тяжёлые JOIN, GROUP BY) Photon часто даёт 3–10×.
Vectorized Hash Tables: революция в JOIN¶
JOIN - это операция, где Spark традиционно больше всего теряет. SortMergeJoin требует сортировки обеих сторон. BroadcastHashJoin строит hash table из меньшей стороны и пробирует её для каждой строки большей стороны. Именно hash probe (поиск в таблице) - самая дорогая часть JOIN.
В JVM Spark hash probe выглядит так:
- Взять строку с левой стороны
- Вычислить hash ключа
- Найти bucket в hash map (объект Java)
- Сравнить ключи (String equals или similar)
- Если нашли - emit row
Проблема: каждый шаг - отдельная операция над одной строкой. Hash probe - это непредсказуемый доступ к памяти (random access по hash bucket), что вызывает cache misses. Cache miss на L3 → DRAM занимает 60–100 ns. При миллионах строк в JOIN - это секунды потерянного времени.
Photon решает это через vectorized hash probing:
// Vectorized hash probe - упрощённый псевдокод
void PhotonHashJoin::probe_batch(
ColumnarBatch* probe_side, // Левая сторона JOIN (батч)
NativeHashTable* hash_table, // Hash table из правой стороны
ColumnarBatch* output
) {
const int64_t* probe_keys = probe_side->column<int64_t>("order_id");
int n = probe_side->num_rows; // 4096 строк в батче
// Шаг 1: Вычисляем хеши для ВСЕХ ключей батча сразу
// SIMD вычисляет hash для 4 int64 за такт
uint32_t hashes[4096];
compute_hashes_vectorized(probe_keys, n, hashes);
// Шаг 2: Lookup buckets для всех ключей батча
// Prefetcher CPU заранее загружает bucket'ы в кэш
uint32_t bucket_ids[4096];
for (int i = 0; i < n; i++) {
bucket_ids[i] = hashes[i] & hash_table->capacity_mask;
// Prefetch следующих bucket'ов, пока обрабатываем текущие
__builtin_prefetch(hash_table->buckets + bucket_ids[i + 16], 0, 1);
}
// Шаг 3: Batch comparison - SIMD сравнивает ключи пачками
// Вместо одного ключа за раз - 8 ключей за такт (AVX-512)
__mmask8 match_mask;
for (int i = 0; i < n; i += 8) {
__m512i probe_batch = _mm512_loadu_si512(probe_keys + i);
__m512i stored_batch = gather_from_hash_table(hash_table, bucket_ids + i);
match_mask = _mm512_cmpeq_epi64_mask(probe_batch, stored_batch);
// match_mask: 8-битная маска совпавших ключей
emit_matched_rows(output, i, match_mask);
}
}
Три оптимизации в этом коде принципиально важны:
Software prefetching: строка __builtin_prefetch(hash_table->buckets + bucket_ids[i + 16], 0, 1) говорит CPU: «через 16 итераций нам понадобится вот этот адрес - начни загружать его из DRAM в кэш прямо сейчас». Это устраняет cache-miss stall на 60–100 ns на строку.
Векторизованное вычисление хешей: вместо одного хеша за раз - сразу для всего батча. SIMD делает это за 4096 / 4 = 1024 итерации вместо 4096.
SIMD batch comparison: _mm512_cmpeq_epi64_mask сравнивает 8 пар int64-ключей за одну инструкцию. Вместо 4096 сравнений - 512 SIMD-итераций.
Aggregation без тормозов: битовые маски для NULL¶
GROUP BY с агрегациями (SUM, COUNT, AVG) - второй главный бенефициар Photon. Рассмотрим GROUP BY region, SUM(amount):
// Vectorized aggregation - упрощённый псевдокод
void PhotonHashAggregate::process_batch(ColumnarBatch* batch) {
const char** regions = batch->column<char*>("region");
const double* amounts = batch->column<double>("amount");
const uint64_t* null_bitmap = batch->null_bitmap("amount");
int n = batch->num_rows;
// Шаг 1: Вычисляем group keys для всего батча
uint32_t group_ids[4096];
lookup_or_create_groups_vectorized(regions, n, group_ids);
// Шаг 2: Accumulate SUM(amount) с учётом NULL
// null_bitmap: бит 1 = NOT NULL, бит 0 = NULL
// SIMD обрабатывает 8 пар (group_id, amount) за такт
for (int i = 0; i < n; i += 8) {
// Загружаем 8 amount-значений
__m512d amounts_vec = _mm512_loadu_pd(amounts + i);
// Загружаем null-маску для 8 значений из bitmap
uint8_t null_bits = extract_null_bits(null_bitmap, i, 8);
__mmask8 not_null = (__mmask8)null_bits;
// VADDPD с маской: добавляем только NOT NULL значения
// Строки с NULL просто пропускаются через маску - без branch!
for (int j = 0; j < 8; j++) {
if (not_null & (1 << j)) {
agg_state[group_ids[i+j]].sum += amounts[i+j];
agg_state[group_ids[i+j]].count++;
}
}
// Photon использует AVX-512 scatter для прямой записи в agg_state
// scatter_accumulate(agg_state, group_ids + i, amounts_vec, not_null);
}
}
Критически важная оптимизация: обработка NULL через bitmasking вместо ветвлений. В JVM Spark каждый NULL требует if (value != null) - это branch (условный переход). Процессор пытается предсказать, будет ли значение NULL. Если данные непредсказуемы (50% NULL) - branch predictor ошибается половину времени, и каждая ошибка стоит 15–20 тактов на сброс pipeline.
Photon хранит NULL-информацию как bitmap (1 бит на строку) и использует маскированные SIMD-инструкции (_mm512_mask_add_pd и аналоги). NULL-строки просто игнорируются через маску без создания branch - branchless execution.
Жизненный цикл запроса: гибридное выполнение¶
Photon не обязан поддерживать 100% операторов Spark SQL - и он не поддерживает. Часть запросов выполняется нативно на C++, часть падает обратно в JVM-Tungsten. Это называется hybrid execution или adaptive fallback.
Как Photon выбирает, что выполнять нативно¶
Когда Catalyst строит физический план, Photon-слой анализирует каждый Physical Plan узел и решает: может ли он выполнить этот оператор нативно?
Поддерживаемые операторы (выполняются в Photon):
PhotonScanMicroBatchRDD/PhotonParquetScan- чтение Parquet, ORC, DeltaPhotonFilter- фильтрация с предикатами (column predicates, BETWEEN, IN)PhotonProject- вычисление выражений (арифметика, string functions, CASE WHEN)PhotonHashAggregate- GROUP BY с SUM, COUNT, AVG, MIN, MAXPhotonBroadcastHashJoin- JOIN с broadcast sidePhotonSortMergeJoin- JOIN с сортировкойPhotonSort- сортировка данныхPhotonWindowAggregate- оконные функции (RANK, ROW_NUMBER, SUM OVER)PhotonShuffleExchangeExec- сериализация shuffle данных (columnar format)
Неподдерживаемые операторы (fallback в JVM):
- Python UDF (row-at-a-time) - требует Python-процесса, разрывает native path
- Pandas UDF (через Arrow) - может работать с
ArrowEvalPython, но выходит из Photon - Сложные nested correlated subqueries - не все варианты поддержаны
GROUPING SETSс некоторыми конфигурациямиRDD API- вне Spark SQL полностью
Механизм Fallback: PhotonToRow и RowToPhoton¶
Когда в плане встречается оператор, который Photon не поддерживает, движок вставляет конвертационные узлы:
PhotonResultStage(илиPhotonToRow) - конвертирует Photon columnar batch обратно в JVM UnsafeRow. Это дорогая операция: нужно создать Java-объекты для каждого значения каждой строки.RowToPhoton- обратная конвертация: JVM UnsafeRow строки упаковываются в columnar C++ буферы.
Каждая граница между Photon и JVM - это overhead. Оптимизация, которую следует отслеживать: минимизировать количество таких переходов в плане.
Пример плана с fallback из-за Python UDF:
PhotonResultStage ← конвертация C++ → JVM строки
PhotonHashAggregate ← нативная агрегация (Photon C++)
RowToPhoton ← конвертация JVM строки → C++
BatchEvalPython [my_udf] ← Python UDF в JVM (НЕ Photon)
PhotonFilter ← нативная фильтрация (Photon C++)
PhotonScan ← нативное чтение Parquet (Photon C++)
Этот план имеет ДВА перехода: Photon→JVM→Photon
Каждый переход = конвертация строк = overhead
Анализ плана: ищем Photon-операторы¶
В Databricks UI физический план показывает префикс Photon у операторов, выполняемых нативно:
# Запрос для анализа
df = spark.read.format("delta").load("/mnt/lakehouse/orders/") \
.filter("amount > 1000") \
.join(
spark.read.format("delta").load("/mnt/lakehouse/customers/"),
on="customer_id",
how="left"
) \
.groupBy("region", "category") \
.agg(
F.sum("amount").alias("total_revenue"),
F.count("*").alias("orders"),
F.avg("amount").alias("avg_order"),
F.rank().over(
Window.partitionBy("region").orderBy(F.desc(F.sum("amount")))
).alias("revenue_rank") # Оконная функция
)
# Анализ плана
df.explain(extended=True)
# С Photon включённым, план выглядит так:
# == Physical Plan ==
# AdaptiveSparkPlan isFinalPlan=true
# +- == Current Plan ==
# PhotonResultStage
# +- PhotonHashAggregate(keys=[region, category],
# functions=[sum(amount), count(*), avg(amount)])
# +- PhotonColumnarExchange (hashpartitioning)
# +- PhotonHashAggregate(keys=[region, category], ...partial...)
# +- PhotonBroadcastHashJoin [customer_id]
# :- PhotonFilter (amount > 1000)
# : +- PhotonScan delta [orders]
# +- PhotonBroadcastExchange
# +- PhotonScan delta [customers]
#
# PhotonWindowAggregate добавляется для RANK() over window
#
# Все операторы с префиксом "Photon" - нативный C++.
# Если видите оператор без "Photon" - это JVM fallback.
Обратите внимание: PhotonHashAggregate появляется дважды - это двухфазная агрегация (partial + final), аналогично Spark SQL. Partial aggregate выполняется до shuffle (уменьшает объём данных), final aggregate - после shuffle.
Photon и Delta Lake: синергия форматов¶
Photon разрабатывался с учётом Delta Lake - open-source transactional storage format от Databricks. Эта синергия даёт дополнительные оптимизации, которые недоступны при работе с обычным Parquet.
Delta Lake Statistics: мгновенный Data Skipping¶
Delta Lake хранит statistics (min/max значения, null counts) для каждого Parquet-файла в метаданных транзакционного лога. При чтении Delta таблицы с фильтром WHERE date = '2024-01-15', PhotonScan использует эти статистики для пропуска файлов без чтения:
Delta Log statistics:
part-00000.parquet: date_min=2024-01-01, date_max=2024-01-10 ← SKIP
part-00001.parquet: date_min=2024-01-11, date_max=2024-01-20 ← READ
part-00002.parquet: date_min=2024-01-21, date_max=2024-01-31 ← SKIP
PhotonScan читает только 1 файл из 3 - без чтения самих данных из S3.
PhotonScan читает эти статистики из Delta Log напрямую в C++ коде - без накладных расходов JVM на парсинг JSON лога.
Deletion Vectors: SIMD-фильтрация удалённых строк¶
Delta Lake 2.0+ поддерживает Deletion Vectors - bitmap-файлы, которые помечают удалённые строки без физической перезаписи Parquet-файлов. PhotonScan применяет Deletion Vector как SIMD-bitmasking операцию при чтении данных:
Parquet файл: 1,000,000 строк
Deletion Vector (bitmap): [0, 1, 0, 0, 1, 1, 0, ...]
(1 = удалена, 0 = активна)
PhotonScan:
Читает Parquet columnar buffers (C++)
Применяет Deletion Vector как SIMD маску
Возвращает только активные строки без создания новых буферов
Альтернатива без Photon: JVM Spark читает Deletion Vector, применяет его строка за строкой через condition check - медленнее в 5–10 раз для таблиц с частыми обновлениями.
Liquid Clustering и Z-Order: Data Locality для SIMD¶
Delta Lake поддерживает Z-order clustering - физическую организацию данных таким образом, чтобы строки с похожими значениями ключевых колонок были рядом в одном Parquet файле. В сочетании с Photon это дает синергию:
- Z-order обеспечивает высокий data skipping (PhotonScan читает меньше файлов)
- Данные внутри файла тоже кластеризованы → лучше min/max statistics для row-group pruning
- Contiguous данные одного значения → лучше SIMD utilization в фильтрах
Практика: включение Photon и анализ производительности¶
Конфигурация Photon в Databricks¶
# ──────────────────────────────────────────────────────────
# Photon включён по умолчанию в Databricks Runtime 9.1+
# для кластеров с типом "Photon-enabled"
# ──────────────────────────────────────────────────────────
# Проверка статуса Photon:
photon_enabled = spark.conf.get(
"spark.databricks.photon.enabled", "false"
)
print(f"Photon enabled: {photon_enabled}")
# Включение/отключение Photon (для бенчмарка):
spark.conf.set("spark.databricks.photon.enabled", "true") # включить
spark.conf.set("spark.databricks.photon.enabled", "false") # отключить
# Дополнительные параметры Photon:
# Размер columnar батча (строк). Default: 4096.
# Большие батчи → лучше SIMD utilization, но больше память.
# Маленькие батчи → меньше overhead при fallback переходах.
spark.conf.set("spark.databricks.photon.batchSize", "4096")
# Включение Photon для конкретных операций:
# (все включены по умолчанию в Databricks Runtime)
spark.conf.set("spark.databricks.photon.hashJoin.enabled", "true")
spark.conf.set("spark.databricks.photon.sort.enabled", "true")
spark.conf.set("spark.databricks.photon.window.enabled", "true")
spark.conf.set("spark.databricks.photon.scan.enabled", "true")
# Парсинг Parquet через Photon (нативный C++ ридер):
spark.conf.set("spark.databricks.photon.parquet.enabled", "true")
Сравнительный бенчмарк: Photon ON vs OFF¶
# ──────────────────────────────────────────────────────────
# Бенчмарк: тяжёлый аналитический запрос
# ──────────────────────────────────────────────────────────
import time
from pyspark.sql import functions as F
from pyspark.sql.window import Window
# Тестовый запрос - типичная BI-аналитика:
# JOIN двух таблиц + тяжёлая агрегация + оконные функции
def heavy_analytics_query(spark):
orders = spark.read.format("delta").load("/mnt/lakehouse/orders/")
products = spark.read.format("delta").load("/mnt/lakehouse/products/")
customers = spark.read.format("delta").load("/mnt/lakehouse/customers/")
# JOIN трёх таблиц
enriched = orders \
.join(products, on="product_id", how="left") \
.join(customers.select("customer_id", "region", "segment"),
on="customer_id", how="left")
# Тяжёлая агрегация
aggregated = enriched \
.filter(F.col("order_date") >= "2023-01-01") \
.filter(F.col("amount") > 0) \
.groupBy("region", "segment", "category") \
.agg(
F.sum("amount").alias("total_revenue"),
F.count("*").alias("order_count"),
F.countDistinct("customer_id").alias("unique_customers"),
F.avg("amount").alias("avg_order_value"),
F.percentile_approx("amount", [0.25, 0.5, 0.75]).alias("quartiles"),
F.sum(F.col("amount") * F.col("quantity")).alias("gross_revenue")
)
# Оконные функции
w = Window.partitionBy("region").orderBy(F.desc("total_revenue"))
result = aggregated \
.withColumn("revenue_rank", F.rank().over(w)) \
.withColumn("revenue_pct",
F.col("total_revenue") / F.sum("total_revenue").over(
Window.partitionBy("region")
) * 100
)
return result
def run_benchmark(spark, label):
start = time.time()
result = heavy_analytics_query(spark)
count = result.count()
elapsed = time.time() - start
print(f"[{label}] {count} rows, {elapsed:.1f}s")
return elapsed
# Бенчмарк с Photon
spark.conf.set("spark.databricks.photon.enabled", "true")
t_photon = run_benchmark(spark, "PHOTON ON")
# Бенчмарк без Photon
spark.conf.set("spark.databricks.photon.enabled", "false")
t_tungsten = run_benchmark(spark, "PHOTON OFF (Tungsten)")
print(f"\nSpeedup: {t_tungsten / t_photon:.1f}×")
print(f"Photon: {t_photon:.1f}s")
print(f"Tungsten: {t_tungsten:.1f}s")
# Типичные результаты на больших данных (100GB+):
# PHOTON ON: 42.3s (3.8× speedup)
# PHOTON OFF: 161.5s
Обнаружение fallback в плане¶
# ──────────────────────────────────────────────────────────
# Диагностика: какие операторы выпали из Photon
# ──────────────────────────────────────────────────────────
from pyspark.sql import functions as F
from pyspark.sql.types import DoubleType
spark.conf.set("spark.databricks.photon.enabled", "true")
# Запрос С Python UDF - часть плана выйдет из Photon
@F.udf(returnType=DoubleType())
def custom_discount(amount, category):
"""Кастомные бизнес-правила скидки - Python UDF."""
if category == "Electronics" and amount > 5000:
return amount * 0.85
elif amount > 10000:
return amount * 0.90
return amount
df = spark.read.format("delta").load("/mnt/lakehouse/orders/")
query_with_udf = df \
.withColumn("discounted_amount", custom_discount("amount", "category")) \
.groupBy("region") \
.agg(F.sum("discounted_amount").alias("discounted_revenue"))
# Анализируем план - ищем Non-Photon операторы
plan = query_with_udf._jdf.queryExecution().executedPlan().toString()
photon_ops = [line for line in plan.split('\n') if 'Photon' in line]
jvm_ops = [line for line in plan.split('\n')
if 'Exec' in line and 'Photon' not in line and line.strip()]
print("Photon операторы:")
for op in photon_ops:
print(f" ✅ {op.strip()}")
print("\nJVM операторы (fallback):")
for op in jvm_ops:
print(f" ❌ {op.strip()}")
# Ожидаемый вывод:
# Photon операторы:
# ✅ PhotonResultStage
# ✅ PhotonHashAggregate (final)
# ✅ PhotonColumnarExchange
# ✅ PhotonHashAggregate (partial)
# ✅ PhotonFilter
# ✅ PhotonScan delta
#
# JVM операторы (fallback):
# ❌ RowToPhoton ← конвертация JVM → Photon (переход!)
# ❌ BatchEvalPython ← Python UDF в JVM
# ❌ RowToPhoton/PhotonToRow ← конвертации на границе
# Решение: заменить UDF на Spark SQL выражение
query_native = df \
.withColumn("discounted_amount",
F.when(
(F.col("category") == "Electronics") & (F.col("amount") > 5000),
F.col("amount") * 0.85
).when(
F.col("amount") > 10000,
F.col("amount") * 0.90
).otherwise(F.col("amount"))
) \
.groupBy("region") \
.agg(F.sum("discounted_amount").alias("discounted_revenue"))
# Теперь весь план - Photon. CASE WHEN поддерживается нативно.
query_native.explain()
# PhotonResultStage
# PhotonHashAggregate
# PhotonColumnarExchange
# PhotonHashAggregate (partial)
# PhotonProject [region, CASE WHEN...] ← нативный CASE WHEN
# PhotonScan delta
Мониторинг через Spark UI метрики¶
# ──────────────────────────────────────────────────────────
# Programmatic доступ к метрикам через REST API
# (работает в Databricks с токеном аутентификации)
# ──────────────────────────────────────────────────────────
import requests
import json
DATABRICKS_HOST = "https://adb-xxxxx.azuredatabricks.net"
DATABRICKS_TOKEN = "dapi..." # из Databricks personal access token
CLUSTER_ID = "0101-xxxxx-xxxxxx"
def get_photon_metrics(app_id):
"""Получает метрики Photon для текущего приложения."""
headers = {"Authorization": f"Bearer {DATABRICKS_TOKEN}"}
# Получаем список stages
url = f"{DATABRICKS_HOST}/api/2.0/spark/{CLUSTER_ID}/api/v1/applications/{app_id}/stages"
response = requests.get(url, headers=headers)
stages = response.json()
photon_time_ms = 0
jvm_time_ms = 0
for stage in stages:
stage_id = stage['stageId']
# Метрики отдельного stage
task_url = f"{DATABRICKS_HOST}/api/2.0/spark/{CLUSTER_ID}/api/v1/applications/{app_id}/stages/{stage_id}"
stage_detail = requests.get(task_url, headers=headers).json()
# В Databricks Photon-метрики видны в taskMetrics
for task_metric in stage_detail.get('taskMetrics', []):
# photonRowsFiltered, photonScanTime - специфичные Photon метрики
photon_scan = task_metric.get('photonScanTime', 0)
photon_time_ms += photon_scan
return {
"photon_scan_time_ms": photon_time_ms,
"photon_pct": photon_time_ms / (photon_time_ms + jvm_time_ms + 1) * 100
}
# В Databricks UI метрики Photon видны напрямую:
# Stage Details → Task Metrics → "Photon Time" vs "JVM Time"
# Job Graph → оранжевые узлы = Photon, синие = JVM fallback
Фактическая производительность: на каких запросах Photon даёт максимум¶
Не все запросы выигрывают от Photon одинаково. Понимание этого позволяет правильно распределять нагрузку.
Запросы с максимальным ускорением (3–10×)¶
Heavy aggregations с GROUP BY на больших данных. Миллионы уникальных ключей GROUP BY - hash table становится узким местом. Photon с cache-optimized C++ hash tables и SIMD-агрегацией выигрывает 4–8×.
-- Этот тип запроса ускоряется Photon в 4–8 раз:
SELECT
region,
DATE_TRUNC('month', order_date) AS month,
category,
subcategory,
SUM(amount) AS revenue,
COUNT(DISTINCT customer_id) AS unique_customers,
PERCENTILE_APPROX(amount, 0.95) AS p95_order_value
FROM orders
WHERE order_date >= '2023-01-01'
AND amount > 0
GROUP BY 1, 2, 3, 4
ORDER BY month, revenue DESC
JOIN-heavy запросы с несколькими таблицами. SortMergeJoin и BroadcastHashJoin - ключевые места, где vectorized hash probing с prefetching и SIMD comparison даёт кратное ускорение.
Scan с предикатами на больших таблицах. PhotonScan + Delta Lake statistics + SIMD filtering. Особенно эффективно при хорошей Z-order кластеризации.
Оконные функции. RANK(), ROW_NUMBER(), SUM OVER(PARTITION BY ...) - PhotonWindowAggregate ускоряет в 2–5×.
Запросы с умеренным ускорением (1.5–3×)¶
Средние таблицы (< 1 GB) - overhead на конвертацию виден. При маленьких данных Photon тратит относительно больше на инициализацию columnar батчей и конвертацию результатов.
Запросы с DISTINCT на высококардинальных колонках. COUNT(DISTINCT) с миллионами уникальных значений - hash table не помещается в L3 кэш, SIMD не может скрыть cache miss latency.
Запросы без ускорения (или замедление)¶
Python UDF в критическом пути. Любой Python UDF создаёт разрыв в Photon pipeline. Если UDF применяется к каждой строке и является критическим оператором - Photon overhead на конвертацию может быть больше выигрыша на других операторах.
# ❌ Антипаттерн: Python UDF разрывает Photon pipeline
@F.udf(returnType=DoubleType())
def bad_udf(amount):
return amount * 1.2 # Тривиальная операция - НИКОГДА не нужен UDF
# ✅ Правильно: Spark SQL expression - нативный Photon
df.withColumn("adjusted", F.col("amount") * 1.2)
Очень короткие запросы (< 1 секунды). Overhead на инициализацию columnar батчей, Photon runtime startup - может быть сравним с временем выполнения. Для OLTP-подобных запросов Photon не предназначен.
Structured Streaming (real-time). Photon оптимизирован для batch analytics. Micro-batch streaming работает с Photon, но выигрыш меньше - задержки на конвертацию между micro-batch'ами.
Photon vs открытые альтернативы¶
Photon - проприетарный движок. Понимание его места в экосистеме нативных engines важно для принятия архитектурных решений.
Сравнительная матрица нативных движков¶
| Критерий | Photon (Databricks) | Gluten + Velox | Apache Comet | Sail |
|---|---|---|---|---|
| Язык | C++ | C++ | Rust | Rust |
| Открытость | Closed-source | Open-source (Apache) | Open-source (Apache) | Open-source (Apache) |
| JVM в runtime | Driver + некоторые Executor | Да (плагин) | Да (плагин) | Нет |
| Совместимость | Только Databricks Runtime | Любой Spark кластер | Любой Spark кластер | Через Spark Connect |
| Delta Lake | Нативная интеграция | Ограниченная | Ограниченная | Нет |
| GPU | Нет | Нет | Нет | Нет |
| AQE | Полная поддержка | Поддерживается | Поддерживается | Через Catalyst |
| Python UDF | Fallback в JVM | Fallback в JVM | Fallback в JVM | Fallback в Python |
| Streaming | Micro-batch (частично) | Через Spark | Через Spark | Нет (beta) |
| Maturity | GA (Databricks) | Beta/GA (Meta backing) | Beta (Apache) | Alpha |
| Vendor lock-in | Только Databricks | Нет | Нет | Нет |
Почему Photon закрытый¶
Databricks сделал стратегическое решение оставить Photon проприетарным. Логика: Photon - это их главное конкурентное преимущество перед AWS EMR, Google Dataproc и Azure HDInsight, которые используют open-source Spark без Photon-оптимизаций. Открытие Photon лишило бы Databricks этого преимущества.
Velox (Meta) и DataFusion (Apache) открыты - компании-создатели не конкурируют на рынке Spark-as-a-service, им не нужно защищать execution engine как конкурентное преимущество.
Когда выбирать Photon¶
Выбирайте Photon если:
- Уже используете Databricks - Photon включён и настроен по умолчанию
- Workload: тяжёлая SQL-аналитика, BI, data warehouse
- Stack: Delta Lake + Spark SQL + Python (без тяжёлых Python UDF)
- Требуется enterprise SLA и поддержка
Рассмотрите альтернативы если:
- Работаете не в Databricks - Photon недоступен
- Vendor lock-in недопустим по политике компании
- Нужна full-stack open source
- GPU acceleration важнее CPU vectorization (RAPIDS cuDF)
Экономика облака: почему Photon экономит деньги¶
Photon имеет прямое финансовое значение, не только техническое. Рассмотрим конкретные механизмы экономии.
DBU (Databricks Units) и скорость выполнения¶
В Databricks стоимость вычислений измеряется в DBU (Databricks Units) - это единица, пропорциональная числу ядер × время работы. Кластер из 16 cores × 2 часа = 32 core-hours × цену DBU.
Если Photon ускоряет запрос в 3× - тот же объём данных обрабатывается за 1/3 времени. При фиксированном размере кластера:
Без Photon: 16 cores × 3 часа = 48 core-hours = $48 (при $1/core-hour)
С Photon: 16 cores × 1 час = 16 core-hours = $16
Экономия: $32 (67%)
Right-sizing: меньший кластер для той же задачи¶
Более мощный эффект: с Photon можно использовать меньший кластер для той же нагрузки:
Требование: обработать 5 TB данных за 2 часа (SLA)
Без Photon: нужно 32 cores (чтобы уложиться в SLA)
С Photon: хватает 12 cores (3× ускорение → 3× меньше cores)
Стоимость/день (8 часов работы):
Без Photon: 32 cores × 8h × $1 = $256/день = $7,680/месяц
С Photon: 12 cores × 8h × $1 = $96/день = $2,880/месяц
Экономия: $4,800/месяц (62%)
Для крупных корпораций, обрабатывающих петабайты данных ежедневно, это экономия миллионов долларов в год. Именно поэтому Databricks активно продвигает Photon как ключевую функцию и включает его по умолчанию.
Домашнее задание¶
Задание 1: Аудит физического плана (основное)¶
Вам предоставлен JSON-вывод df.explain(extended=True) для следующего запроса в Databricks:
SELECT
c.region,
p.category,
DATE_TRUNC('week', o.order_date) AS week,
SUM(o.amount) AS weekly_revenue,
COUNT(DISTINCT o.customer_id) AS unique_buyers,
SUM(o.amount) / LAG(SUM(o.amount), 1, NULL) OVER (
PARTITION BY c.region, p.category
ORDER BY DATE_TRUNC('week', o.order_date)
) - 1 AS wow_growth
FROM orders o
JOIN customers c ON o.customer_id = c.customer_id
JOIN products p ON o.product_id = p.product_id
WHERE o.order_date >= '2023-01-01'
AND my_custom_python_udf(o.amount, p.category) > 0.5
GROUP BY 1, 2, 3
ORDER BY week, weekly_revenue DESC
Выполните архитектурный аудит:
- Перечислите все операторы в плане с указанием:
Photon (нативный C++)илиJVM fallback - Найдите все переходы
PhotonToRow/RowToPhoton- сколько их, где они расположены? - Рассчитайте «коэффициент Photon-утилизации»: количество операторов в Photon / общее число операторов × 100%
- Предложите конкретное изменение SQL или Python-кода, которое максимизирует Photon-утилизацию. Объясните почему это изменение помогает.
Задание 2: Бенчмарк ключевых операторов (практическое)¶
Если у вас есть доступ к Databricks Trial или workspace:
- Создайте Delta таблицу с минимум 50 миллионами строк (можно использовать
spark.range()+withColumn()) - Напишите 3 запроса: чистый
GROUP BY,JOIN двух таблиц,оконные функции - Запустите каждый запрос с
photon.enabled=trueиphoton.enabled=false, измерьте время черезtime.time() - Для каждого запроса: какое ускорение дал Photon? Совпадает ли это с теоретическими ожиданиями?
- Запустите тот же
GROUP BYзапрос сPython UDFвместоCASE WHEN. Сравните с нативным вариантом. Насколько UDF «портит» Photon?
Задание 3: Архитектурный выбор (проектировочное)¶
Ваша команда оценивает переход с self-managed Apache Spark (на AWS EMR) на Databricks с Photon. Текущая инфраструктура:
- 10 daily batch jobs, каждая 2–4 часа, heavy SQL аналитика
- 3 streaming jobs (Spark Structured Streaming, micro-batch 1 минута)
- 2 ML jobs (scikit-learn inference через Python UDF)
- Итого: ~$15,000/месяц на AWS EMR
Напишите архитектурный анализ (1–2 страницы):
- Для каких jobs Photon даст максимальный эффект и почему?
- Для каких jobs Photon не поможет или поможет мало?
- Какие изменения в коде нужны для максимизации Photon-утилизации (особенно для ML jobs)?
- Оцените потенциальную экономию с Photon (в процентах и абсолютных числах), используя метрики из урока
Резюме¶
Photon - это конкретный ответ Databricks на вопрос «что будет, если написать Spark execution engine с нуля, используя всё, что мы знаем о современном аппаратном обеспечении».
Ключевые архитектурные решения Photon и их обоснование: C++ вместо JVM даёт детерминированную SIMD-генерацию через #pragma omp simd и явные intrinsics без зависимости от JIT; columnar-native формат вместо UnsafeRow даёт cache-friendly последовательный доступ к памяти; batch-at-a-time вместо row-at-a-time даёт 1024–4096 строк за вызов вместо одной; vectorized hash tables с prefetching устраняют cache miss stall в JOIN; bitmasking для NULL делает обработку NULL branchless.
Практическое следствие: на heavy SQL workloads (JOIN, GROUP BY, оконные функции) Photon ускоряет выполнение в 3–10 раз. Это прямая экономия облачных расходов через либо меньшее время выполнения (меньше DBU), либо меньший кластер (меньше cores).
Главное ограничение: Python UDF разрывает native path. Каждый Python UDF создаёт PhotonToRow → JVM → RowToPhoton переход, уничтожая часть выигрыша. Иерархия выбора неизменна: Native Spark SQL → Spark SQL expressions → Pandas UDF → Python UDF - и с Photon эта иерархия становится ещё более критичной.
Photon остаётся закрытым продуктом Databricks. Для тех, кто работает вне Databricks, эквивалентные подходы - Apache Comet (DataFusion/Rust), Gluten+Velox (C++) или Sail (полная замена JVM) - реализуют те же архитектурные принципы в открытом виде.