Kappa-архитектура на Spark Structured Streaming: стриминг как единый источник правды
Kappa-архитектура на Spark Structured Streaming: стриминг как единый источник правды
В предыдущем уроке мы разобрали 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:
- V1 работает в production, обрабатывает real-time события, пишет в
gold_revenue_v1. Аналитики читают эту таблицу. - V2 стартует параллельно с
startingOffsets="earliest"и новым checkpoint-путём. Kafka хранит полную историю (365 дней), V2 читает её на максимальной скорости кластера - это и есть «батч» в Kappa. - V2 догоняет real-time: скорость чтения Kafka намного выше скорости поступления новых событий, поэтому V2 быстро сокращает отставание и начинает читать те же текущие offsets, что и V1.
- Атомарное переключение:
ALTER VIEWв Iceberg меняет указатель с V1-таблицы на V2-таблицу за один атомарный коммит. Аналитики переключаются мгновенно, без 404. - 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 - и поможет принять обоснованный архитектурный выбор для вашей платформы.