Kappa-архитектура на Spark Structured Streaming: стриминг как единый источник правды

Kappa-архитектура на Spark Structured Streaming: стриминг как единый источник правды

datamodeling

В предыдущем уроке мы разобрали Lambda Architecture - элегантное решение проблемы «точность vs скорость» через разделение на два слоя обработки. Но у Lambda есть системный недостаток: любое изменение бизнес-логики требует синхронного обновления двух кодовых баз, а Serving Layer несёт сложную логику объединения без двойного счёта.

Kappa Architecture, предложенная Джеем Кребсом (Jay Kreps, создатель Apache Kafka) в 2014 году, ставит радикальный вопрос: а зачем вообще Batch Layer? Если Kafka хранит события долго, а Spark умеет читать топик с самого начала - не является ли «пересчёт батчем» просто частным случаем «воспроизведения потока»? Этот вопрос привёл к появлению архитектуры, где весь стек - от инжестии до витрин - строится на едином стриминговом пайплайне.


1. Философия Kappa: «Всё есть поток»

1.1 Переосмысление природы батча

В Lambda Architecture батч и стриминг считаются принципиально разными вещами: разные движки, разные модели исполнения, разный код. Kappa предлагает другую интерпретацию:

Батч - это поток с конечными границами.

Подумайте об этом: чем отличается «обработка данных за 1 января» от «обработки событий, поступивших между midnight и 23:59»? Только тем, что у батча известна правая граница. Если стриминговый движок умеет «обрабатывать события от offset 0 до offset N» - это и есть батч. После достижения offset N он продолжает работу в real-time режиме.

Диаграмма показывает смещение парадигмы: в традиционном взгляде batch и stream - это разные миры с разными технологиями. В Kappa оба «мира» - это один и тот же Kafka-лог, просто читаемый с разных позиций. Один и тот же Spark Streaming Job сначала «догоняет» историю (читает с offset 0), затем переходит в режим постоянной обработки новых событий.

1.2 Три столпа Kappa Architecture

Immutable Event Log - это сам Kafka (или совместимый брокер: Redpanda, Apache Pulsar). Ключевое отличие от классического использования Kafka: в Kappa настраивается долгосрочное или бесконечное хранение. Kafka больше не является «шиной» с retention в несколько дней - он становится первичным хранилищем событий. Любой consumer может читать с любого offset в любое время.

Unified Processing Engine - Spark Structured Streaming с единственным кодом, который работает как на исторических данных (replay от offset 0), так и на текущем потоке. Никаких if batch_mode else stream_mode. Одна функция, один пайплайн.

ACID Serving Layer - Iceberg или Delta Lake. Streaming Job непрерывно пишет materialized views: текущие состояния, агрегаты, сессионизированные данные. Аналитики делают SQL-запросы к этим таблицам, не зная ничего о том, что под капотом - стриминг.

1.3 Lambda vs Kappa: принципиальное различие

Аспект Lambda Kappa
Кодовых баз 2 (batch + stream) 1 (stream)
Source of Truth Master Dataset (S3/HDFS) Event Log (Kafka)
Пересчёт Batch Job с нуля Stream Replay от earliest
Serving Layer Объединение Batch+Speed View Только Streaming Materialized Views
Задержка данных Секунды (Speed) + часы (Batch) Секунды (единый пайплайн)
Стоимость хранения S3 дёшево Kafka дороже S3
Replay скорость Быстро (Batch full scan) Зависит от throughput кластера
Оперативная сложность Высокая (два стека) Средняя (один стек, но stateful)
Отладка Проще (детерминированный batch) Сложнее (stateful stream)

Вывод: Kappa уменьшает дублирование кода и упрощает архитектуру Serving Layer, но предъявляет более высокие требования к Kafka-инфраструктуре (долгое хранение дорого) и к команде (stateful streaming сложнее операционно).


2. Event Log как фундаментальный источник правды

2.1 Kafka как первичное хранилище

В классическом использовании Kafka - это шина сообщений: данные хранятся 3–7 дней, после чего удаляются. Kafka считается «транспортным» слоем, а не хранилищем. В Kappa Architecture этот взгляд меняется кардинально.

Kafka становится первичным хранилищем - equivalent HDFS в Lambda. Но вместо файлов хранятся события, вместо партиций HDFS - партиции Kafka-топика. Ретентшн настраивается на месяцы или на «никогда» (size-based retention).

2.2 Два режима хранения: Time-based и Log Compaction

Kafka предоставляет два механизма хранения, и в Kappa Architecture оба используются для разных типов топиков.

Time-based Retention - хранение всех событий за фиксированный период. Подходит для событийных потоков (кликстрим, транзакции, логи), где нужна полная история всех событий за последние N дней. Старые события удаляются по истечении retention.ms.

Log Compaction - Kafka хранит только последнее значение для каждого ключа сообщения. Старые значения для того же ключа удаляются (compacted out). Подходит для «состояний» сущностей: профили пользователей, статусы заказов, конфигурации сервисов. Лог компактен в смысле объёма, но потребитель видит актуальное состояние каждой сущности.

Для Kappa Architecture типичная конфигурация топиков:

  • Событийные топики (purchases, views, clicks): retention.ms = 31536000000 (365 дней), cleanup.policy = delete
  • Сущностные топики (user_profiles, product_catalog): cleanup.policy = compact, min.compaction.lag.ms = 86400000

2.3 Настройка Kafka для Kappa: конкретные параметры

# Создание событийного топика с годовым retention (для Kappa)
kafka-topics.sh --create \
  --bootstrap-server kafka:9092 \
  --topic ecommerce.purchases \
  --partitions 12 \           # Параллелизм = количество партиций
  --replication-factor 3 \    # Надёжность: 3 реплики
  --config retention.ms=31536000000 \   # 365 дней
  --config segment.bytes=536870912 \    # Сегмент 512MB (оптимально для долгого хранения)
  --config compression.type=lz4         # Сжатие: lz4 быстрее gzip, лучше чем none

# Создание компактного топика для состояний сущностей
kafka-topics.sh --create \
  --bootstrap-server kafka:9092 \
  --topic ecommerce.user_profiles \
  --partitions 12 \
  --replication-factor 3 \
  --config cleanup.policy=compact \
  --config min.compaction.lag.ms=3600000 \   # Compaction не раньше чем через час
  --config max.compaction.lag.ms=86400000    # Compaction не позже чем через сутки

Почему важен segment.bytes: Kafka хранит данные в сегментах. При удалении по retention удаляются целые сегменты, а не отдельные сообщения. Если сегмент 1GB и retention 7 дней, но за 7 дней накопилось 500MB - сегмент ещё не закрыт, и retention не срабатывает. Для долгосрочного хранения лучше делать сегменты меньше (256-512MB), чтобы ротация работала предсказуемо.


3. Spark Structured Streaming как основа Kappa Pipeline

3.1 Replay: как один пайплайн обрабатывает историю и real-time

Главный технический механизм Kappa Architecture - replay (воспроизведение). Когда нужно пересчитать данные с новой бизнес-логикой, не запускается отдельный Batch Job - запускается тот же Streaming Job, но с startingOffsets = "earliest".

Sequence Diagram объясняет процедуру migration без downtime:

  1. V1 работает в production, обрабатывает real-time события, пишет в gold_revenue_v1. Аналитики читают эту таблицу.
  2. V2 стартует параллельно с startingOffsets="earliest" и новым checkpoint-путём. Kafka хранит полную историю (365 дней), V2 читает её на максимальной скорости кластера - это и есть «батч» в Kappa.
  3. V2 догоняет real-time: скорость чтения Kafka намного выше скорости поступления новых событий, поэтому V2 быстро сокращает отставание и начинает читать те же текущие offsets, что и V1.
  4. Атомарное переключение: ALTER VIEW в Iceberg меняет указатель с V1-таблицы на V2-таблицу за один атомарный коммит. Аналитики переключаются мгновенно, без 404.
  5. V1 останавливается: больше не нужен.

3.2 Управление смещениями: startingOffsets как рычаг Kappa

В Spark Structured Streaming параметр startingOffsets определяет, откуда начинать чтение при первом запуске. После первого запуска позиция чтения хранится в checkpoint и startingOffsets игнорируется.

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

spark = SparkSession.builder \
    .appName("Kappa-Architecture-Pipeline") \
    .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://catalog:8181") \
    .config("spark.sql.catalog.lakehouse.warehouse", "s3://bucket/warehouse") \
    # RocksDB State Store: гораздо эффективнее default HashMapStateStore
    # для больших stateful операций (деdup, windowed agg)
    .config("spark.sql.streaming.stateStore.providerClass",
            "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider") \
    # Компактим RocksDB при каждом commit для предотвращения раздувания
    .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),
    StructField("product_id", StringType(),    True),
    StructField("amount",     DoubleType(),    True),
    StructField("event_ts",   TimestampType(), False),
])

# === РЕЖИМ 1: Начальный Replay (пересчёт истории) ===
# startingOffsets="earliest": читаем Kafka с самого начала
# Используется при:
# - первом запуске нового пайплайна
# - смене бизнес-логики (новая версия V2)
# - восстановлении после потери checkpoint
REPLAY_MODE = True   # Управляется конфигурацией или ENV-переменной

starting_offsets = "earliest" if REPLAY_MODE else "latest"

# Для точного управления можно передать JSON со смещениями по партициям:
# starting_offsets = json.dumps({
#     "ecommerce.purchases": {"0": 0, "1": 0, "2": 0}  # offset 0 в каждой партиции
# })

raw_kafka_stream = (
    spark.readStream
         .format("kafka")
         .option("kafka.bootstrap.servers", "kafka:9092")
         .option("subscribe", "ecommerce.purchases")
         .option("startingOffsets", starting_offsets)
         # В Replay режиме снимаем ограничение: читаем как можно быстрее
         # В production режиме ставим 50000 для защиты от OOM
         .option("maxOffsetsPerTrigger", "500000" if REPLAY_MODE else "50000")
         # failOnDataLoss=false для replay: старые сегменты могут быть удалены,
         # если retention не настроен достаточно долго
         .option("failOnDataLoss", "false" if REPLAY_MODE else "true")
         .load()
         .select(
             F.col("partition"),
             F.col("offset"),
             F.col("timestamp").alias("kafka_ts"),
             F.from_json(F.col("value").cast("string"), event_schema).alias("d")
         )
         .select("partition", "offset", "kafka_ts", "d.*")
)

maxOffsetsPerTrigger в Replay режиме можно ставить значительно выше, чем в real-time режиме - 500 000 или даже 1 000 000 сообщений за микро-батч. Ограничений по latency нет (мы обрабатываем историю, не real-time), поэтому единственное ограничение - объём памяти executor'ов. Это позволяет «догнать» год истории за часы, а не дни.

3.3 Event Time vs Processing Time: критическое различие для Kappa

В Kappa Architecture все агрегации строятся по Event Time - времени, когда событие произошло на стороне клиента или источника. Использование Processing Time (времени поступления в Spark) сделало бы replay бессмысленным: повторная обработка исторических событий дала бы другое время, а значит другие результаты.

Processing Time в Kappa - это архитектурная ошибка. Если группировать события по F.current_timestamp() или F.now(), каждый запуск Replay будет давать другие временны́е окна. «Выручка за 10:00–11:00 23 мая» при Replay превратится в «выручка за 14:00–15:00» (время запуска Replay). Историческая аналитика становится бессмысленной.

Event Time гарантирует детерминизм: какой бы момент мы ни запустили Replay, события с event_ts = 2026-05-23 10:00 всегда попадут в окно 10:00–11:00. Результат воспроизводим.


4. Stateful Processing: время, водяные знаки и оконные функции

4.1 Механизм Watermarking в деталях

Watermark (водяной знак) - это граница, ниже которой Spark перестаёт ждать опоздавших событий. Формально: если последнее увиденное событие имело event_ts = T, то watermark равен T - watermark_delay. Все события с event_ts < watermark считаются опоздавшими и игнорируются.

Ключевое свойство watermark: он решает две проблемы одновременно.

Первая - корректность: Spark знает, когда временно́е окно «закрыто» и его результат финален. Окно 10:00–10:10 становится финальным, когда watermark превышает 10:10. До этого момента Spark держит окно «открытым» - готов принять опоздавшие события в это окно.

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

Компромисс: чем больше watermark_delay, тем больше опоздавших событий принимается (лучше для completeness), но тем дольше Spark держит окна открытыми (больше памяти) и тем выше latency финальных результатов. Для производственных систем рекомендуем: max expected late arrival × 1.5.

4.2 Tumbling и Sliding Windows: реализация на PySpark

Tumbling Window (кувыркающееся / фиксированное окно) - непересекающиеся окна фиксированного размера. Каждое событие попадает ровно в одно окно. Используется для «метрик за период»: выручка за час, количество заказов за день.

Sliding Window (скользящее окно) - окна перекрываются. Одно событие может попадать сразу в несколько окон. Используется для «скользящих средних»: среднее количество событий за последние 15 минут, обновляемое каждые 5 минут.

from pyspark.sql import functions as F

# === Tumbling Window (фиксированное, непересекающееся) ===
# Каждая покупка попадает ровно в один час.
# Результат: выручка за час - обновляется каждые 30 секунд (trigger interval)
tumbling_revenue = (
    raw_kafka_stream
    # ОБЯЗАТЕЛЬНО: watermark ПЕРЕД stateful операцией
    # 10 минут - максимальное ожидаемое опоздание покупок
    .withWatermark("event_ts", "10 minutes")
    .groupBy(
        # window(timeColumn, windowDuration)
        # Размер окна: 1 час. Смещение: нет (окна по UTC-часам: 10:00-11:00, 11:00-12:00)
        F.window("event_ts", "1 hour"),
        "event_type"
    )
    .agg(
        F.sum("amount").alias("revenue"),
        F.count("*").alias("tx_count"),
        F.approx_count_distinct("user_id", rsd=0.01).alias("unique_buyers")
    )
    # Разворачиваем структуру window {start, end} в отдельные колонки
    .select(
        F.col("window.start").alias("hour_start"),
        F.col("window.end").alias("hour_end"),
        "event_type", "revenue", "tx_count", "unique_buyers"
    )
)

# === Sliding Window (скользящее, пересекающееся) ===
# Окно 15 минут, смещение 5 минут: 10:00-10:15, 10:05-10:20, 10:10-10:25...
# Каждое событие попадает в 3 окна (15/5 = 3)
# Используется для: скользящее среднее, moving KPI

sliding_activity = (
    raw_kafka_stream
    .withWatermark("event_ts", "20 minutes")   # Больше delay = больше late data принимаем
    .groupBy(
        # window(timeColumn, windowDuration, slideDuration)
        # Окно 15 минут, обновляется каждые 5 минут
        F.window("event_ts", "15 minutes", "5 minutes"),
        "event_type"
    )
    .agg(
        F.count("*").alias("event_count"),
        F.sum("amount").alias("window_revenue")
    )
    .select(
        F.col("window.start").alias("window_start"),
        F.col("window.end").alias("window_end"),
        "event_type", "event_count", "window_revenue"
    )
)

Стоимость Sliding Window: каждое событие обрабатывается windowDuration / slideDuration раз. При окне 15 минут и слайде 5 минут каждое событие обновляет 3 окна. State Store хранит в 3 раза больше состояния, чем для Tumbling Window того же периода. Это нужно учитывать при capacity planning.

4.3 Сессионизация: группировка событий в сессии пользователя

Одна из самых важных stateful операций в e-commerce аналитике - сессионизация: объединение последовательных действий пользователя в «сессию» с максимальным gap между событиями (например, 30 минут без активности = конец сессии).

В Spark Structured Streaming это реализуется через Session Window (Spark 3.2+) или через кастомный flatMapGroupsWithState.

from pyspark.sql.streaming.state import GroupState, GroupStateTimeout
from pyspark.sql.types import (StructType, StructField,
                                StringType, TimestampType, LongType, IntegerType)

# Схема состояния сессии (хранится в State Store для каждого user_id)
session_state_schema = StructType([
    StructField("session_id",     StringType(),    False),
    StructField("user_id",        StringType(),    False),
    StructField("session_start",  TimestampType(), False),
    StructField("last_event_ts",  TimestampType(), False),
    StructField("event_count",    IntegerType(),   False),
    StructField("total_amount",   DoubleType(),    False),
])

# Схема выходного события (когда сессия завершается)
session_output_schema = StructType([
    StructField("session_id",    StringType(),    False),
    StructField("user_id",       StringType(),    False),
    StructField("session_start", TimestampType(), False),
    StructField("session_end",   TimestampType(), False),
    StructField("duration_sec",  LongType(),      False),
    StructField("event_count",   IntegerType(),   False),
    StructField("total_amount",  DoubleType(),    False),
])

import uuid
from datetime import datetime, timedelta

SESSION_TIMEOUT_MINUTES = 30   # Сессия завершается через 30 минут тишины

def update_session_state(user_id: str, events, state: GroupState):
    """
    flatMapGroupsWithState функция: вызывается для каждого user_id
    при каждом новом микро-батче.

    Логика:
    - Если есть открытая сессия и новые события укладываются в timeout - расширяем сессию
    - Если timeout истёк - закрываем сессию, возвращаем её как результат, открываем новую
    - При истечении GroupState timeout - тоже закрываем сессию
    """
    SESSION_TIMEOUT = timedelta(minutes=SESSION_TIMEOUT_MINUTES)
    sorted_events = sorted(events, key=lambda e: e.event_ts)

    # Обработка timeout (GroupState timeout истёк - пользователь молчал долго)
    if state.hasTimedOut:
        if state.exists:
            # Сессия была открыта - закрываем её и возвращаем как результат
            s = state.get
            yield (s.session_id, s.user_id, s.session_start,
                   s.last_event_ts, int((s.last_event_ts - s.session_start).total_seconds()),
                   s.event_count, s.total_amount)
        state.remove()
        return

    completed_sessions = []
    current_state = state.get if state.exists else None

    for event in sorted_events:
        if current_state is None:
            # Нет открытой сессии - создаём новую
            current_state = {
                "session_id": str(uuid.uuid4()),
                "user_id": user_id,
                "session_start": event.event_ts,
                "last_event_ts": event.event_ts,
                "event_count": 1,
                "total_amount": event.amount or 0.0,
            }
        else:
            gap = event.event_ts - current_state["last_event_ts"]
            if gap > SESSION_TIMEOUT:
                # Gap слишком большой - закрываем текущую сессию
                completed_sessions.append(current_state)
                # Открываем новую сессию
                current_state = {
                    "session_id": str(uuid.uuid4()),
                    "user_id": user_id,
                    "session_start": event.event_ts,
                    "last_event_ts": event.event_ts,
                    "event_count": 1,
                    "total_amount": event.amount or 0.0,
                }
            else:
                # Продолжаем текущую сессию
                current_state["last_event_ts"] = event.event_ts
                current_state["event_count"] += 1
                current_state["total_amount"] += event.amount or 0.0

    # Сохраняем открытую сессию в State Store
    if current_state:
        state.update((
            current_state["session_id"],
            current_state["user_id"],
            current_state["session_start"],
            current_state["last_event_ts"],
            current_state["event_count"],
            current_state["total_amount"],
        ))
        # Устанавливаем timeout: если 30 минут не будет событий - закроем сессию
        state.setTimeoutDuration(f"{SESSION_TIMEOUT_MINUTES} minutes")

    # Возвращаем все завершённые сессии этого микро-батча
    for s in completed_sessions:
        yield (s["session_id"], s["user_id"], s["session_start"],
               s["last_event_ts"],
               int((s["last_event_ts"] - s["session_start"]).total_seconds()),
               s["event_count"], s["total_amount"])


# Применяем сессионизацию
sessions_stream = (
    raw_kafka_stream
    .withWatermark("event_ts", "30 minutes")
    .groupBy("user_id")
    .applyInPandasWithState(
        update_session_state,
        outputStructType=session_output_schema,
        stateStructType=session_state_schema,
        outputMode="append",
        timeoutConf=GroupStateTimeout.EventTimeTimeout
    )
)

Важное замечание о flatMapGroupsWithState: эта операция требует тщательного проектирования State Schema. Если в будущем потребуется добавить поле в session_state_schema, нужно либо создавать новый checkpoint (потеря state), либо реализовывать миграцию state вручную. Это один из главных операционных рисков тяжёлой stateful логики в Kappa.


5. Replayability: батч через стриминговый Replay

5.1 Почему Replay - это не просто «запустить заново»

В Lambda Architecture «пересчёт» означает запустить Spark Batch Job, который делает spark.read.parquet(...) и обрабатывает всё с нуля. Быстро, детерминировано, параллельно.

В Kappa Architecture «пересчёт» означает запустить Streaming Job с startingOffsets="earliest". Принципиальные отличия:

Стоимость Replay в Kappa: каждое историческое событие нужно прочитать из Kafka, десериализовать, обработать через stateful операции и записать в Iceberg. При объёме 365 дней × 10 000 событий/сек = 315 миллиардов событий - это серьёзная нагрузка на кластер, занимающая часы или дни.

Параллелизм ограничен числом партиций Kafka: Spark может читать максимум столько executor'ов параллельно, сколько партиций в топике. Если топик 12 партиций - максимум 12 параллельных потоков чтения, независимо от размера кластера.

State Store при Replay: если пайплайн stateful (window aggregations, sessionization), Spark должен построить весь промежуточный state по пути от offset 0 до текущего момента. Для сессионизации с 10 млн активных пользователей это сотни гигабайт State Store.

5.2 Паттерн «Stateless First, Stateful Second»

Для оптимизации долгого Replay рекомендуется разбить пайплайн на два этапа:

# === Этап 1: Stateless Replay (быстрый, не требует State Store) ===
# Читает Kafka → применяет чистые трансформации → пишет в промежуточный Bronze Iceberg
# Скорость: ограничена только Kafka throughput и CPU на десериализацию
# Параллелизм: partitions * executors

def stateless_transform(df):
    """Чистые, безstate трансформации. Детерминированы, быстры."""
    return (
        df
        .filter(F.col("amount") > 0)
        .withColumn("net_amount",
            F.when(F.col("discount_pct") <= 50,
                   F.col("amount") * (1 - F.col("discount_pct") / 100))
            .otherwise(F.lit(0.0))
        )
        .withColumn("event_date", F.to_date("event_ts"))
        .withColumn("event_hour", F.date_trunc("hour", F.col("event_ts")))
    )

bronze_replay_query = (
    raw_kafka_stream                       # startingOffsets="earliest"
    .transform(stateless_transform)        # Stateless: быстро
    .writeStream
    .trigger(processingTime="10 seconds")  # Частый триггер для быстрого replay
    .outputMode("append")
    .format("iceberg")
    .option("checkpointLocation", "s3://bucket/checkpoints/replay-bronze-v2")
    .toTable("lakehouse.bronze.purchases_v2")
)

# Запускаем Этап 1, ждём завершения (когда job догонит current offsets)
bronze_replay_query.awaitTermination(timeout=86400)  # timeout 24 часа

# === Этап 2: Stateful Aggregation (медленнее, но работает с уже загруженными данными) ===
# Читает из Bronze Iceberg (batch!), строит stateful агрегаты
# Преимущество: Spark Batch параллелизм не ограничен числом Kafka-партиций
silver_from_bronze = (
    spark.read                           # BATCH read, не streaming!
         .format("iceberg")
         .load("lakehouse.bronze.purchases_v2")
         .transform(stateless_transform)  # Применяем ту же логику
)

# Stateful Window Aggregation - теперь это просто Spark Batch
hourly_revenue = (
    silver_from_bronze
    .groupBy(
        "event_hour",
        "event_type"
    )
    .agg(
        F.sum("net_amount").alias("net_revenue"),
        F.count("*").alias("tx_count"),
        F.countDistinct("user_id").alias("unique_buyers")
    )
)

hourly_revenue.writeTo("lakehouse.gold.revenue_by_hour_v2").overwritePartitions()

«Stateless First, Stateful Second» - это практический компромисс для Kappa при больших объёмах исторических данных. Stateless Replay работает на максимальной скорости Kafka без ограничений State Store. После завершения Replay - Stateful агрегации выполняются как обычный Batch Job из уже записанного Bronze Iceberg, где нет ограничения по числу партиций Kafka.

5.3 Дедупликация при Replay: защита от двойного счёта

При Replay существует опасность двойного счёта: если новый V2-пайплайн запущен с earliest, а часть событий уже была записана V1 в Bronze-таблицу - без дедупликации данные окажутся в таблице дважды.

def kappa_foreachbatch_with_dedup(batch_df, batch_id: int) -> None:
    """
    Идемпотентная запись с дедупликацией при Replay.
    Использует MERGE INTO для обнаружения уже существующих event_id.
    """
    if batch_df.isEmpty():
        return

    # Дедупликация внутри батча (могут быть дубли из Kafka при at-least-once)
    deduped = batch_df.dropDuplicates(["event_id"])
    deduped.createOrReplaceTempView("incoming")

    # MERGE INTO Bronze: если event_id уже есть - пропускаем (не обновляем!)
    # При Replay это защищает от двойного счёта
    # WHEN MATCHED: ничего не делаем (событие уже записано корректно)
    # WHEN NOT MATCHED: вставляем новое событие
    spark.sql("""
        MERGE INTO lakehouse.bronze.purchases_v2 AS target
        USING incoming AS source
        ON target.event_id = source.event_id
        WHEN NOT MATCHED THEN INSERT *
    """)
    # Нет WHEN MATCHED: существующие записи неизменны (иммутабельность Bronze)

    # Альтернатива для чистого Bronze (только APPEND):
    # Держите отдельную таблицу processed_event_ids и фильтруйте перед записью
    # Это дороже, но не требует MERGE INTO для Bronze

6. Lakehouse + Kappa: стриминговая запись в Iceberg

6.1 Как Iceberg обеспечивает надёжность стриминговой записи

Apache Iceberg разработан с учётом streaming workloads. Каждый микро-батч Structured Streaming создаёт новый Iceberg Snapshot - атомарный коммит в metadata. Это обеспечивает несколько критических свойств:

Exactly-Once: Spark записывает в Iceberg вместе с идентификатором транзакции (query ID + batch ID). При повторном коммите (после сбоя) Iceberg обнаруживает уже существующий batch ID и молча пропускает повторную запись. Это нативная exactly-once семантика без дополнительной дедупликации в foreachBatch.

Concurrent Readers: пока Streaming Job пишет новый снапшот, аналитики читают предыдущий. Чтение никогда не видит «половинного» состояния - только завершённые снапшоты.

Time Travel: каждый снапшот сохраняется как отдельная версия. Аналитик может запросить данные «на 10 минут назад» - до последнего обновления Streaming Job.

# === Стандартная стриминговая запись в Iceberg ===
# Spark + Iceberg нативно поддерживают exactly-once через WAL Iceberg
purchase_stream_query = (
    raw_kafka_stream
    .transform(stateless_transform)
    .writeStream
    .trigger(processingTime="30 seconds")
    .outputMode("append")
    .format("iceberg")
    # checkpoint хранит: последние обработанные Kafka offsets + Iceberg commit ID
    # Вместе они гарантируют exactly-once при перезапуске
    .option("checkpointLocation", "s3://bucket/checkpoints/iceberg-stream-v1")
    # fanout_enabled: разрешает запись в несколько партиций за один task
    # Нужен при стриминге, когда один executor пишет в несколько partition key
    .option("fanout-enabled", "true")
    .toTable("lakehouse.bronze.purchases")
)

fanout-enabled = true критически важен для стриминга. В batch-режиме Spark гарантирует, что каждый task пишет только в одну partition. В streaming-режиме один task может получить события из разных партиций Kafka, которые относятся к разным event_date. Без fanout-enabled Iceberg выбросит исключение при попытке записать в несколько partition key из одного task.

6.2 Stream-to-Stream Join: объединение двух потоков

В e-commerce типичная задача - соединить поток показов рекламы (impressions) с потоком кликов (clicks) для расчёта CTR в реальном времени. Это Stream-to-Stream Join - оба источника потоковые.

# Схема потока показов рекламы
impression_schema = StructType([
    StructField("impression_id", StringType(), False),
    StructField("ad_id",         StringType(), False),
    StructField("user_id",       StringType(), False),
    StructField("campaign_id",   StringType(), True),
    StructField("shown_at",      TimestampType(), False),
])

# Схема потока кликов
click_schema = StructType([
    StructField("click_id",  StringType(),    False),
    StructField("ad_id",     StringType(),    False),
    StructField("user_id",   StringType(),    False),
    StructField("clicked_at", TimestampType(), False),
])

# Читаем поток показов
impressions_stream = (
    spark.readStream
         .format("kafka")
         .option("kafka.bootstrap.servers", "kafka:9092")
         .option("subscribe", "ad_impressions")
         .option("startingOffsets", "latest")
         .load()
         .select(F.from_json(F.col("value").cast("string"), impression_schema).alias("d"))
         .select("d.*")
         # Watermark: ждём опоздавшие показы до 2 часов
         .withWatermark("shown_at", "2 hours")
)

# Читаем поток кликов
clicks_stream = (
    spark.readStream
         .format("kafka")
         .option("kafka.bootstrap.servers", "kafka:9092")
         .option("subscribe", "ad_clicks")
         .option("startingOffsets", "latest")
         .load()
         .select(F.from_json(F.col("value").cast("string"), click_schema).alias("d"))
         .select("d.*")
         # Watermark: ждём опоздавшие клики до 2 часов
         .withWatermark("clicked_at", "2 hours")
)

# Stream-to-Stream Join
# Inner Join: выдаём только пары показ+клик (не показы без клика)
# State: Spark держит в памяти показы и клики до истечения watermark
ad_performance = (
    impressions_stream.alias("imp")
    .join(
        clicks_stream.alias("clk"),
        on=[
            # Условие JOIN: тот же ad_id, тот же пользователь
            F.col("imp.ad_id")   == F.col("clk.ad_id"),
            F.col("imp.user_id") == F.col("clk.user_id"),
            # Временное ограничение: клик должен быть после показа и не позже чем через 2 часа
            # Это критически важно: без временного ограничения State Store бесконечно растёт
            F.col("clk.clicked_at") >= F.col("imp.shown_at"),
            F.col("clk.clicked_at") <= F.col("imp.shown_at") + F.expr("INTERVAL 2 HOURS")
        ],
        how="inner"
    )
    .select(
        F.col("imp.impression_id"),
        F.col("imp.ad_id"),
        F.col("imp.user_id"),
        F.col("imp.campaign_id"),
        F.col("imp.shown_at"),
        F.col("clk.click_id"),
        F.col("clk.clicked_at"),
        # Время от показа до клика
        (F.col("clk.clicked_at").cast("long") - F.col("imp.shown_at").cast("long"))
            .alias("click_latency_sec")
    )
)

# Записываем в Iceberg: каждый клик с информацией о показе
(
    ad_performance
    .writeStream
    .trigger(processingTime="1 minute")
    .outputMode("append")
    .format("iceberg")
    .option("checkpointLocation", "s3://bucket/checkpoints/ad-performance")
    .toTable("lakehouse.gold.ad_click_performance")
)

Временное ограничение в JOIN - обязательно. Без условия clicked_at <= shown_at + INTERVAL 2 HOURS State Store будет хранить каждый показ, пока не придёт соответствующий клик - бесконечно. С временным ограничением Spark может очистить State Store после истечения max(watermark_impressions, watermark_clicks).


7. Операционные вызовы Kappa Architecture

7.1 Мониторинг State Store: как понять, что скоро OOM

State Store - главный операционный риск Kappa. В отличие от batch-обработки, где объём обрабатываемых данных конечен, streaming State Store может расти бесконечно при неправильно настроенном watermark.

7.2 Программный мониторинг streaming metrics

import time

def monitor_streaming_job(query, check_interval_sec: int = 60):
    """
    Мониторинг здоровья Streaming Job в production.
    Логирует ключевые метрики и поднимает алерты при аномалиях.

    query: активный StreamingQuery (результат .start())
    check_interval_sec: интервал проверки метрик
    """
    while query.isActive:
        progress = query.lastProgress

        if progress is None:
            time.sleep(check_interval_sec)
            continue

        # === Базовые метрики производительности ===
        input_rate   = progress.get("inputRowsPerSecond", 0)
        process_rate = progress.get("processedRowsPerSecond", 0)
        batch_dur_ms = progress.get("batchDuration", 0)
        num_rows     = progress.get("numInputRows", 0)

        print(f"\n=== Streaming Health Check ===")
        print(f"Input Rate:     {input_rate:.0f} строк/сек")
        print(f"Process Rate:   {process_rate:.0f} строк/сек")
        print(f"Batch Duration: {batch_dur_ms}ms")
        print(f"Rows in batch:  {num_rows}")

        # === Алерт: Consumer Lag растёт ===
        if input_rate > 0 and process_rate < input_rate * 0.9:
            lag_ratio = input_rate / process_rate
            print(f"WARN: Process Rate отстаёт от Input Rate в {lag_ratio:.1f}x!")
            print("      Возможная причина: медленный foreachBatch, медленный S3, GC паузы")

        # === Алерт: Батч слишком долгий ===
        # Если батч длится дольше trigger interval - следующий батч задерживается
        # Это первый признак растущего Consumer Lag
        trigger_interval_ms = 30_000  # наш trigger: 30 секунд
        if batch_dur_ms > trigger_interval_ms * 1.5:
            print(f"WARN: Batch Duration {batch_dur_ms}ms > trigger * 1.5 ({trigger_interval_ms * 1.5}ms)")
            print("      Оптимизируйте foreachBatch или увеличьте trigger interval")

        # === Метрики State Store ===
        state_operators = progress.get("stateOperators", [])
        for i, so in enumerate(state_operators):
            rows_total   = so.get("numRowsTotal", 0)
            memory_bytes = so.get("memoryUsedBytes", 0)
            mem_mb = memory_bytes / 1024 / 1024

            print(f"\nState Operator [{i}]:")
            print(f"  Строк в state:  {rows_total:,}")
            print(f"  Память:         {mem_mb:.1f} MB")

            # Алерт: State Store занимает много памяти
            executor_heap_mb = 4096   # executor.memory = 4GB
            if mem_mb > executor_heap_mb * 0.7:
                print(f"  КРИТИЧНО: State Store > 70% heap executor!")
                print(f"  Действие: уменьшить watermark_delay или добавить executor'ы")

        # === Watermark Gap ===
        event_time_watermark = progress.get("eventTime", {}).get("watermark", "N/A")
        print(f"\nWatermark: {event_time_watermark}")

        time.sleep(check_interval_sec)


# Запускаем мониторинг в отдельном потоке
import threading
monitor_thread = threading.Thread(
    target=monitor_streaming_job,
    args=(purchase_stream_query, 60),
    daemon=True
)
monitor_thread.start()

7.3 Schema Evolution в Kappa: осторожность с checkpoint

Schema Evolution - изменение схемы данных в потоке - особенно опасна в Kappa Architecture, потому что checkpoint хранит информацию о схеме State Store. Несовместимое изменение ломает checkpoint.

Практическое правило: любое изменение State Schema (структуры данных, хранимых в State Store) требует нового checkpoint. Изменения схемы входящих данных (Kafka JSON) или выходных таблиц (Iceberg) - в большинстве случаев совместимы и не требуют пересоздания checkpoint.


8. Анти-паттерны Kappa Architecture

8.1 «Стриминг везде» - самый дорогой анти-паттерн

Kappa Architecture соблазняет инженеров применять streaming-first подход везде. Но это ошибка.

Streaming дороже batch по нескольким причинам: постоянно работающий кластер (vs кластер по расписанию), overhead на checkpoint каждые N секунд, State Store занимает RAM постоянно, debugging stateful streaming в разы сложнее.

Если ежемесячный финансовый отчёт можно запустить батчем один раз - не нужен Kappa.

8.2 Короткий retention Kafka: уничтожение Replay

# АНТИПАТТЕРН: Kafka с коротким retention для Kappa Source of Truth
# 7 дней retention = нельзя сделать Replay за месяц
# Бизнес решает пересчитать данные за прошлый квартал - невозможно

# СЛЕДСТВИЕ: приходится держать отдельный S3 архив событий
# Это фактически возвращает к Lambda Architecture (S3 = Master Dataset)

# ПРАВИЛЬНО: retention должен покрывать максимальный горизонт Replay
# Если пересчёты нужны за год - retention минимум 365 дней
# Стоимость хранения 1TB в Kafka ≈ $30-100/месяц (зависит от репликации)
# vs $23/TB/месяц в S3 Standard
# Для больших объёмов: Kafka → S3 через MirrorMaker2 или Kafka S3 Sink
# Тогда Replay читает S3, а не Kafka - снимает ограничение по партициям

8.3 Нет идемпотентности в Sink

# АНТИПАТТЕРН: неидемпотентный sink при Replay
# При Replay пайплайн читает те же события снова
# Если sink не идемпотентен - данные дублируются

def bad_foreachbatch(batch_df, batch_id):
    # ОШИБКА: простой INSERT без дедупликации
    # При Replay те же event_id записываются второй раз
    batch_df.write.format("iceberg").mode("append").saveAsTable("bronze.events")
    # Результат: 2x дубликаты в таблице

def good_foreachbatch(batch_df, batch_id):
    # ПРАВИЛЬНО: MERGE INTO с дедупликацией по event_id
    # При Replay: WHEN MATCHED - пропускаем, WHEN NOT MATCHED - вставляем
    batch_df.createOrReplaceTempView("incoming")
    spark.sql("""
        MERGE INTO lakehouse.bronze.events AS t
        USING (SELECT DISTINCT * FROM incoming) AS s
        ON t.event_id = s.event_id
        WHEN NOT MATCHED THEN INSERT *
    """)
    # Результат: идемпотентно, Replay не создаёт дубликатов

8.4 State без ограничений - гарантированный OOM

# АНТИПАТТЕРН: stateful операция без watermark
# Spark хранит state для КАЖДОГО ключа и КАЖДОГО окна навсегда

bad_agg = (
    stream_df
    # НЕТ withWatermark → нет очистки State Store
    .groupBy(F.window("event_ts", "1 hour"), "user_id")
    .agg(F.sum("amount"))
    # После 7 дней: State Store хранит 7 * 24 = 168 часовых окон × N пользователей
    # После месяца: 720 окон × N пользователей → OOM
)

# ПРАВИЛЬНО: watermark освобождает State Store по мере закрытия окон
good_agg = (
    stream_df
    .withWatermark("event_ts", "30 minutes")    # Ждём опоздавших 30 минут
    .groupBy(F.window("event_ts", "1 hour"), "user_id")
    .agg(F.sum("amount"))
    # После watermark > window.end: Spark удаляет окно из State Store
    # Стабильный объём state ≈ (watermark_delay / trigger_interval) × active_users_in_window
)

8.5 Разные checkpoint для одного логического джоба

# АНТИПАТТЕРН: несколько стриминговых джобов читают один топик
# с разными checkpoint - нет гарантии согласованности

query_revenue = (stream_df.writeStream
    .option("checkpointLocation", "s3://bucket/checkpoints/revenue")
    .toTable("gold.revenue"))

query_users = (stream_df.writeStream
    .option("checkpointLocation", "s3://bucket/checkpoints/users")
    .toTable("gold.users"))

# Проблема: при сбое revenue job откатится к своему checkpoint (offset 1000),
# а users job продолжит с offset 1100.
# Таблицы revenue и users теперь рассогласованы по времени.

# ПРАВИЛЬНО: один foreachBatch пишет в несколько таблиц атомарно
def unified_foreachbatch(batch_df, batch_id):
    # Один checkpoint, один микро-батч, несколько целевых таблиц
    batch_df.createOrReplaceTempView("batch")
    spark.sql("INSERT INTO gold.revenue SELECT ... FROM batch")
    spark.sql("INSERT INTO gold.users SELECT ... FROM batch")
    # Если любая запись падает - весь batch_id откатывается

unified_query = (
    stream_df.writeStream
             .foreachBatch(unified_foreachbatch)
             .option("checkpointLocation", "s3://bucket/checkpoints/unified")
             .start()
)

9. Производственный кейс: Kappa Pipeline для e-commerce

9.1 Требования и архитектурный выбор

Рассмотрим e-commerce платформу, перешедшую с Lambda на Kappa:

  • 500 000 пользователей, 2 000 событий в секунду в пике
  • Требования: live dashboard (задержка < 30 сек), сессионизация, CTR рекламы
  • Проблема с Lambda: синхронизация двух кодовых баз при каждом изменении бизнес-правил

Решение: Kappa на Spark Structured Streaming + Apache Iceberg + Kafka с retention 90 дней.

9.2 Полный код основного Kappa Pipeline

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

spark = SparkSession.builder \
    .appName("Kappa-ECommerce-Pipeline-V2") \
    .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://catalog:8181") \
    .config("spark.sql.catalog.lakehouse.warehouse", "s3://datalake/warehouse") \
    .config("spark.sql.streaming.stateStore.providerClass",
            "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider") \
    .config("spark.sql.streaming.stateStore.rocksdb.compactOnCommit", "true") \
    # Оптимизация: не записывать лишние метаданные при каждом микро-батче
    .config("spark.sql.streaming.minBatchesToRetain", "2") \
    .getOrCreate()

# === Схема события ===
event_schema = StructType([
    StructField("event_id",    StringType(),    False),
    StructField("user_id",     StringType(),    False),
    StructField("session_id",  StringType(),    True),
    StructField("event_type",  StringType(),    False),
    StructField("product_id",  StringType(),    True),
    StructField("category",    StringType(),    True),
    StructField("amount",      DoubleType(),    True),
    StructField("discount_pct", DoubleType(),   True),
    StructField("event_ts",    TimestampType(), False),
    StructField("device_type", StringType(),    True),
])

# === Чтение из Kafka ===
# PIPELINE_VERSION управляет: latest для нового деплоя, earliest для replay
PIPELINE_VERSION = "V2"
IS_REPLAY = True   # True = читаем историю, False = только real-time

raw_stream = (
    spark.readStream
         .format("kafka")
         .option("kafka.bootstrap.servers", "kafka:9092")
         .option("subscribe", "ecommerce.events")
         .option("startingOffsets", "earliest" if IS_REPLAY else "latest")
         # При Replay: большой батч для скорости. При real-time: защита от OOM
         .option("maxOffsetsPerTrigger", "500000" if IS_REPLAY else "50000")
         .option("failOnDataLoss", "false")
         .load()
         .select(
             F.col("partition").alias("kafka_partition"),
             F.col("offset").alias("kafka_offset"),
             F.col("timestamp").alias("kafka_ts"),
             F.from_json(F.col("value").cast("string"), event_schema).alias("d")
         )
         .select("kafka_partition", "kafka_offset", "kafka_ts", "d.*")
)


# === Единая функция трансформации (Kappa принцип) ===
def apply_business_rules_v2(df):
    """
    V2 бизнес-логика. Отличие от V1: добавлен category_revenue_share.
    Используется ОДИНАКОВО при Replay и при real-time processing.
    Детерминирована: результат зависит только от event_ts (не от current_timestamp).
    """
    return (
        df
        # Фильтрация невалидных записей
        .filter(F.col("amount").isNull() | (F.col("amount") >= 0))
        # V2 правило: скидки > 50% не учитываются
        .withColumn(
            "net_amount",
            F.when(
                (F.col("event_type") == "purchase") &
                (F.col("discount_pct") <= 50),
                F.col("amount") * (1 - F.col("discount_pct") / 100)
            ).when(
                F.col("event_type") == "purchase",
                F.lit(0.0)   # Скидка > 50%: не считаем
            ).otherwise(F.lit(0.0))
        )
        # НОВОЕ в V2: категория вклада в выручку
        .withColumn(
            "revenue_tier",
            F.when(F.col("net_amount") == 0,     "no_revenue")
             .when(F.col("net_amount") < 1000,   "micro")
             .when(F.col("net_amount") < 10000,  "mid")
             .otherwise("large")
        )
        # Производные временны́е поля
        .withColumn("event_date", F.to_date("event_ts"))
        .withColumn("event_hour", F.date_trunc("hour", F.col("event_ts")))
        .withColumn("event_minute", F.date_trunc("minute", F.col("event_ts")))
    )


# === foreachBatch: единый обработчик ===
def kappa_main_batch(batch_df, batch_id: int) -> None:
    """
    Центральный foreachBatch для Kappa Pipeline.
    Обрабатывает каждый микро-батч: Bronze → Silver → Gold.
    Идемпотентен: повторный вызов с теми же данными даёт тот же результат.
    """
    if batch_df.isEmpty():
        return

    # Дедупликация: защита от Kafka at-least-once и повторных foreachBatch вызовов
    deduped = batch_df.dropDuplicates(["event_id"]).cache()
    record_count = deduped.count()

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

    print(f"[Pipeline-{PIPELINE_VERSION} | Batch {batch_id}] {record_count} событий")

    # Применяем V2 бизнес-правила
    enriched = apply_business_rules_v2(deduped)
    enriched.cache()

    # --- Bronze: идемпотентная запись (MERGE IN, игнорируем дубли по event_id) ---
    enriched.createOrReplaceTempView("incoming_events")

    spark.sql(f"""
        MERGE INTO lakehouse.bronze.events AS t
        USING (SELECT * FROM incoming_events) AS s
        ON t.event_id = s.event_id
        WHEN NOT MATCHED THEN INSERT *
    """)

    # --- Silver: текущие агрегаты по пользователям (MERGE INTO, обновляем) ---
    purchases = enriched.filter(F.col("event_type") == "purchase")

    if not purchases.isEmpty():
        purchases.createOrReplaceTempView("incoming_purchases")

        spark.sql("""
            MERGE INTO lakehouse.silver.user_purchase_summary AS t
            USING (
                SELECT
                    user_id,
                    SUM(net_amount)   AS batch_revenue,
                    COUNT(*)          AS batch_count,
                    MAX(event_ts)     AS last_purchase_ts,
                    MAX(revenue_tier) AS max_tier
                FROM incoming_purchases
                GROUP BY user_id
            ) AS s
            ON t.user_id = s.user_id
            WHEN MATCHED THEN UPDATE SET
                t.total_revenue   = t.total_revenue + s.batch_revenue,
                t.total_purchases = t.total_purchases + s.batch_count,
                t.last_purchase_ts = s.last_purchase_ts,
                t.updated_at      = current_timestamp()
            WHEN NOT MATCHED THEN INSERT
                (user_id, total_revenue, total_purchases,
                 last_purchase_ts, created_at, updated_at)
            VALUES (s.user_id, s.batch_revenue, s.batch_count,
                    s.last_purchase_ts, current_timestamp(), current_timestamp())
        """)

    # --- Gold: 1-минутные агрегаты для live дашборда ---
    gold_agg = (
        enriched
        .groupBy("event_minute", "event_type", "revenue_tier", "device_type")
        .agg(
            F.sum("net_amount").alias("net_revenue"),
            F.count("*").alias("event_count"),
            F.approx_count_distinct("user_id", rsd=0.02).alias("unique_users"),
            F.approx_count_distinct("session_id", rsd=0.02).alias("unique_sessions")
        )
    )

    # Append: Gold-таблица читается через Serving View с ограничением по времени
    # Компакция удалит старые мелкие файлы
    (
        gold_agg
        .writeTo("lakehouse.gold.realtime_revenue")
        .append()
    )

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

    # Очищаем кеш
    enriched.unpersist()
    deduped.unpersist()


# === Запуск основного Kappa Pipeline ===
# При IS_REPLAY=True: читает историю 90 дней, затем переходит в real-time
# При IS_REPLAY=False: читает только новые события

trigger_config = {"processingTime": "10 seconds"} if IS_REPLAY else {"processingTime": "30 seconds"}

main_query = (
    raw_stream
    .writeStream
    .trigger(**trigger_config)
    .foreachBatch(kappa_main_batch)
    .option("checkpointLocation",
            f"s3://bucket/checkpoints/kappa-main-{PIPELINE_VERSION.lower()}")
    .start()
)

print(f"Kappa Pipeline {PIPELINE_VERSION} запущен")
print(f"Режим: {'Replay (история + real-time)' if IS_REPLAY else 'Real-time'}")
print(f"Checkpoint: s3://bucket/checkpoints/kappa-main-{PIPELINE_VERSION.lower()}")

9.3 Процедура переключения V1 → V2 без downtime

# === Алгоритм migration: V1 → V2 без остановки сервиса ===

# Шаг 1: V2 запускается с IS_REPLAY=True, пишет в gold.realtime_revenue_v2
# Параллельно V1 продолжает работать: gold.realtime_revenue_v1

# Шаг 2: Ждём, пока V2 «догонит» V2 до real-time
# Мониторим отставание через query.lastProgress
def wait_for_catchup(query_v2, max_lag_sec: int = 60):
    """Ждём, пока V2 догонит real-time поток."""
    import time
    while True:
        progress = query_v2.lastProgress
        if progress:
            # watermarkGap: разница между последним event_ts и watermark
            # Когда V2 в real-time: watermarkGap ≈ watermark_delay
            batch_dur = progress.get("batchDuration", 999999)
            trigger_ms = 30_000   # наш trigger

            # Если batch duration < trigger: V2 обрабатывает меньше данных чем приходит
            # Это означает: V2 догнал или почти догнал real-time
            if batch_dur < trigger_ms:
                print("V2 догнал real-time поток. Готов к переключению.")
                break

        print(f"V2 ещё догоняет историю... Ждём {max_lag_sec}s")
        time.sleep(max_lag_sec)


wait_for_catchup(main_query)

# Шаг 3: Атомарное переключение Serving View
# Аналитики читают через view gold.realtime_serving
# Переключение мгновенное, нет downtime
spark.sql("""
    CREATE OR REPLACE VIEW lakehouse.gold.realtime_serving AS
    SELECT * FROM lakehouse.gold.realtime_revenue_v2
    WHERE event_minute >= current_timestamp() - INTERVAL 48 HOURS
""")

print("Serving View переключён на V2")
print("Аналитики теперь видят данные из V2 пайплайна")

# Шаг 4: Останавливаем V1 (после переключения View)
# query_v1.stop()
# V1 checkpoint можно удалить через 7 дней (если Replay V2 подтверждён)

10. Итоги: когда выбирать Kappa

10.1 Kappa как философия и инструмент

Kappa Architecture - это не серебряная пуля. Это архитектурный выбор с конкретными trade-offs, подходящий для определённых условий:

10.2 Пять принципов production-grade Kappa

Принцип 1: Kafka - это ваш Master Dataset. Настройте retention на весь горизонт возможного Replay. Короткий retention уничтожает главное преимущество Kappa - возможность пересчёта.

Принцип 2: Весь код детерминирован по Event Time. Никаких current_timestamp(), now(), date.today() в логике трансформации. Только event_ts из события. Иначе Replay даёт другие результаты.

Принцип 3: Watermark - всегда, для каждой stateful операции. Без watermark State Store - бомба замедленного действия. Выбирайте watermark_delay = max_expected_late_arrival × 1.5.

Принцип 4: Один checkpoint на один логический пайплайн. Несколько writeStream из одного источника - это несколько независимых потоков с независимыми checkpoint. Используйте foreachBatch для атомарной записи в несколько таблиц.

Принцип 5: Migration через Blue-Green Streaming. При смене логики - V2 запускается параллельно с новым checkpoint. Serving View переключается атомарно после догона. V1 остановить после подтверждения. Нет downtime, нет потери данных.

Следующий урок завершит трилогию: сравним Lambda и Kappa по конкретным критериям - стоимостной модели, операционной сложности и эволюции в Lakehouse - и поможет принять обоснованный архитектурный выбор для вашей платформы.