Lambda vs Kappa Architecture concept page. Lambda (Marz 2011): three layers — batch (Spark, точно/медленно) + speed (Storm/Flink, быстро/приближённо) + serving (merge views). Kappa (Kreps 2014): один streaming pipeline через Kafka log + Flink EOS, replay через blue-green datasource swap (new consumer group offset=earliest, отдельный output namespace, atomic swap queries). Lakehouse hybrid (2026): Iceberg/Delta как unified primitive, Spark Structured Streaming + batch backfill пишут в одну table, Materialize держит incremental view. Сценарии: Lambda dual pipeline, Kappa happy path, Kappa replay, Lakehouse hybrid, Lambda merge bug failure mode. ADR-001: когда Kappa default, когда Lambda оправдан. ADR-002: blue-green datasource swap для replay.
В 2010 Nathan Marz описал Lambda Architecture — единственный тогда способ совместить streaming-скорость и batch-точность в одной системе. К 2014 Jay Kreps (создатель Kafka) ответил эссе «Questioning the Lambda Architecture» и предложил Kappa — один streaming pipeline, replay через перечитывание Kafka log. Десятилетие индустрия медленно мигрировала с Lambda на Kappa, а к 2026 lakehouse-гибриды (Iceberg/Delta + Flink/Spark Structured Streaming) стали де-факто default.
Понимать разницу важно потому что:
Lambda — параллельно batch (точно, медленно) + speed (быстро, приближённо), serving мержит обе view. Kappa — только streaming, replay через consumer group reset на Kafka log от offset=earliest. Lakehouse — unified table (Iceberg/Delta), batch и stream пишут в одно и то же.
Lambda — два кодовых движка с разной семантикой времени. Kappa — один движок, который умеет делать и то и другое за счёт log replay. Lakehouse — один storage primitive, любой compute поверх.
Три горизонтальных слоя из одного evt источника:
master dataset (S3) + параллельная speed ветка; batch layer (Spark 6h) пишет в batch view, speed (Flink) в realtime view, serving мержит обе. Видна та самая «двойная» топология — каждый event идёт сразу в две системы.kafka log → flink-v1 → druid-v1 → serving. Один путь. Рядом теневой flink-v2 + druid-v2 для replay-сценария.lakehouse) — iceberg как unified primitive; параллельно spark-stream (CDC) и spark-batch (backfill), оба пишут в materialize incremental view.ADR на master: ADR-001 (когда выбирать Lambda vs Kappa) и ADR-002 (blue-green datasource swap для replay).
Каждый event форкается в master dataset и speed layer одновременно. Spark cron каждые 6 часов пересчитывает всю историю в batch view (ground truth). Flink держит approximate realtime view для последних 6h. Query до 6h ago — точный ответ из batch view, окно 0-6h — merge обеих. Боль: та же логика реализована дважды в разных engines.
Один pipeline: events → Kafka (30d retention + tiered S3 для дешёвого long-term) → Flink с exactly-once → Druid datasource с idempotent upsert → serving. Один codebase, один compute, один store. Real-time данные сразу точные.
Day 5: найден баг в transformation. Деплоится Flink v2 с фиксом в shadow cluster, new consumer group от offset=earliest, output в отдельный Druid datasource v2. Production v1 продолжает untouched — zero downtime. Когда v2 догнал production offset + watermark OK — атомарный blue-green swap queries на v2, потом drop v1 и retire старый consumer group. Cost: Nx compute на window replay.
2026 default. Events пишутся в Iceberg table (S3 + Parquet + manifest). Spark Structured Streaming читает её snapshots инкрементально через CDC. Для backfill старше Kafka retention запускается Spark batch на тот же S3 — пишет в ту же materialize view. SQL queries поверх. Один codebase (Spark SQL), один store, best of both worlds.
Real failure mode Lambda: batch run завершился с timestamp T, но realtime view не expired events до T → event_X засчитан и в batch view, и в realtime view → double counting на merge. Лечение: explicit expiration realtime entries по batch run timestamp.
ADR-001: Lambda vs Kappa — когда какая выигрывает. Default = Kappa (или lakehouse-вариант). Один codebase, один compute, replay через consumer reset, output stores с idempotent upsert (Druid dedup_column, ClickHouse ReplacingMergeTree, Cassandra INSERT IF NOT EXISTS). Lambda оправдан только когда: (1) регуляторы требуют независимый batch reconciliation (банкинг, payments, медицина); (2) в стеке нет EOS-capable streaming engine; (3) replay full history стабильно дороже maintenance двух pipelines. Hybrid (Kappa + occasional batch backfill для данных старше Kafka retention) — частый и здоровый компромисс.
ADR-002: Replay strategy = blue-green datasource swap. Naive replay в тот же sink даёт дубли или corrupt state; «остановить prod и replay» даёт часы downtime. Правильно: shadow Flink с new consumer group от offset=earliest, отдельный output namespace, ждём полный catch-up + watermark check, atomic swap queries, потом cleanup. Обязательно: scheduled replays в off-peak, capacity planning заранее (replay storm может потребовать 5-10× compute), idempotent upsert на sink, периодические snapshots output чтобы не реплеить всю историю каждый раз.
dedup_column, ClickHouse ReplacingMergeTree, Cassandra INSERT IF NOT EXISTS, Iceberg MERGE INTO.Не Kappa, когда:
Не Lambda, когда:
Не Lakehouse hybrid, когда: