Lambda-архитектура в Lakehouse: интеграция Batch и Real-time потоков на Spark

Lambda-архитектура в Lakehouse: интеграция Batch и Real-time потоков на Spark

datamodeling

Современная аналитическая платформа должна одновременно отвечать на два противоречивых запроса бизнеса: «покажи мне точные данные за прошлый год» и «покажи мне, что происходит прямо сейчас». Исторически эти задачи решались разными технологическими стеками, разными командами и разными языками программирования. Lambda Architecture была первым серьёзным ответом на вопрос: «как объединить точность исторической обработки с низкой задержкой потоков?»

Сегодня Spark и форматы Lakehouse (Apache Iceberg, Delta Lake) позволяют реализовать Lambda-подход на единой платформе, избавившись от большинства её классических болей. Но чтобы понять, что именно современный стек решает, нужно сначала разобраться, почему эта архитектура появилась и какие фундаментальные компромиссы она устанавливает.


1. Почему Batch-only архитектуры перестали справляться

1.1 Ночной ETL как единственный источник аналитики

До середины 2010-х годов стандартная аналитическая платформа работала по простой схеме: каждую ночь Spark-джоб или Hive-скрипт читал данные из OLTP-базы, трансформировал их и обновлял витрины данных в DWH. Утром аналитики приходили на работу и видели данные за вчера.

Для многих задач это было нормально: финансовые отчёты за прошлый месяц, недельные KPI, квартальные дашборды. Бизнес не требовал большего. Но несколько категорий задач принципиально изменили требования:

Fraud Detection (Обнаружение мошенничества). Транзакцию нужно заблокировать в течение сотен миллисекунд - задолго до того, как деньги покинут счёт. Ночной ETL здесь бессмысленен: к утру транзакция уже проведена, деньги выведены, мошенник недосягаем. Системы защиты от фрода требуют анализа текущего паттерна поведения пользователя в реальном времени.

Recommendation Systems (Рекомендательные системы). Пользователь смотрит товар прямо сейчас. Рекомендации, основанные на его сессии за последние 10 минут, конвертируют в 3–5 раз лучше, чем рекомендации на основе истории за прошлый месяц. Netflix, Amazon и YouTube давно это поняли и перестроили свои стеки под real-time scoring.

Operational Dashboards (Операционные дашборды). Менеджер склада хочет видеть не «сколько заказов было вчера», а «сколько заказов ожидает обработки прямо сейчас». Call-центр мониторит загруженность операторов в реальном времени. Эти задачи требуют задержки в секунды, а не часы.

Event-Driven Systems. Микросервисная архитектура генерирует тысячи событий в секунду: order_placed, payment_initiated, item_reserved, shipment_dispatched. Аналитика этих событий нужна немедленно - для мониторинга аномалий, расчёта SLA и реагирования на сбои.

1.2 Фундаментальный конфликт: точность против скорости

Когда инженеры попытались просто «ускорить» batch ETL, они столкнулись с фундаментальным ограничением: Hadoop/Hive не был создан для низкой задержки. Минимальное время выполнения MapReduce-джоба - несколько минут, даже для маленьких датасетов. Причина - overhead на инициализацию YARN-контейнеров, планирование задач и запись промежуточных результатов в HDFS.

Параллельно развивались потоковые движки - Apache Storm, затем Apache Flink. Они обеспечивали задержки в миллисекунды, но приходилось писать совершенно другой код: event-at-a-time обработка вместо SQL, stateful операторы вместо GROUP BY, сложное управление состоянием вместо простых агрегаций.

Пропасть между мирами была реальной инженерной проблемой. Batch-пайплайн на Hive считал «выручку за неделю» правильно, но вчера. Streaming-пайплайн на Storm давал приблизительный ответ сейчас, но результаты расходились с batch-расчётом из-за разной логики обработки. Аналитики и бизнес не могли доверять ни одному из источников.

Именно в этот момент Натан Марц (Nathan Marz) сформулировал Lambda Architecture как ответ на вопрос: «можно ли иметь и точность, и скорость одновременно?»


2. Lambda Architecture: архитектура двух слоёв

2.1 Исторический контекст и философия

Lambda Architecture была описана Натаном Марцем в 2011–2012 годах на основе опыта работы с Twitter и BackType. Книга «Big Data: Principles and best practices of scalable realtime data systems» (2015) стала каноническим текстом эпохи.

Ключевая идея: не пытайся получить точность и скорость в одной системе. Вместо этого - раздели задачи. Позволь одному слою отвечать за историческую точность, другому - за мгновенный отклик. Покупатель BI-дашборда видит результат, объединяющий оба слоя.

Философски Lambda опирается на принцип иммутабельности: исходные данные никогда не изменяются. Любая ошибка в логике обработки исправляется пересчётом от сырых данных, а не патчингом значений в таблицах. Это гарантирует воспроизводимость - hallmark зрелой data-платформы.

2.2 Три слоя классической Lambda Architecture

Data Sources - всё, что генерирует данные: транзакционные СУБД, Kafka-топики, REST API, IoT-устройства, S3-логи. Все события записываются как иммутабельные факты в Batch Layer.

Batch Layer состоит из двух компонентов. Master Dataset - неизменяемое хранилище всех сырых данных с самого начала. Это append-only файлы в HDFS или S3, партиционированные по дате. Никакого UPDATE, никакого DELETE - только добавление. Batch Views - предварительно вычисленные агрегаты: выручка по дням, количество активных пользователей по неделям. Они пересчитываются полностью каждые несколько часов и обеспечивают абсолютную точность, потому что работают со всем историческим датасетом.

Speed Layer компенсирует задержку Batch Layer. Пока Batch-джоб пересчитывает данные за последние 6 часов, Speed Layer уже обрабатывает события, поступившие 30 секунд назад. Его задача - только «горячий хвост»: события, которые ещё не вошли в Batch Views. Точность немного хуже (нет полной истории для контекста), но задержка - секунды.

Serving Layer получает запрос от аналитика и объединяет ответы двух источников: исторически точные Batch Views + свежие Realtime Views. Для запроса «выручка сегодня» Serving Layer возвращает сумму значений из обоих слоёв.

2.3 Иммутабельная модель данных и философия перепроведения

Центральное свойство Lambda Architecture - иммутабельность исходных данных. Master Dataset - это не таблица с текущим состоянием, а append-only журнал всего, что когда-либо произошло. Эта идея прямо заимствована из Event Sourcing.

Почему это важно? Представьте, что в Batch View обнаружена ошибка: бизнес-логика неверно суммировала заказы с промокодами. В мутабельной системе данные в витринах уже испорчены, исходные значения перезаписаны UPDATE-ами, восстановить правду невозможно. В Lambda Architecture:

  1. Исправляем код трансформации
  2. Удаляем неверные Batch Views
  3. Запускаем пересчёт с нуля от Master Dataset
  4. Через несколько часов имеем правильные данные за всё историческое время

Master Dataset при этом нетронут - он хранит каждое событие, которое было получено от источников. Это свойство называется replayability (воспроизводимость) и является фундаментом надёжной data-платформы.


3. Batch Layer: Source of Truth и историческая обработка

3.1 Структура Master Dataset

Master Dataset в классической Lambda Architecture - это партиционированные файлы в HDFS. В современной Lakehouse-реализации - это Bronze-слой в Iceberg или Delta Lake. Принцип один: никогда не изменяй уже записанные данные.

Партиционирование по event_date - это не просто оптимизация, это архитектурное требование. Batch-пересчёт работает по конкретным временным диапазонам. Без партиционирования каждый запрос «данные за эту неделю» читает весь архив. С партиционированием Spark применяет Static Partition Pruning и читает ровно 7 партиций.

Только additive schema evolution: новые колонки добавлять можно, существующие удалять нельзя. Это гарантирует, что исторические данные читаются корректно даже после изменения схемы.

3.2 Детерминизм трансформаций: почему это критично

Batch Views должны быть детерминированы: для одних и тех же входных данных они всегда дают одинаковый результат. Это кажется очевидным, но на практике часто нарушается.

Примеры недетерминированных трансформаций, которые ломают Batch Layer:

# ПЛОХО: результат зависит от времени запуска
df.withColumn("is_recent", F.col("event_ts") > F.current_timestamp() - F.expr("INTERVAL 30 DAYS"))

# ХОРОШО: детерминированно, зависит только от данных
df.withColumn("is_recent", F.col("event_ts") > F.lit("2026-04-23").cast("timestamp"))

# ПЛОХО: случайная выборка меняется при каждом запуске
df.sample(fraction=0.1, seed=None)

# ХОРОШО: фиксированный seed гарантирует воспроизводимость
df.sample(fraction=0.1, seed=42)

# ПЛОХО: зависимость от внешнего состояния (текущего курса валют)
df.withColumn("usd_amount", F.col("rub_amount") / get_current_exchange_rate())

# ХОРОШО: исторический курс валют читается из справочника
exchange_rates = spark.table("lakehouse.reference.exchange_rates")
df.join(exchange_rates, on=["date", "currency"], how="left")

Нарушение детерминизма означает, что два запуска одного пайплайна на одних и тех же данных дадут разные результаты. Это делает debugging невозможным: воспроизвести инцидент не получится, потому что повторный запуск уже даёт другой ответ.

3.3 Batch пересчёт: полный vs инкрементальный

В классической Lambda Architecture Batch Layer выполняет полный пересчёт (full recomputation): каждые N часов джоб читает весь Master Dataset от начала времён и пересчитывает все Batch Views заново. Это максимально надёжно, но дорого при больших объёмах данных.

В современных Lakehouse-реализациях применяется инкрементальный пересчёт: читаются только новые партиции (через Iceberg Incremental Read или Delta Lake DESCRIBE HISTORY), и результат мержится с предыдущими Batch Views через MERGE INTO. При этом принцип иммутабельности Master Dataset сохраняется - новые данные всегда APPEND-only, обновляются только выходные витрины.

Инкрементальный пересчёт - это не отступление от философии Lambda. Это применение той же идеи (точность от полной истории) с практической оптимизацией: вместо полного пересчёта используем Iceberg Snapshots для чтения только новых данных. Результат математически эквивалентен.


4. Speed Layer: Real-time Ingestion и Structured Streaming

4.1 Роль Speed Layer в Lambda Architecture

Speed Layer существует для одной цели: покрыть временной лаг Batch Layer. Если Batch-пайплайн пересчитывает данные раз в 6 часов, то последние 6 часов данных отсутствуют в Batch Views. Speed Layer заполняет этот пробел, обрабатывая события по мере их поступления.

Принципиальное отличие от Batch Layer: Speed Layer не хранит полную историю. Он работает только с «горячим хвостом» - событиями, которые ещё не вошли в Batch Views. Как только Batch-джоб завершается и включает в себя очередной временной диапазон, соответствующие данные из Speed Layer становятся ненужными и могут быть удалены.

Это разделение ответственности означает, что Speed Layer может пожертвовать точностью ради скорости. Например, он может не учитывать некоторые типы late-arriving событий или делать приближённые агрегации - потому что в конечном счёте Batch Layer пересчитает всё точно.

4.2 Spark Structured Streaming: единый API для Batch и Stream

Революция Spark Structured Streaming (появился в Spark 2.0) заключается в том, что для streaming-обработки используется тот же DataFrame API, что и для batch. Это прямое воплощение идеи Lambda Architecture на уровне кода: одна и та же трансформация может работать как на историческом датасете (Batch Layer), так и на потоке событий (Speed Layer).

Ключевая абстракция - Unbounded Table (бесконечная таблица). Spark представляет поток данных как таблицу, в которую постоянно добавляются новые строки. Каждый микро-батч читает очередную «порцию» из этой таблицы и обрабатывает её как обычный DataFrame.

Unbounded Table растёт бесконечно слева направо. Каждый микро-батч читает только «новую» часть таблицы - ту, что появилась после предыдущего триггера. Spark хранит в checkpoint, какие смещения Kafka уже обработаны, чтобы следующий микро-батч начинал ровно с того места, где остановился предыдущий.

4.3 Настройка Kafka-источника: управление смещениями

from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.types import StructType, StructField, StringType, LongType, TimestampType

spark = SparkSession.builder \
    .appName("Lambda-SpeedLayer") \
    .config("spark.sql.extensions",
            "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
    .config("spark.sql.catalog.lakehouse", "org.apache.iceberg.spark.SparkCatalog") \
    .config("spark.sql.catalog.lakehouse.type", "hadoop") \
    .config("spark.sql.catalog.lakehouse.warehouse", "s3://my-bucket/warehouse") \
    .getOrCreate()

# Схема события из Kafka
# Kafka хранит сообщения как байтовые строки: value - JSON с бизнес-данными
event_schema = StructType([
    StructField("event_id",    StringType(),    nullable=False),
    StructField("user_id",     StringType(),    nullable=False),
    StructField("event_type",  StringType(),    nullable=False),
    StructField("amount",      LongType(),      nullable=True),
    StructField("event_ts",    TimestampType(), nullable=False),
])

# Читаем из Kafka
# startingOffsets="latest": при первом запуске начинаем с конца топика.
# Это правильно для Speed Layer: нам нужны только свежие события,
# история уже обработана Batch Layer.
# При перезапуске checkpoint автоматически восстановит позицию.
raw_kafka_stream = (
    spark.readStream
         .format("kafka")
         .option("kafka.bootstrap.servers", "kafka:9092")
         .option("subscribe",               "ecommerce.events")
         .option("startingOffsets",         "latest")
         # Максимальное количество сообщений за один микро-батч.
         # Защищает от ситуации, когда consumer lag вырос и один микро-батч
         # пытается обработать миллионы сообщений, вызывая OOM.
         .option("maxOffsetsPerTrigger",    "50000")
         # Если топик не существует - упасть, а не ждать молча
         .option("failOnDataLoss",          "true")
         .load()
)

# Десериализуем JSON из поля value
# Kafka хранит value как BINARY, поэтому CAST к STRING обязателен
parsed_stream = (
    raw_kafka_stream
    .select(
        # Метаданные Kafka: всегда полезны для debugging
        F.col("topic"),
        F.col("partition"),
        F.col("offset"),
        F.col("timestamp").alias("kafka_ts"),      # Время записи в Kafka
        # Десериализуем бизнес-данные
        F.from_json(
            F.col("value").cast("string"),
            event_schema
        ).alias("data")
    )
    .select(
        "topic", "partition", "offset", "kafka_ts",
        "data.*"   # Разворачиваем структуру: event_id, user_id, event_type, amount, event_ts
    )
)

Обратите внимание на maxOffsetsPerTrigger. Без этого ограничения при перезапуске после длительного простоя один микро-батч попытается обработать весь накопившийся consumer lag - миллионы сообщений. Это приведёт к OOM или к очень долгому первому микро-батчу, который заблокирует всё остальное. Ограничение позволяет системе постепенно «догонять» отставание.

4.4 Выбор триггера: компромисс между задержкой и стоимостью

Trigger - это стратегия запуска микро-батчей. Выбор триггера - это баланс между задержкой (latency) и стоимостью (cost).

# Вариант 1: Fixed Interval - триггер каждые N секунд
# Используется: когда нужна предсказуемая задержка (например, обновление дашборда раз в 30 сек)
# Поведение: если обработка предыдущего батча заняла больше N секунд,
#            следующий запускается сразу после завершения, без ожидания
query_fixed = (
    parsed_stream.writeStream
                 .trigger(processingTime="30 seconds")
                 .format("iceberg")
                 .outputMode("append")
                 .toTable("lakehouse.silver.speed_layer_events")
)

# Вариант 2: Once - запустить один микро-батч и остановиться
# Используется: в Lambda Architecture как "pseudo-streaming" —
# запускается по cron каждые 5 минут, обрабатывает накопившееся, завершается.
# Преимущество: не нужен постоянно работающий Spark-процесс.
# Недостаток: задержка = интервал cron (5-15 минут).
query_once = (
    parsed_stream.writeStream
                 .trigger(once=True)
                 .format("iceberg")
                 .outputMode("append")
                 .toTable("lakehouse.silver.speed_layer_events")
)

# Вариант 3: AvailableNow - обработать всё доступное сейчас, затем остановиться
# Spark 3.3+: аналог Once, но разбивает работу на несколько микро-батчей,
# что даёт более эффективное использование ресурсов при большом lag.
query_available_now = (
    parsed_stream.writeStream
                 .trigger(availableNow=True)
                 .format("iceberg")
                 .outputMode("append")
                 .toTable("lakehouse.silver.speed_layer_events")
)

Рекомендация для Lambda Speed Layer: используйте processingTime="30 seconds" или processingTime="1 minute" для задержки в секунды, или availableNow=True по расписанию для задержки в минуты. Trigger continuous (экспериментальный, истинный streaming без микро-батчей) в production практически не используется из-за ограниченной поддержки операций.

4.5 Watermark: управление опоздавшими событиями

В реальных системах события приходят с опозданием. Мобильное приложение работало офлайн и отправило пачку событий через 30 минут. Сетевой сбой задержал сообщение. Время на стороне клиента было неправильно настроено.

Watermark - это механизм, который говорит Spark: «события старше X минут от последнего обработанного события мы считаем опоздавшими и не ждём их». Без watermark Spark вынужден хранить весь state бесконечно, ожидая теоретически возможных опоздавших событий - это приводит к OOM.

# Watermark: ждём опоздавшие события до 10 минут
# Это означает: если последнее событие имело event_ts = 10:00,
# мы ещё принимаем события с event_ts >= 09:50.
# Событие с event_ts = 09:40 считается опоздавшим и игнорируется.
windowed_stream = (
    parsed_stream
    .withWatermark("event_ts", "10 minutes")   # Ждём опоздавших до 10 минут
    .groupBy(
        F.window("event_ts", "1 minute"),       # Группируем по 1-минутным окнам
        "event_type"
    )
    .agg(
        F.count("*").alias("event_count"),
        F.sum("amount").alias("total_amount"),
        F.approx_count_distinct("user_id").alias("unique_users")
    )
)

# Записываем агрегированные Realtime Views
query = (
    windowed_stream
    .writeStream
    .trigger(processingTime="30 seconds")
    .outputMode("update")            # update: обновляем только изменившиеся строки окон
    .format("iceberg")
    .option("checkpointLocation", "s3://my-bucket/checkpoints/speed-layer-agg")
    .toTable("lakehouse.gold.realtime_views")
)

query.awaitTermination()

outputMode("update") означает: записываем только строки, которые изменились в этом микро-батче. Это оптимально для агрегированных Realtime Views, где большинство временных окон уже закрыты и не меняются.


5. Serving Layer: объединение Batch и Realtime Views

5.1 Главная сложность Serving Layer

Serving Layer - самая инженерно сложная часть Lambda Architecture. Аналитический запрос должен одновременно видеть:

  • Исторически точные данные из Batch Views (за все прошедшие периоды)
  • Свежие данные из Realtime Views (за текущий «горячий хвост»)
  • Без дублирования: одно и то же событие не должно суммироваться дважды

Сложность возникает из-за временного перекрытия: в момент, когда Batch Layer завершает очередной пересчёт и «захватывает» последние 6 часов данных, Speed Layer ещё хранит Realtime Views за эти же 6 часов. Если просто сложить оба результата - получим двойной счёт.

Временное перекрытие - системная проблема Serving Layer. В момент, когда новый Batch View только что записан, он перекрывается с Realtime View за тот же период. Serving Layer должен знать: «после 18:00 используй Batch для периода 12:00-18:00, и только Speed Layer - для 18:00-сейчас».

5.2 Паттерн «Query-time Union»: объединение в момент запроса

В Lakehouse-реализации Serving Layer реализуется через SQL-представления (Views), которые объединяют Batch и Speed результаты:

-- Представление Serving Layer, объединяющее Batch и Speed данные
-- Ключевой принцип: Speed Layer покрывает только периоды ПОЗЖЕ последнего Batch run
CREATE OR REPLACE VIEW lakehouse.gold.serving_revenue_by_hour AS

WITH batch_coverage AS (
    -- Определяем, до какого момента покрыт Batch Layer
    -- Читаем метаданные последнего Batch View
    SELECT MAX(processed_up_to_ts) as batch_max_ts
    FROM lakehouse.gold.batch_view_metadata
),

batch_data AS (
    -- Данные из Batch Layer: полностью точные, за все завершённые периоды
    SELECT
        hour_ts,
        event_type,
        total_amount,
        event_count,
        'batch' as source_layer
    FROM lakehouse.gold.batch_revenue_by_hour
    WHERE hour_ts < (SELECT batch_max_ts FROM batch_coverage)
),

speed_data AS (
    -- Данные из Speed Layer: только за период ПОСЛЕ последнего Batch run
    -- Это предотвращает двойной счёт
    SELECT
        window.start as hour_ts,
        event_type,
        total_amount,
        event_count,
        'speed' as source_layer
    FROM lakehouse.gold.realtime_views
    WHERE window.start >= (SELECT batch_max_ts FROM batch_coverage)
)

-- UNION ALL без дублирования: batch и speed покрывают непересекающиеся временные периоды
SELECT * FROM batch_data
UNION ALL
SELECT * FROM speed_data
ORDER BY hour_ts DESC;

Этот запрос атомарно определяет, где заканчивается Batch и начинается Speed, и возвращает объединённые данные без двойного счёта. При каждом вызове batch_max_ts пересчитывается, поэтому представление автоматически адаптируется к прогрессу Batch Layer.

5.3 Практическая реализация: PySpark Serving Query

def serve_unified_query(metric: str, start_date: str, end_date: str) -> None:
    """
    Пример Serving Layer: возвращает метрику за период,
    объединяя Batch Views и Realtime Views без дублирования.

    metric: имя метрики (например, 'revenue', 'orders_count')
    start_date: начало периода, формат 'YYYY-MM-DD'
    end_date: конец периода, формат 'YYYY-MM-DD'
    """
    # 1. Определяем временную границу Batch Layer
    # Читаем метаданные, которые Batch-джоб обновляет при каждом запуске
    batch_metadata = spark.table("lakehouse.gold.batch_view_metadata")
    batch_max_ts = batch_metadata.select(F.max("processed_up_to_ts")).collect()[0][0]

    print(f"Batch Layer покрывает данные до: {batch_max_ts}")
    print(f"Speed Layer покрывает данные с: {batch_max_ts} до сейчас")

    # 2. Читаем Batch Views: точные исторические данные
    batch_views = (
        spark.table("lakehouse.gold.batch_revenue_by_hour")
        .filter(F.col("hour_ts").between(start_date, end_date))
        # Берём ТОЛЬКО данные до границы Batch Layer
        .filter(F.col("hour_ts") < F.lit(batch_max_ts))
        .select(
            "hour_ts", "event_type",
            F.col("total_amount"),
            F.col("event_count"),
            F.lit("batch").alias("source_layer")
        )
    )

    # 3. Читаем Realtime Views: свежие данные
    realtime_views = (
        spark.table("lakehouse.gold.realtime_views")
        .filter(F.col("window.start").between(start_date, end_date))
        # Берём ТОЛЬКО данные ПОСЛЕ границы Batch Layer - избегаем дублирования
        .filter(F.col("window.start") >= F.lit(batch_max_ts))
        .select(
            F.col("window.start").alias("hour_ts"),
            "event_type",
            F.col("total_amount"),
            F.col("event_count"),
            F.lit("speed").alias("source_layer")
        )
    )

    # 4. Объединяем два источника
    # UNION ALL корректен: временные диапазоны не пересекаются
    unified = batch_views.unionByName(realtime_views, allowMissingColumns=True)

    # 5. Финальная агрегация по дням для дашборда
    result = (
        unified
        .groupBy(F.to_date("hour_ts").alias("date"), "event_type")
        .agg(
            F.sum("total_amount").alias("total_amount"),
            F.sum("event_count").alias("event_count"),
            F.countDistinct("source_layer").alias("layers_used")  # 1=только batch, 2=batch+speed
        )
        .orderBy("date", "event_type")
    )

    print("Результат unified query (Batch + Speed):")
    result.show(truncate=False)

6. Structured Streaming и Lakehouse: реализация Lambda без «пропасти»

6.1 Как Spark преодолел «пропасть» между Batch и Stream

Классическая проблема Lambda Architecture - дублирование бизнес-логики. Нужно написать одинаковую трансформацию дважды: на Hive/SQL для Batch Layer и на Storm/Flink для Speed Layer. Любое изменение бизнес-правила требует синхронного обновления обоих пайплайнов. Рассинхронизация неизбежна.

Spark Structured Streaming решает это радикально: один и тот же код DataFrame работает в обоих режимах.

def compute_revenue_metrics(df):
    """
    Эта функция работает ОДИНАКОВО для:
    1. Batch DataFrame (исторические данные)
    2. Streaming DataFrame (поток событий в реальном времени)
    Бизнес-логика написана один раз.
    """
    return (
        df
        # Фильтрация: только успешные транзакции с положительной суммой
        .filter(F.col("event_type") == "purchase")
        .filter(F.col("amount") > 0)
        # Обогащение: категоризация суммы
        .withColumn("amount_tier",
            F.when(F.col("amount") < 1000, "micro")
             .when(F.col("amount") < 10000, "mid")
             .otherwise("large")
        )
        # Агрегация по временному окну
        .groupBy(
            F.window("event_ts", "1 hour").alias("hour_window"),
            "amount_tier"
        )
        .agg(
            F.sum("amount").alias("revenue"),
            F.count("*").alias("transactions"),
            F.approx_count_distinct("user_id", rsd=0.01).alias("unique_buyers")
        )
    )

# --- BATCH режим (Batch Layer) ---
historical_df = (
    spark.table("lakehouse.bronze.events")
         .filter(F.col("event_date").between("2026-01-01", "2026-05-22"))
)
# Используем ту же функцию
batch_result = compute_revenue_metrics(historical_df)
batch_result.write.format("iceberg").mode("overwrite").saveAsTable("lakehouse.gold.batch_revenue")

# --- STREAMING режим (Speed Layer) ---
# Тот же вызов функции, но df - это streaming DataFrame
streaming_df = (
    spark.readStream
         .format("kafka")
         .option("kafka.bootstrap.servers", "kafka:9092")
         .option("subscribe", "ecommerce.events")
         .option("startingOffsets", "latest")
         .load()
         .selectExpr("CAST(value AS STRING) as json_payload")
         .select(F.from_json("json_payload", event_schema).alias("data"))
         .select("data.*")
)

# Используем ту же функцию - она работает и для streaming!
streaming_result = compute_revenue_metrics(streaming_df.withWatermark("event_ts", "10 minutes"))

streaming_query = (
    streaming_result
    .writeStream
    .trigger(processingTime="30 seconds")
    .outputMode("update")
    .format("iceberg")
    .option("checkpointLocation", "s3://my-bucket/checkpoints/speed-revenue")
    .toTable("lakehouse.gold.realtime_revenue")
)

Единая функция compute_revenue_metrics - это квинтэссенция того, что Spark принёс в мир Lambda. Больше нет двух разных кодовых баз. Изменение фильтра или логики агрегации применяется одновременно к обоим слоям.

6.2 Lakehouse как фундамент Lambda

Apache Iceberg и Delta Lake принесли в Lakehouse то, чего не хватало классической Lambda: ACID-транзакции для конкурентной записи. Это решает ключевую проблему: что происходит, когда Batch-джоб и Streaming-джоб одновременно пишут в одну таблицу?

Snapshot Isolation в Iceberg работает так: каждый джоб видит консистентный снапшот данных на момент начала своей транзакции. Streaming-джоб делает коммиты каждые 30 секунд - создаёт новые снапшоты. Batch-джоб читает все данные в рамках одного снапшота (того, что был в начале), не видит промежуточных коммитов Streaming и в конце атомарно создаёт свой снапшот. Конфликта нет, данные консистентны.

Без ACID-транзакций (как в классическом Hadoop) параллельная запись могла бы привести к «половинным коммитам»: аналитик видит данные, которые Batch-джоб записал наполовину, пока Streaming продолжает добавлять свои строки.


7. Паттерны записи на Silver: foreachBatch и MERGE

7.1 Проблема мелких файлов в Speed Layer

Каждый микро-батч Structured Streaming пишет в Iceberg/Delta новые файлы. При триггере processingTime="30 seconds" за сутки создаётся 86400 / 30 = 2880 микро-батчей. Если в каждом - хотя бы 1 файл на партицию, за сутки накапливается тысячи мелких файлов. Запрос аналитика, читающий все партиции, откроет тысячи файлов вместо десятков - это overhead на metadata операции и на JVM GC.

Компакция - это процесс слияния мелких файлов в крупные. В Iceberg это операция rewrite_data_files, в Delta Lake - OPTIMIZE. Ключевое свойство обоих: компакция работает конкурентно со стримингом - Streaming пишет новые файлы, пока компакция перепаковывает старые, без взаимной блокировки.

7.2 foreachBatch: мощный паттерн произвольной логики записи

foreachBatch позволяет выполнять произвольный Python-код для каждого микро-батча. Вместо стандартного format("iceberg").toTable(...) вы получаете DataFrame микро-батча и можете делать с ним всё что угодно: несколько записей, MERGE INTO, обновление нескольких таблиц.

def process_microbatch(batch_df, batch_id: int) -> None:
    """
    foreachBatch функция: вызывается для каждого микро-батча.
    batch_df: обычный (не streaming) DataFrame с данными текущего микро-батча
    batch_id: уникальный ID батча, монотонно возрастающий

    Идемпотентность: функция должна давать одинаковый результат при повторном вызове
    с теми же данными. Spark может вызвать её повторно при сбое.
    Идемпотентность гарантируется через дедупликацию по batch_id и event_id.
    """
    if batch_df.isEmpty():
        print(f"Батч {batch_id}: пустой, пропускаем")
        return

    # 1. Дедупликация: защита от повторного вызова (at-least-once semantics foreachBatch)
    deduped_df = batch_df.dropDuplicates(["event_id"])

    # 2. Кешируем DataFrame: он будет использован несколько раз
    # Без cache() Spark читал бы Kafka повторно для каждого .write
    deduped_df.cache()

    record_count = deduped_df.count()
    print(f"Батч {batch_id}: {record_count} уникальных событий после дедупликации")

    # 3. Запись в Bronze (append-only)
    # Все сырые события сохраняются в неизменяемый лог
    (
        deduped_df
        .withColumn("event_date", F.to_date("event_ts"))
        .writeTo("lakehouse.bronze.events")
        .append()
    )

    # 4. MERGE INTO Silver: обновление текущего состояния пользователей
    # Регистрируем временное представление для использования в SQL
    deduped_df.createOrReplaceTempView("incoming_events_view")

    spark.sql("""
        MERGE INTO lakehouse.silver.user_profiles AS target
        USING (
            SELECT
                user_id,
                SUM(amount) as session_revenue,
                COUNT(*) as session_events,
                MAX(event_ts) as last_seen_at
            FROM incoming_events_view
            WHERE event_type = 'purchase'
            GROUP BY user_id
        ) AS source
        ON target.user_id = source.user_id
        WHEN MATCHED THEN UPDATE SET
            target.total_revenue  = target.total_revenue + source.session_revenue,
            target.total_events   = target.total_events + source.session_events,
            target.last_seen_at   = source.last_seen_at,
            target.updated_at     = current_timestamp()
        WHEN NOT MATCHED THEN INSERT
            (user_id, total_revenue, total_events, last_seen_at, created_at, updated_at)
            VALUES (source.user_id, source.session_revenue, source.session_events,
                    source.last_seen_at, current_timestamp(), current_timestamp())
    """)

    # 5. Запись в Gold Realtime View (агрегаты для дашборда)
    (
        deduped_df
        .groupBy(
            F.window("event_ts", "1 minute").alias("time_window"),
            "event_type"
        )
        .agg(
            F.sum("amount").alias("total_amount"),
            F.count("*").alias("event_count")
        )
        .writeTo("lakehouse.gold.realtime_views")
        .append()
    )

    # 6. Фоновая компакция: каждые 60 батчей (~30 минут при 30-секундном триггере)
    if batch_id % 60 == 0:
        print(f"Батч {batch_id}: запускаем компакцию Bronze...")
        spark.sql("""
            CALL lakehouse.system.rewrite_data_files(
                table => 'lakehouse.bronze.events',
                strategy => 'sort',
                sort_order => 'event_date ASC, user_id ASC',
                options => map(
                    'target-file-size-bytes', '134217728',
                    'partial-progress.enabled', 'true'
                )
            )
        """)

    deduped_df.unpersist()


# Запускаем стриминговый пайплайн
streaming_query = (
    parsed_stream
    .writeStream
    .trigger(processingTime="30 seconds")
    .foreachBatch(process_microbatch)
    .option("checkpointLocation", "s3://my-bucket/checkpoints/lambda-speed-layer")
    .start()
)

streaming_query.awaitTermination()

Почему foreachBatch предпочтительнее стандартного writeStream.format(...)? Потому что он позволяет:

  • Выполнять MERGE INTO (недоступен в стандартном writeStream)
  • Писать одновременно в несколько таблиц (Bronze, Silver, Gold) за один микро-батч
  • Добавлять произвольную логику (компакцию, метрики, алерты)
  • Делать дедупликацию с кешированием до первого .count()

8. Гарантии доставки и контроль состояния

8.1 Exactly-Once Semantics: три уровня защиты

В распределённых системах «ровно один раз» - это не тривиально. Сеть падает, исполнители умирают, контейнеры рестартуют. Structured Streaming обеспечивает Exactly-Once через комбинацию трёх механизмов:

Checkpoint - это папка на S3, где Spark хранит прогресс обработки: для каждой партиции Kafka - последний успешно обработанный offset. Checkpoint обновляется после каждого успешного микро-батча. При рестарте Spark читает checkpoint и продолжает ровно с того места, где остановился.

Write-Ahead Log в Iceberg означает, что метаданные нового снапшота записываются на диск ДО того, как снапшот считается зафиксированным. Если процесс умирает во время записи данных - метаданные не обновились, и аналитик видит предыдущий консистентный снапшот. Никаких «половинных коммитов».

Идемпотентность - это защита от сценария: Spark успешно записал данные в Iceberg, но упал до обновления checkpoint. При рестарте Spark прочитает старый checkpoint и заново обработает тот же микро-батч. Без дедупликации данные продублируются. С dropDuplicates(["event_id"]) повторная обработка идемпотентна.

8.2 Управление состоянием: Watermark и State Store

Stateful операции (window aggregations, deduplication, stream-stream joins) требуют хранения промежуточного состояния между микро-батчами. Spark хранит это состояние в State Store - RocksDB или HDFS-backend.

Без watermark State Store накапливает состояние бесконечно: для каждого временного окна, которое когда-либо открывалось, Spark хранит промежуточный агрегат - ожидая, вдруг придут опоздавшие события. Через неделю работы State Store занимает сотни гигабайт и стриминг умирает от OOM.

# Критически важно: всегда устанавливайте watermark ПЕРЕД stateful операцией
# Watermark говорит Spark: "окна старше X минут от последнего события можно закрыть
# и выбросить из State Store"

revenue_by_window = (
    streaming_df
    # Шаг 1: watermark на колонке event_ts
    # "10 minutes" означает: опоздавшие события принимаем до 10 минут,
    # более старые игнорируем
    .withWatermark("event_ts", "10 minutes")
    # Шаг 2: stateful агрегация
    .groupBy(
        F.window("event_ts", "1 hour"),   # Часовые окна
        "event_type"
    )
    .agg(
        F.sum("amount").alias("revenue"),
        F.count("*").alias("tx_count")
    )
)
# Без watermark: State Store хранит ВСЕ окна с начала времён → OOM через N дней
# С watermark: Spark удаляет из State Store окна старше
#              (last_event_ts - 10 minutes - window_duration)

Практическое правило: watermark должен быть равен максимальному ожидаемому опозданию событий в вашей системе. Для мобильных приложений - 15-30 минут (пользователь мог быть офлайн). Для IoT-датчиков - 1-5 минут. Для Kafka-потока в датацентре - 1-2 минуты.

8.3 Мониторинг через Spark UI: диагностика Consumer Lag

Spark UI предоставляет вкладку «Structured Streaming» с ключевыми метриками:

# Программный доступ к метрикам стриминга
query = streaming_df.writeStream \
    .trigger(processingTime="30 seconds") \
    .foreachBatch(process_microbatch) \
    .option("checkpointLocation", "s3://bucket/checkpoints/speed-layer") \
    .start()

# После запуска - получаем метрики текущего состояния
status = query.status
progress = query.lastProgress

print("=== Streaming Metrics ===")
print(f"Статус: {status['message']}")
print(f"Триггер активен: {status['isTriggerActive']}")

if progress:
    # Input Rate: сколько сообщений поступает в Kafka в секунду
    input_rate = progress.get("inputRowsPerSecond", 0)
    # Process Rate: сколько сообщений Spark обрабатывает в секунду
    process_rate = progress.get("processedRowsPerSecond", 0)

    print(f"Input Rate:   {input_rate:.1f} строк/сек (приходит в Kafka)")
    print(f"Process Rate: {process_rate:.1f} строк/сек (обрабатывает Spark)")

    # КЛЮЧЕВОЙ СИГНАЛ: если Input Rate > Process Rate - Consumer Lag растёт
    if input_rate > process_rate * 1.1:  # 10% допустимое отклонение
        print("ПРЕДУПРЕЖДЕНИЕ: Consumer Lag растёт!")
        print("Причины: недостаточно ресурсов исполнителей, долгая foreachBatch, I/O S3")

    # Дополнительные метрики
    batch_duration_ms = progress.get("batchDuration", 0)
    num_input_rows = progress.get("numInputRows", 0)
    print(f"Длительность батча: {batch_duration_ms}ms")
    print(f"Строк в последнем батче: {num_input_rows}")

Consumer Lag - главный сигнал неблагополучия Speed Layer. Если скорость поступления событий в Kafka (Input Rate) превышает скорость их обработки Spark'ом (Process Rate), очередь накапливается. Через несколько часов lag составит миллионы сообщений, и «реальное время» скатится в «45 минут назад».

Причины роста Consumer Lag:

  • Недостаточно Spark executors - нужно добавить ресурсы
  • foreachBatch включает слишком долгий MERGE INTO - нужно оптимизировать
  • Медленный S3 I/O - нужно проверить пропускную способность
  • GC паузы на JVM - нужно настроить память executor'ов

9. Lambda vs Kappa: сравнение архитектур

9.1 Критика Lambda Architecture

Lambda Architecture решила проблему Batch-только систем, но породила новые. Джей Кребс (Jay Kreps, создатель Kafka) в 2014 году сформулировал главные претензии:

Дублирование кода. Даже если Spark позволяет использовать один DataFrame API, Batch Layer и Speed Layer всё равно требуют отдельных пайплайнов: разных конфигураций, разных trigger-стратегий, разных checkpoint-директорий. При масштабировании (десятки источников × десятки витрин) инженерная сложность растёт квадратично.

Проблема реконсиляции. Serving Layer должен «склеивать» данные из двух источников так, чтобы не было двойного счёта. Это нетривиальная логика, которая усложняется при добавлении временных зон, DST-переводов и late-arriving событий.

Latency Batch Layer. Batch пересчёт раз в 6 часов означает, что аналитик видит точные данные с задержкой до 6 часов. Для многих сценариев это приемлемо, но не для всех.

9.2 Kappa Architecture: стриминг как единственный источник правды

Kappa Architecture (Jay Kreps, 2014) предлагает радикальное упрощение: уберите Batch Layer. Только Speed Layer, только Kafka как Master Dataset, только стриминговая обработка.

Ключевая идея Kappa: иммутабельный Master Dataset - это сам Kafka-топик с retention forever (или очень долгим). Для «пересчёта» с новой бизнес-логикой запускается новый стриминговый джоб, читающий топик с самого начала (startingOffsets=earliest). После завершения пересчёта старый джоб заменяется новым.

Преимущества Kappa:

  • Единая кодовая база: один стриминговый пайплайн
  • Нет проблемы реконсиляции Serving Layer
  • Проще операционно: меньше джобов, меньше зависимостей

Ограничения Kappa:

  • Kafka как «вечный» storage - дорого при больших объёмах и длинных retention
  • Пересчёт большой истории через стриминг - медленно (часы/дни)
  • Сложные join с историческими данными ограничены в Structured Streaming
  • Stateful операции с большим state (годы данных) требуют огромного State Store

9.3 Практическое сравнение: что выбрать

Критерий Lambda Kappa Hybrid Lakehouse
Кодовая база Двойная Единая Единая (Spark API)
Сложность Serving Высокая Низкая Средняя
Latency Секунды-часы Секунды Секунды-минуты
Стоимость хранения S3 (дёшево) Kafka (дорого) S3 (дёшево)
Пересчёт истории Быстро (Batch) Медленно (Stream replay) Быстро (Batch + Iceberg)
ACID транзакции Нет (классика) Частично Да (Iceberg/Delta)
Сложность state Нет Высокая Средняя
Зрелость 2012 2014 2020+

Hybrid Lakehouse подход - то, что большинство modern data platforms реализуют сегодня:

  • S3 + Iceberg/Delta как Master Dataset (дёшево, надёжно, ACID)
  • Spark Batch для исторических пересчётов
  • Spark Structured Streaming для Speed Layer
  • Единый DataFrame API для обоих
  • Iceberg Snapshot Isolation для конкурентной записи

10. Анти-паттерны Lambda в production

10.1 Dual Pipeline Nightmare

Самый опасный анти-паттерн: Batch и Streaming пайплайны эволюционируют независимо и постепенно начинают давать разные результаты для одних и тех же данных.

Решение: единая функция трансформации, как показано в разделе 6.1. Один def compute_revenue_metrics(df), вызываемый и в Batch, и в Streaming контексте. Любое изменение применяется атомарно к обоим слоям.

10.2 Broken Checkpoint Recovery

Ещё один опасный сценарий: checkpoint устарел или стал несовместим со схемой данных.

# ПРОБЛЕМА: изменили схему входных данных, но checkpoint хранит старую схему
# Spark не может прочитать старый checkpoint с новой схемой → ошибка запуска

# СИМПТОМ: AnalysisException при старте стриминга после изменения схемы
# "Input schema changes are not supported in streaming. ..."

# АНТИПАТТЕРН: удалить checkpoint и начать заново
# Это приведёт к потере прогресса и к повторной обработке с нуля
# или к пропуску данных (если startingOffsets=latest)

# ПРАВИЛЬНОЕ РЕШЕНИЕ 1: Schema Evolution в Iceberg
# Добавляем колонку в таблицу ПРЕЖДЕ, чем менять стриминг
spark.sql("""
    ALTER TABLE lakehouse.bronze.events
    ADD COLUMN new_column STRING
    AFTER existing_column
""")
# Теперь стриминг может продолжить с тем же checkpoint

# ПРАВИЛЬНОЕ РЕШЕНИЕ 2: Миграция checkpoint
# Если схема изменилась несовместимо - создаём новый стриминговый джоб
# с новым checkpoint path и startingOffsets=earliest
# параллельно со старым, затем переключаем Serving Layer

10.3 State Explosion

При длинных временных окнах в Structured Streaming State Store может вырасти до неприемлемых размеров:

# АНТИПАТТЕРН: очень длинное окно без агрессивного watermark
problematic_query = (
    streaming_df
    .withWatermark("event_ts", "7 days")    # Spark хранит state 7 дней!
    .groupBy(F.window("event_ts", "30 days")) # Окно 30 дней × миллионы ключей
    .agg(F.sum("amount"))
    # Проблема: State Store = (30 days / trigger interval) × keys_count × state_size
    # Через неделю работы: гигабайты RAM на каждый executor
)

# ПРАВИЛЬНЫЙ ПОДХОД: разбить на несколько шагов
# Шаг 1: короткие окна в стриминге (watermark 10 минут, окно 1 час)
hourly_stream = (
    streaming_df
    .withWatermark("event_ts", "10 minutes")
    .groupBy(F.window("event_ts", "1 hour"))
    .agg(F.sum("amount").alias("hourly_revenue"))
    .writeStream.toTable("lakehouse.silver.hourly_revenue")
)

# Шаг 2: длинные агрегации в Batch Layer (читает Silver, агрегирует по 30 дням)
# Нет ограничений на длину периода, нет State Store
monthly_batch = (
    spark.table("lakehouse.silver.hourly_revenue")
         .groupBy(F.window("window.start", "30 days"))
         .agg(F.sum("hourly_revenue").alias("monthly_revenue"))
)

10.4 Checkpoint на локальной FS

# АНТИПАТТЕРН: checkpoint в /tmp или на локальном диске
bad_query = (
    streaming_df.writeStream
                .option("checkpointLocation", "/tmp/checkpoint")  # ОШИБКА!
                .start()
)
# При рестарте контейнера или перезапуске Spark на другом узле - checkpoint потерян
# Стриминг начнёт с нуля → потеря прогресса или дубликаты

# ПРАВИЛЬНО: checkpoint на S3 или HDFS
good_query = (
    streaming_df.writeStream
                .option("checkpointLocation", "s3://my-bucket/checkpoints/job-name-v1")
                .start()
)
# S3 доступен из любого узла кластера, переживает рестарты

# ВАЖНО: включайте версию в имя checkpoint
# При несовместимом изменении джоба - меняйте v1 на v2
# Это создаёт новый checkpoint (джоб стартует заново) без риска конфликта со старым

11. Производственный кейс: e-commerce Lambda Pipeline на Spark + Iceberg

11.1 Архитектура системы

Рассмотрим реальный сценарий: e-commerce платформа с 1 млн пользователей. Kafka получает 10 000 событий в секунду: просмотры товаров, добавления в корзину, покупки, отмены. Требования:

  • Live dashboard: обновление раз в 30 секунд, задержка < 1 минуты
  • Дневные отчёты: точные данные за любой день, любой период
  • Fraud detection: алерт в течение 10 секунд при аномальной активности

11.2 Полный код Speed Layer с Fraud Detection

from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.types import (StructType, StructField,
                                StringType, LongType, DoubleType, TimestampType)

# === Инициализация Spark Session ===
spark = SparkSession.builder \
    .appName("ECommerce-Lambda-SpeedLayer") \
    .config("spark.sql.extensions",
            "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
    .config("spark.sql.catalog.lakehouse", "org.apache.iceberg.spark.SparkCatalog") \
    .config("spark.sql.catalog.lakehouse.type", "rest") \
    .config("spark.sql.catalog.lakehouse.uri", "http://iceberg-catalog:8181") \
    .config("spark.sql.catalog.lakehouse.warehouse", "s3://my-bucket/warehouse") \
    # Оптимизации для Structured Streaming
    .config("spark.sql.streaming.stateStore.providerClass",
            "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider") \
    .config("spark.sql.streaming.stateStore.rocksdb.compactOnCommit", "true") \
    .getOrCreate()

# === Схема событий ===
event_schema = StructType([
    StructField("event_id",    StringType(),    False),
    StructField("user_id",     StringType(),    False),
    StructField("session_id",  StringType(),    True),
    StructField("event_type",  StringType(),    False),  # purchase, view, cart_add, cart_remove
    StructField("product_id",  StringType(),    True),
    StructField("amount",      DoubleType(),    True),
    StructField("discount_pct", DoubleType(),   True),   # процент скидки 0-100
    StructField("event_ts",    TimestampType(), False),
    StructField("ip_address",  StringType(),    True),
])

# === Читаем из Kafka ===
raw_stream = (
    spark.readStream
         .format("kafka")
         .option("kafka.bootstrap.servers", "kafka:9092")
         .option("subscribe", "ecommerce.purchases,ecommerce.views,ecommerce.carts")
         .option("startingOffsets", "latest")
         .option("maxOffsetsPerTrigger", "100000")  # защита от OOM при большом lag
         .option("kafka.group.id", "spark-lambda-speed-layer")
         .load()
         .select(
             F.col("topic"),
             F.col("partition"),
             F.col("offset"),
             F.col("timestamp").alias("kafka_ingestion_ts"),
             F.from_json(F.col("value").cast("string"), event_schema).alias("d")
         )
         .select("topic", "partition", "offset", "kafka_ingestion_ts", "d.*")
)


# === Общая функция обогащения (используется в Batch И в Streaming) ===
def enrich_events(df):
    """
    Обогащение событий: одинаковая логика для Batch и Speed Layer.
    Добавляет бизнес-правила: фильтрация невалидных записей,
    категоризация, расчёт net_amount.
    """
    return (
        df
        # Фильтрация: только валидные суммы для финансовых событий
        .filter(
            (F.col("event_type") != "purchase") | (F.col("amount") > 0)
        )
        # Бизнес-правило: скидки > 50% не учитываются в выручке
        # Это применяется ОДИНАКОВО и в Batch, и в Speed
        .withColumn(
            "net_amount",
            F.when(
                (F.col("event_type") == "purchase") & (F.col("discount_pct") <= 50),
                F.col("amount") * (1 - F.col("discount_pct") / 100)
            ).otherwise(F.lit(0.0))
        )
        # Категоризация суммы (для сегментации аналитики)
        .withColumn(
            "amount_tier",
            F.when(F.col("amount") < 500,   "micro")
             .when(F.col("amount") < 5000,  "mid")
             .when(F.col("amount") < 50000, "large")
             .otherwise("enterprise")
        )
    )


# === foreachBatch функция ===
def process_ecommerce_batch(batch_df, batch_id: int) -> None:
    """
    Центральная функция обработки каждого микро-батча Speed Layer.
    Выполняет запись в 3 слоя Lakehouse параллельно.
    """
    if batch_df.isEmpty():
        return

    # 1. Дедупликация по event_id (защита от at-least-once Kafka)
    deduped = batch_df.dropDuplicates(["event_id"]).cache()
    count = deduped.count()

    if count == 0:
        deduped.unpersist()
        return

    print(f"[Batch {batch_id}] Обрабатываем {count} уникальных событий")

    # 2. Обогащаем (та же функция, что и в Batch Layer)
    enriched = enrich_events(deduped)

    # 3. Bronze: пишем сырые события (append-only)
    # Сохраняем ВСЕ поля, включая topic, partition, offset для дебаггинга
    (
        enriched
        .withColumn("event_date", F.to_date("event_ts"))
        .writeTo("lakehouse.bronze.ecommerce_events")
        .append()
    )

    # 4. Silver: обновляем профили пользователей через MERGE INTO
    # Только для событий типа purchase
    enriched.filter(F.col("event_type") == "purchase") \
            .createOrReplaceTempView("speed_purchases")

    spark.sql("""
        MERGE INTO lakehouse.silver.user_purchase_profiles AS t
        USING (
            SELECT
                user_id,
                COUNT(*)         AS purchase_count,
                SUM(net_amount)  AS total_net_revenue,
                MAX(event_ts)    AS last_purchase_at,
                MAX(amount_tier) AS highest_tier
            FROM speed_purchases
            GROUP BY user_id
        ) AS s
        ON t.user_id = s.user_id
        WHEN MATCHED THEN UPDATE SET
            t.purchase_count     = t.purchase_count + s.purchase_count,
            t.total_net_revenue  = t.total_net_revenue + s.total_net_revenue,
            t.last_purchase_at   = s.last_purchase_at,
            t.updated_at         = current_timestamp()
        WHEN NOT MATCHED THEN INSERT *
    """)

    # 5. Gold Realtime View: агрегаты для live дашборда
    (
        enriched
        .groupBy(
            F.date_trunc("minute", F.col("event_ts")).alias("minute_ts"),
            "event_type",
            "amount_tier"
        )
        .agg(
            F.sum("net_amount").alias("net_revenue"),
            F.count("*").alias("event_count"),
            F.approx_count_distinct("user_id").alias("unique_users"),
            F.approx_count_distinct("session_id").alias("unique_sessions")
        )
        .writeTo("lakehouse.gold.realtime_revenue_by_minute")
        .append()
    )

    # 6. Фоновая компакция каждые 120 батчей (~1 час при 30-секундном триггере)
    if batch_id % 120 == 0:
        try:
            spark.sql("""
                CALL lakehouse.system.rewrite_data_files(
                    table => 'lakehouse.bronze.ecommerce_events',
                    where => 'event_date = current_date()',
                    options => map(
                        'target-file-size-bytes', '134217728',
                        'max-concurrent-file-group-rewrites', '4'
                    )
                )
            """)
            print(f"[Batch {batch_id}] Компакция Bronze завершена")
        except Exception as ex:
            # Компакция не критична: если упала, логируем и продолжаем
            print(f"[Batch {batch_id}] Компакция не удалась: {ex}")

    deduped.unpersist()


# === Fraud Detection: отдельный стриминговый джоб с коротким триггером ===
def detect_fraud_batch(batch_df, batch_id: int) -> None:
    """
    Обнаружение аномалий: пользователи с > 10 покупками за 2 минуты
    или с суммой > 100 000 руб за 5 минут.
    """
    if batch_df.isEmpty():
        return

    batch_df.createOrReplaceTempView("recent_purchases")

    # Детектируем аномальную активность
    alerts = spark.sql("""
        SELECT
            user_id,
            COUNT(*) as purchase_count_2min,
            SUM(amount) as total_amount_2min,
            CASE
                WHEN COUNT(*) > 10             THEN 'HIGH_FREQUENCY'
                WHEN SUM(amount) > 100000      THEN 'HIGH_AMOUNT'
                ELSE 'NORMAL'
            END as alert_type,
            current_timestamp() as detected_at
        FROM recent_purchases
        WHERE event_type = 'purchase'
          AND event_ts >= current_timestamp() - INTERVAL 2 MINUTES
        GROUP BY user_id
        HAVING COUNT(*) > 10 OR SUM(amount) > 100000
    """)

    alert_count = alerts.count()
    if alert_count > 0:
        print(f"[Fraud Batch {batch_id}] АЛЕРТ: {alert_count} подозрительных пользователей!")
        (
            alerts
            .writeTo("lakehouse.gold.fraud_alerts")
            .append()
        )


# === Запуск Speed Layer ===
main_query = (
    raw_stream
    .writeStream
    .trigger(processingTime="30 seconds")
    .foreachBatch(process_ecommerce_batch)
    .option("checkpointLocation", "s3://my-bucket/checkpoints/ecommerce-speed-v2")
    .start()
)

# Fraud detection - отдельный джоб с более частым триггером
fraud_query = (
    raw_stream
    .filter(F.col("event_type") == "purchase")
    .withWatermark("event_ts", "5 minutes")
    .writeStream
    .trigger(processingTime="5 seconds")
    .foreachBatch(detect_fraud_batch)
    .option("checkpointLocation", "s3://my-bucket/checkpoints/fraud-detection-v1")
    .start()
)

# Ждём завершения (в production работает постоянно)
spark.streams.awaitAnyTermination()

11.3 Batch Layer: исторический пересчёт

# === Batch Layer: ежечасный инкрементальный пересчёт ===
# Запускается по расписанию (Apache Airflow, Prefect, cron)

from datetime import datetime, timedelta

def run_hourly_batch_job(target_hour: datetime) -> None:
    """
    Инкрементальный Batch Job: обрабатывает данные за конкретный час.
    Используется для ежечасного обновления Batch Views.

    target_hour: datetime объект часа для обработки (напр., 2026-05-23 14:00)
    """
    hour_start = target_hour.replace(minute=0, second=0, microsecond=0)
    hour_end   = hour_start + timedelta(hours=1)

    print(f"Batch Job: обрабатываем {hour_start} - {hour_end}")

    # Читаем из Bronze только нужный час (Partition Pruning по event_date + фильтр)
    bronze_hour = (
        spark.table("lakehouse.bronze.ecommerce_events")
             .filter(F.col("event_date") == hour_start.date())       # Partition Pruning
             .filter(F.col("event_ts").between(hour_start, hour_end)) # Точный диапазон
    )

    # Применяем ту же функцию обогащения (единая бизнес-логика!)
    enriched_hour = enrich_events(bronze_hour)

    # Агрегация за час
    hourly_aggregates = (
        enriched_hour
        .groupBy(
            F.date_trunc("hour", F.col("event_ts")).alias("hour_ts"),
            "event_type",
            "amount_tier"
        )
        .agg(
            F.sum("net_amount").alias("net_revenue"),
            F.count("*").alias("event_count"),
            F.countDistinct("user_id").alias("unique_users"),
            F.countDistinct("session_id").alias("unique_sessions"),
            F.lit("batch").alias("source_layer")
        )
    )

    # Записываем Batch View: перезаписываем партицию нужного часа
    # Это идемпотентно: повторный запуск за тот же час даст тот же результат
    (
        hourly_aggregates
        .writeTo("lakehouse.gold.batch_revenue_by_hour")
        .overwritePartitions()
    )

    # Обновляем метаданные: до какого момента покрыт Batch Layer
    spark.sql(f"""
        MERGE INTO lakehouse.gold.batch_view_metadata AS t
        USING (SELECT '{hour_end}' as processed_up_to_ts) AS s
        ON 1 = 1
        WHEN MATCHED AND t.processed_up_to_ts < s.processed_up_to_ts
            THEN UPDATE SET t.processed_up_to_ts = s.processed_up_to_ts,
                            t.updated_at = current_timestamp()
        WHEN NOT MATCHED THEN INSERT *
    """)

    print(f"Batch Job завершён: {hour_start} - {hour_end}")
    print(f"Batch Views покрывают данные до: {hour_end}")


# === Полный исторический пересчёт (при исправлении бизнес-логики) ===
def run_full_historical_recompute(start_date: str, end_date: str) -> None:
    """
    Полный пересчёт Batch Views за исторический период.
    Запускается при изменении бизнес-логики (например, пересмотр правила скидок).

    Ключевой принцип Lambda: Bronze нетронут, пересчитываем только Gold Batch Views.
    """
    print(f"Исторический пересчёт: {start_date} - {end_date}")
    print("Bronze Layer не изменяется - только Gold Batch Views пересчитываются")

    # Читаем весь исторический Bronze за период
    bronze_historical = (
        spark.table("lakehouse.bronze.ecommerce_events")
             .filter(F.col("event_date").between(start_date, end_date))
    )

    # Применяем ТЕКУЩУЮ (исправленную) бизнес-логику
    enriched_historical = enrich_events(bronze_historical)  # новая версия функции

    # Агрегируем по дням для исторических отчётов
    daily_aggregates = (
        enriched_historical
        .groupBy(
            F.to_date("event_ts").alias("date"),
            "event_type"
        )
        .agg(
            F.sum("net_amount").alias("net_revenue"),
            F.count("*").alias("event_count"),
            F.countDistinct("user_id").alias("unique_users")
        )
    )

    # Перезаписываем все партиции за период
    (
        daily_aggregates
        .writeTo("lakehouse.gold.batch_revenue_by_day")
        .overwritePartitions()
    )

    print(f"Исторический пересчёт завершён: {daily_aggregates.count()} строк записано")

11.4 Serving Layer: финальная сборка для дашборда

# === Serving Layer: объединение Batch и Speed ===
# Этот запрос используется Live Dashboard (обновляется каждые 30 секунд)

def get_dashboard_data(period_start: str, period_end: str):
    """
    Возвращает данные для Live Dashboard за период.
    Автоматически объединяет точные Batch данные и свежие Speed данные.
    Нет двойного счёта.
    """
    # Определяем границу Batch Layer
    batch_meta = spark.table("lakehouse.gold.batch_view_metadata")
    batch_max_ts_row = batch_meta.select(F.max("processed_up_to_ts")).first()
    batch_max_ts = batch_max_ts_row[0] if batch_max_ts_row[0] else "2000-01-01"

    print(f"Dashboard: Batch покрывает до {batch_max_ts}, Speed - остальное")

    # Batch данные: точные исторические агрегаты (до batch_max_ts)
    batch_data = (
        spark.table("lakehouse.gold.batch_revenue_by_hour")
             .filter(F.col("hour_ts").between(period_start, period_end))
             .filter(F.col("hour_ts") < F.lit(batch_max_ts))
             .select(
                 "hour_ts", "event_type", "amount_tier",
                 "net_revenue", "event_count", "unique_users",
                 F.lit("batch").alias("data_source")
             )
    )

    # Speed данные: свежие агрегаты (начиная с batch_max_ts)
    speed_data = (
        spark.table("lakehouse.gold.realtime_revenue_by_minute")
             .filter(F.col("minute_ts").between(period_start, period_end))
             .filter(F.col("minute_ts") >= F.lit(batch_max_ts))
             .withColumn("hour_ts", F.date_trunc("hour", F.col("minute_ts")))
             .groupBy("hour_ts", "event_type", "amount_tier")
             .agg(
                 F.sum("net_revenue").alias("net_revenue"),
                 F.sum("event_count").alias("event_count"),
                 F.sum("unique_users").alias("unique_users")  # приблизительно
             )
             .withColumn("data_source", F.lit("speed"))
    )

    # Объединяем (непересекающиеся периоды, нет двойного счёта)
    unified = (
        batch_data.unionByName(speed_data, allowMissingColumns=True)
                  .orderBy("hour_ts", "event_type")
    )

    return unified


# Пример вызова для дашборда
dashboard_df = get_dashboard_data(
    period_start="2026-05-23 00:00:00",
    period_end="2026-05-23 23:59:59"
)

dashboard_df.show(truncate=False)
# hour_ts             | event_type | net_revenue | data_source
# 2026-05-23 00:00:00 | purchase   | 45 000.00   | batch
# 2026-05-23 01:00:00 | purchase   | 38 500.00   | batch
# ...
# 2026-05-23 14:00:00 | purchase   | 72 100.00   | batch
# 2026-05-23 15:00:00 | purchase   | 12 300.00   | speed  <- свежие данные
# 2026-05-23 15:30:00 | purchase   | 5 800.00    | speed  <- за последние минуты

12. Итоги: Lambda Architecture как фундаментальная идея

Lambda Architecture - это не «устаревший Hadoop-паттерн», а фундаментальная идея интеграции historical correctness и real-time analytics. Конкретные технологии менялись (Hadoop → Spark, Hive → Iceberg, Storm → Structured Streaming), но архитектурный принцип остался:

Главный takeaway для Senior Data Engineer: Lambda Architecture - это не техническое решение, а архитектурная философия. Технологии реализации будут меняться. Появятся новые форматы хранения, новые движки. Но принцип «иммутабельный источник истины + два слоя обработки с разными trade-offs» останется актуальным, пока существует конфликт между correctness и freshness данных.

Следующий урок рассмотрит Kappa Architecture - альтернативный подход, который ставит под сомнение необходимость Batch Layer и предлагает «streaming-first» мышление.