Photon (Databricks): SIMD-движок на C++, что заменяет в Tungsten, 3–10× на Join и агрегациях

Архитектура Photon: почему Databricks переписали Spark на C++, как работают векторизованные операторы с AVX-512, fallback в JVM, и когда ускорение достигает 10×.

optimization

Рождение 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 ЗАМЕНЯЕТ

  • Физические операторы: HashAggregateExecPhotonHashAggregateExec, SortMergeJoinExecPhotonSortMergeJoinExec и т.д.
  • 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, Delta
  • PhotonFilter - фильтрация с предикатами (column predicates, BETWEEN, IN)
  • PhotonProject - вычисление выражений (арифметика, string functions, CASE WHEN)
  • PhotonHashAggregate - GROUP BY с SUM, COUNT, AVG, MIN, MAX
  • PhotonBroadcastHashJoin - JOIN с broadcast side
  • PhotonSortMergeJoin - 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

Выполните архитектурный аудит:

  1. Перечислите все операторы в плане с указанием: Photon (нативный C++) или JVM fallback
  2. Найдите все переходы PhotonToRow / RowToPhoton - сколько их, где они расположены?
  3. Рассчитайте «коэффициент Photon-утилизации»: количество операторов в Photon / общее число операторов × 100%
  4. Предложите конкретное изменение SQL или Python-кода, которое максимизирует Photon-утилизацию. Объясните почему это изменение помогает.

Задание 2: Бенчмарк ключевых операторов (практическое)

Если у вас есть доступ к Databricks Trial или workspace:

  1. Создайте Delta таблицу с минимум 50 миллионами строк (можно использовать spark.range() + withColumn())
  2. Напишите 3 запроса: чистый GROUP BY, JOIN двух таблиц, оконные функции
  3. Запустите каждый запрос с photon.enabled=true и photon.enabled=false, измерьте время через time.time()
  4. Для каждого запроса: какое ускорение дал Photon? Совпадает ли это с теоретическими ожиданиями?
  5. Запустите тот же 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 страницы):

  1. Для каких jobs Photon даст максимальный эффект и почему?
  2. Для каких jobs Photon не поможет или поможет мало?
  3. Какие изменения в коде нужны для максимизации Photon-утилизации (особенно для ML jobs)?
  4. Оцените потенциальную экономию с 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) - реализуют те же архитектурные принципы в открытом виде.