Сравнение Lambda vs Kappa: критерии выбора, стоимостная модель и эволюция в Lakehouse

Сравнение Lambda vs Kappa: критерии выбора, стоимостная модель и эволюция в Lakehouse

datamodeling

Три урока этого блока провели нас через два архитектурных мира: Lambda - двухслойная система с чёткими ролями Batch и Speed, Kappa - радикальная ставка на единый стриминговый пайплайн. Теперь пора сделать главное: сравнить их по критериям, которые важны на реальных проектах - деньги, операционная сложность, размер команды, объём данных, требования SLA.

Этот урок - не «кто лучше». Это инструментарий для принятия обоснованного архитектурного решения, которое вы сможете защитить перед командой, стейкхолдерами и собственной совестью в три часа ночи, когда production упадёт.


1. Исторический контекст: откуда взялся конфликт парадигм

1.1 Эпоха «нет-стримингового» Data Engineering

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

Первые попытки ускорить аналитику велись через nearline системы: рядом с основным DWH поднимался in-memory кеш (Memcached, Redis) с «свежими» данными. Бизнес-логика дублировалась вручную: OLAP-запрос читал из DWH, а real-time дашборд - из кеша. Синхронизация не гарантировалась никем.

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

1.2 Spark Streaming и первые стриминговые фреймворки

Apache Storm (2011), затем Apache Spark Streaming (2013, DStream API) принесли событийную обработку в мейнстрим. Инженеры наконец могли обрабатывать данные по мере поступления, без ожидания ночного ETL.

Но тут возникло новое противоречие: стриминг не заменял батч, он дополнял его. Kafka не хранила данные дольше нескольких дней. Исторических данных у стриминга не было. Для аудита, исправления ошибок, пересчёта - нужен был всё тот же батч на HDFS. Так родилась Lambda Architecture: ответ на вопрос «как жить, когда у тебя два разных хранилища и два разных движка?»

Kappa появилась позже (2014) как реакция на уже накопленную операционную боль Lambda-систем: команды, поддерживающие два пайплайна на разных языках, получали нескончаемые инциденты из-за расхождения бизнес-логики между слоями. Джей Кребс предложил: «Что если избавиться от батчевого слоя и сделать Kafka хранилищем навсегда?»

1.3 Lakehouse как третья волна

После 2018 года появились транзакционные форматы поверх object storage: Delta Lake (Databricks, 2019), Apache Iceberg (Netflix/Apple, 2020 GA), Apache Hudi (Uber). Они принесли ACID-транзакции, Schema Evolution и Time Travel в S3. Это изменило правила игры.

Три волны показывают эволюцию подходов. Каждая волна не отменяла предыдущую полностью - где-то в мире до сих пор работают MapReduce-джобы на Hadoop. Но каждая волна смещала «стандарт отрасли» и давала инженерам новые инструменты для решения старых компромиссов.

Lakehouse - это точка конвергенции: Iceberg-таблица одновременно является и «Master Dataset» для Batch Layer, и «Serving Layer» для Speed Layer, и выходным хранилищем Kappa Pipeline. Граница между парадигмами размывается.


2. Анатомия архитектур: ДНК Lambda и Kappa

2.1 Lambda: двухслойная защита от нестабильности

Lambda Architecture возникла в эпоху, когда стриминговые системы были ненадёжны. Storm падал. Данные теряли. Checkpointing не существовал. Batch Layer был страховкой: «если Speed Layer врёт - подождём ночного батча, и всё будет правильно».

Три свойства, за которые Lambda платит высокую цену:

Correctness - Batch Layer пересчитывает данные полностью от Master Dataset. Никаких приближений, никаких эффектов late data. Финальный ответ на вопрос «выручка за январь» будет точным.

Replayability - при обнаружении бага в логике Batch Job запускается заново от Master Dataset. Через несколько часов все витрины пересчитаны с правильной логикой. Master Dataset не тронут.

Fault Isolation - если Speed Layer упал, аналитики продолжают видеть исторически точные данные из Batch Views. Платформа деградирует, но не падает полностью.

За эти свойства Lambda расплачивается: дублированием кода (два пайплайна с одной бизнес-логикой), сложностью Serving Layer (граничная логика между Batch и Speed без двойного счёта) и операционными расходами двух независимых стеков.

2.2 Kappa: ставка на единый поток

Kappa отказывается от Batch Layer как концепции. Единственный «источник правды» - иммутабельный Kafka-лог. «Пересчёт» означает replay этого лога. «Батч» - это replay с startingOffsets=earliest.

Три свойства, за которые платит Kappa:

Unified Codebase - одна функция трансформации работает и при Replay, и в real-time. Изменение бизнес-правила применяется один раз.

Simplicity of Serving - нет граничной логики между Batch и Speed. Одна таблица в Iceberg, один SQL-запрос для аналитика.

Reduced Operational Surface - один Spark-кластер вместо двух, одна кодовая база вместо двух.

За это Kappa расплачивается: дорогим долгосрочным хранением в Kafka, сложностью Stateful Processing (watermarks, state explosion), дорогостоящим Replay при больших объёмах и высокими требованиями к команде (debugging streaming != debugging batch).


3. Технические критерии выбора: чеклист архитектора

3.1 Критерий 1: Replay Factor - сколько стоит пересчёт истории

Первый и важнейший технический вопрос при выборе архитектуры: как часто вам нужно пересчитывать исторические данные, и за какой период?

В Lambda пересчёт - это Spark Batch Job, читающий партиционированные файлы с S3 параллельно на весь кластер. При 100 executor'ах скорость чтения - сотни гигабайт в минуту. Год истории на 10TB может быть пересчитан за 2–4 часа.

В Kappa пересчёт - это Streaming Job с startingOffsets=earliest, ограниченный числом партиций Kafka (максимум одна задача чтения на партицию). Год истории на 10TB займёт значительно больше времени при тех же ресурсах - Kafka не оптимизирована для batch-read сотнями параллельных воркеров.

Практическое правило: если объём данных за типичный период пересчёта превышает 5TB, и Kafka-топик имеет менее 96 партиций - Lambda Batch Layer будет пересчитывать историю существенно быстрее, чем Kappa Replay. Это может означать разницу между «пересчёт за 2 часа» и «пересчёт за 20 часов».

3.2 Критерий 2: Сложность бизнес-логики и джойны

Не вся бизнес-логика одинаково хорошо реализуется в стриминге:

В стриминге хорошо работают:

  • Stateless трансформации (фильтрация, обогащение, маппинг)
  • Агрегации по коротким временным окнам (минуты, часы)
  • Дедупликация в пределах watermark-окна
  • Stream-to-Stream Join с ограниченным временным диапазоном
  • CDC-инжестия и простые MERGE INTO

В стриминге сложно или неэффективно:

  • JOIN с большими справочниками (dimension tables), которые меняются - broadcast join устаревает через микро-батч
  • Расчёт LTV за три года - требует State Store с трёхлетней историей или Replay
  • Сложные многоуровневые агрегации (Gold → Platinum → Diamond сегментация)
  • Retrospective анализ: «какой был статус пользователя 6 месяцев назад?»
from pyspark.sql import functions as F

# Пример: расчёт 90-дневного скользящего LTV в батче (просто)
# В Spark Batch: читаем всю историю, делаем оконную функцию - нет ограничений
def compute_90day_ltv_batch(purchases_df):
    """
    Расчёт LTV за 90 дней для каждого пользователя на каждый день.
    Оконная функция по event_date, 90-дневное скользящее окно.

    В Batch: эта функция отрабатывает на всей истории за минуты.
    В Streaming: хранить 90 дней state для миллионов пользователей = сотни GB RAM.
    """
    from pyspark.sql.window import Window

    window_90d = (
        Window
        .partitionBy("user_id")
        .orderBy(F.col("event_ts").cast("long"))
        # Окно: от 90 дней назад (90 * 86400 секунд) до текущего момента
        .rangeBetween(-90 * 86400, 0)
    )

    return (
        purchases_df
        .withColumn(
            "ltv_90d",
            F.sum("net_amount").over(window_90d)
        )
        .withColumn(
            "purchases_90d",
            F.count("*").over(window_90d)
        )
    )

# В Streaming (Kappa): то же самое требует State хранить 90 дней истории
# для каждого пользователя - это неприемлемо при > 1M активных пользователей.
# Правильный паттерн в Kappa:
# 1. Streaming пишет сырые покупки в Bronze Iceberg (append)
# 2. Batch Job (trigger=once, каждую ночь) читает Bronze и считает 90d LTV
# 3. Результат пишется в Gold таблицу
# Это гибридный подход: Kappa для инжестии, Batch для тяжёлой аналитики

3.3 Критерий 3: SLA и реальные требования к задержке

Один из самых распространённых архитектурных промахов - избыточная real-time: команда выбирает Kappa «потому что real-time лучше», хотя бизнес на самом деле не нуждается в задержке меньше 15 минут.

Важный вывод: для задержек от 15 минут до 1 часа Incremental Batch (паттерн trigger=AvailableNow по расписанию) часто оказывается дешевле и проще полноценного Kappa Streaming. Кластер не крутится 24/7, нет State Store, нет watermark - а данные обновляются каждые 15 минут.

3.4 Критерий 4: Операционная зрелость команды

Kappa Architecture предъявляет более высокие требования к команде, чем Lambda. Это не значит «Kappa сложнее» в абстрактном смысле - это значит, что конкретные инженерные навыки, необходимые для эксплуатации Kappa, встречаются реже.

Практическое наблюдение: команды, переходящие с Lambda на Kappa без достаточной подготовки, часто обнаруживают, что избавились от «проблемы двух кодовых баз» и приобрели «проблему State OOM в production в пятницу вечером». Операционная зрелость в streaming - это инвестиция, которую нельзя пропустить.


4. Стоимостная модель: деньги как критерий архитектуры

4.1 Compute: постоянный vs периодический кластер

Это наиболее очевидное экономическое различие между Kappa и Lambda:

Kappa: Spark Streaming Job работает 24/7. Кластер постоянно занят, постоянно потребляет ресурсы. Даже в период ночного low-traffic кластер крутится - ждёт новых событий, обновляет State Store, пишет checkpoint каждые N секунд.

Lambda: Batch Job запускается по расписанию. Ночью - один крупный джоб на несколько часов, днём - несколько инкрементальных, по расписанию. Spark-кластер живёт ровно столько, сколько нужно для обработки, и гасится.

Числовой пример для платформы обработки 1 000 событий/сек в среднем:

Статья Kappa (24/7 streaming) Lambda (batch + scheduled stream)
Spark executors (средн.) 20 × $0.10/ч × 24ч = $48/сут Batch: 40 × $0.10 × 3ч + Speed: 5 × $0.10 × 8ч = $16/сут
Driver node $0.20/ч × 24ч = $4.80/сут Пропорционально = $1.00/сут
Compute итого/сутки ~$53 ~$17
Compute итого/год ~$19 000 ~$6 200

Это упрощённая модель. Реальность добавляет нюансы: Kappa при K8s autoscaling может масштабироваться до нуля в период тишины (Serverless Spark). Lambda в свою очередь несёт overhead на запуск/остановку кластера (5–10 минут) при каждом запуске. Но порядок величин сохраняется: постоянный streaming обычно дороже scheduled batch при равном объёме данных.

4.2 Storage: S3 vs Kafka - кардинальная разница в стоимости

Хранение данных - второй ключевой компонент TCO:

Стоимость хранения 1TB данных за год:

  • S3 Standard: ~$276/год ($23 × 12)
  • S3 Infrequent Access: ~$150/год (для холодных исторических данных)
  • Self-hosted Kafka (3x replication): ~$6 500–10 800/год ($450–900 × 12)
  • Confluent Cloud: ~$3 240/год ($270 × 12)

Kafka в ~10–40 раз дороже S3 за один и тот же объём данных. Для небольших объёмов (< 1TB) это несущественно. Для платформ с терабайтами в день и retention в месяцы - это миллионы долларов разницы в год.

4.3 Скрытые затраты: операционная сложность

TCO включает не только инфраструктуру, но и человеческое время:

Вывод о TCO: в большинстве случаев Kappa имеет более высокую стоимость хранения (Kafka retention) и сопоставимую или ниже стоимость compute (один кластер vs два). Lambda имеет более низкую стоимость хранения (S3) и выше стоимость разработки (двойная кодовая база). Какая итоговая сумма больше - зависит от объёма данных, частоты изменений логики и зрелости команды.


5. Паттерн Incremental Batch: лучшее из двух миров

5.1 Trigger.AvailableNow: стриминговый код, батчевая экономика

Один из самых мощных паттернов современного Data Engineering - Incremental Batch (его также называют «микро-Lambda» или «Kappa по расписанию»). Идея: пишем код в стиле Structured Streaming (с checkpoint, watermark, stateful операциями), но запускаем его не 24/7, а по расписанию через Airflow или cron.

Ключевой инструмент - trigger(availableNow=True): Spark обрабатывает все данные, накопившиеся с последнего запуска, и завершается. Checkpoint сохраняет прогресс - следующий запуск начнётся точно с того места, где остановился предыдущий.

Incremental Batch экономит деньги и сохраняет гибкость: кластер живёт только во время обработки (обычно 2–5 минут при нормальном lag), checkpoint хранит прогресс точно как в стриминге, exactly-once гарантируется механизмом Iceberg+WAL.

5.2 Код Incremental Batch: практический пример

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

spark = SparkSession.builder \
    .appName("IncrementalBatch-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") \
    .getOrCreate()

event_schema = StructType([
    StructField("event_id",   StringType(),    False),
    StructField("user_id",    StringType(),    False),
    StructField("event_type", StringType(),    False),
    StructField("amount",     DoubleType(),    True),
    StructField("event_ts",   TimestampType(), False),
])

def apply_business_rules(df):
    """
    Единая функция трансформации.
    Работает одинаково в ЛЮБОМ режиме:
    - Continuous Streaming (Kappa 24/7)
    - Incremental Batch (trigger=availableNow, запуск каждые 15 мин)
    - Full Replay (startingOffsets=earliest, однократно)
    - Обычный Batch (spark.read.format("kafka"))
    """
    return (
        df
        .filter(F.col("amount") > 0)
        .withColumn("net_amount",
            F.col("amount") * F.lit(1.0)   # упрощено; реальная логика здесь
        )
        .withColumn("event_date", F.to_date("event_ts"))
    )


def process_microbatch(batch_df, batch_id: int) -> None:
    """foreachBatch: одинаков для Continuous и Incremental Batch режимов."""
    if batch_df.isEmpty():
        return

    enriched = apply_business_rules(batch_df.dropDuplicates(["event_id"]))
    enriched.cache()

    count = enriched.count()
    print(f"[Batch {batch_id}] {count} событий")

    # Bronze: append-only
    (
        enriched
        .writeTo("lakehouse.bronze.events")
        .append()
    )

    # Gold: агрегаты
    (
        enriched
        .groupBy("event_date", "event_type")
        .agg(F.sum("net_amount").alias("revenue"), F.count("*").alias("count"))
        .writeTo("lakehouse.gold.daily_revenue")
        .overwritePartitions()
    )

    enriched.unpersist()


# ============================================================
# РЕЖИМ 1: Continuous Streaming (Kappa 24/7)
# Используется когда: задержка < 1 минуты критична
# Стоимость: кластер работает 24/7
# ============================================================
def run_continuous():
    query = (
        spark.readStream
             .format("kafka")
             .option("kafka.bootstrap.servers", "kafka:9092")
             .option("subscribe", "ecommerce.events")
             .option("startingOffsets", "latest")
             .option("maxOffsetsPerTrigger", "50000")
             .load()
             .select(F.from_json(F.col("value").cast("string"), event_schema).alias("d"))
             .select("d.*")
             .writeStream
             .trigger(processingTime="30 seconds")   # Непрерывная обработка
             .foreachBatch(process_microbatch)
             .option("checkpointLocation", "s3://bucket/cp/continuous-v1")
             .start()
    )
    query.awaitTermination()


# ============================================================
# РЕЖИМ 2: Incremental Batch (запускается по расписанию Airflow)
# Используется когда: задержка 5-30 минут приемлема
# Стоимость: кластер живёт только во время обработки (2-10 минут)
# Запуск: airflow dag → spark-submit → этот скрипт → завершается
# ============================================================
def run_incremental_batch():
    """
    Запускается Airflow DAG каждые 15 минут.
    Spark читает из Kafka всё, что накопилось с последнего запуска,
    обрабатывает и завершается.
    Checkpoint хранит прогресс между запусками.
    """
    query = (
        spark.readStream
             .format("kafka")
             .option("kafka.bootstrap.servers", "kafka:9092")
             .option("subscribe", "ecommerce.events")
             # startingOffsets игнорируется после первого запуска:
             # checkpoint хранит точную позицию последнего обработанного offset
             .option("startingOffsets", "earliest")  # Только при первом запуске
             .option("maxOffsetsPerTrigger", "500000")  # Крупный батч: все за 15 минут
             .load()
             .select(F.from_json(F.col("value").cast("string"), event_schema).alias("d"))
             .select("d.*")
             .writeStream
             # КЛЮЧЕВОЕ ОТЛИЧИЕ: availableNow = обработать всё накопленное и завершиться
             # Spark разбивает на несколько микро-батчей при большом объёме
             .trigger(availableNow=True)
             .foreachBatch(process_microbatch)
             # Тот же checkpoint path - Spark продолжает с последней позиции Kafka
             .option("checkpointLocation", "s3://bucket/cp/incremental-v1")
             .start()
    )

    # awaitTermination() завершится автоматически после обработки всех доступных данных
    query.awaitTermination()
    print("Incremental Batch завершён. Кластер можно гасить.")


# ============================================================
# РЕЖИМ 3: Full Replay (при смене бизнес-логики)
# Используется: при критическом изменении бизнес-правил
# Запускается: вручную, один раз
# ============================================================
def run_full_replay(new_version: str = "v2"):
    """
    Пересчёт истории с нуля с новой бизнес-логикой.
    Пишет в новые таблицы (_v2), потом переключается View.
    """
    query = (
        spark.readStream
             .format("kafka")
             .option("kafka.bootstrap.servers", "kafka:9092")
             .option("subscribe", "ecommerce.events")
             .option("startingOffsets", "earliest")   # С самого начала
             .option("maxOffsetsPerTrigger", "1000000")  # Максимальная скорость
             .load()
             .select(F.from_json(F.col("value").cast("string"), event_schema).alias("d"))
             .select("d.*")
             .writeStream
             .trigger(availableNow=True)  # Обработать всю историю и остановиться
             .foreachBatch(lambda df, bid: process_microbatch_v2(df, bid, new_version))
             # НОВЫЙ checkpoint: начинаем с нуля, не используем старый прогресс
             .option("checkpointLocation", f"s3://bucket/cp/replay-{new_version}")
             .start()
    )
    query.awaitTermination()
    print(f"Full Replay завершён. Переключаем View на {new_version}")


def process_microbatch_v2(batch_df, batch_id: int, version: str) -> None:
    """Аналог process_microbatch, пишет в таблицы с суффиксом _v2."""
    if batch_df.isEmpty():
        return
    enriched = apply_business_rules(batch_df.dropDuplicates(["event_id"]))
    (
        enriched
        .writeTo(f"lakehouse.bronze.events_{version}")
        .append()
    )


# Выбор режима через переменную окружения (управляется Airflow)
import os
MODE = os.environ.get("PIPELINE_MODE", "incremental")

if MODE == "continuous":
    run_continuous()
elif MODE == "incremental":
    run_incremental_batch()
elif MODE == "replay":
    run_full_replay(os.environ.get("REPLAY_VERSION", "v2"))

Единая кодовая база apply_business_rules и process_microbatch работает во всех трёх режимах. Это квинтэссенция Modern Data Engineering: не Lambda, не Kappa, а адаптивный пайплайн, который выбирает режим исполнения в зависимости от требований.

5.3 Сравнение трёх режимов в таблице

Параметр Continuous (Kappa) Incremental Batch Full Replay
Задержка данных Секунды Минуты (интервал cron) Часы (однократно)
Compute стоимость 24/7 × полный кластер N мин/час × кластер Одноразовый большой кластер
Сложность state Высокая (watermarks) Средняя Низкая (нет state)
Надёжность Непрерывный процесс Запуск-остановка Запуск-остановка
Когда использовать SLA < 1 мин SLA 5–30 мин Смена логики
Checkpoint Обязателен, долгоживущий Обязателен, обновляется каждый запуск Новый путь каждый раз

6. Lakehouse как точка конвергенции

6.1 Как Iceberg размывает границу между Batch и Stream

До появления транзакционных форматов (Iceberg, Delta, Hudi) архитектурный выбор Batch vs Stream был жёстким: Batch пишет в HDFS-файлы, Stream пишет в HBase или Cassandra, и они никогда не встречаются в одной таблице.

Iceberg изменил это. Одна Iceberg-таблица может одновременно:

  • Получать данные от Streaming Job (новые снапшоты каждые 30 секунд)
  • Читаться аналитиком через Spark SQL или Trino
  • Получать данные от Batch Job (MERGE INTO раз в сутки)
  • Читаться инкрементально другим Streaming Job (через Iceberg Incremental Read)

Iceberg Snapshot Isolation позволяет Streaming Job и Batch Job писать в одну таблицу параллельно без блокировок и corruption. Это устраняет главное техническое ограничение классической Lambda: в Hadoop-эпоху параллельная запись была небезопасной и требовала разделения таблиц (Speed Layer → HBase, Batch Layer → Hive).

6.2 Iceberg Incremental Read: стриминг читает из Iceberg

Iceberg поддерживает инкрементальное чтение: Streaming Job может читать только снапшоты, появившиеся после последнего запуска. Это позволяет выстраивать цепочки стриминговых пайплайнов через Iceberg-таблицы:

# Streaming Job 1: читает из Kafka, пишет в Bronze Iceberg (append)
# (код как в предыдущих примерах)

# Streaming Job 2: читает из Bronze Iceberg, делает агрегации, пишет в Silver
# Используется Iceberg Incremental Read: читает только НОВЫЕ снапшоты
silver_from_bronze = (
    spark.readStream
         .format("iceberg")
         .option("stream-from-timestamp",
                 "2026-05-23T00:00:00")   # Стартовый момент (только при первом запуске)
         # После первого запуска checkpoint хранит ID последнего прочитанного снапшота
         .table("lakehouse.bronze.ecommerce_events")
)

# Применяем Silver трансформации
silver_query = (
    silver_from_bronze
    .filter(F.col("event_type") == "purchase")
    .groupBy(
        F.to_date("event_ts").alias("date"),
        "event_type"
    )
    .agg(
        F.sum("amount").alias("daily_revenue"),
        F.count("*").alias("order_count")
    )
    .writeStream
    .trigger(processingTime="5 minutes")  # Silver обновляется каждые 5 минут
    .outputMode("complete")
    .format("iceberg")
    .option("checkpointLocation", "s3://bucket/cp/silver-from-bronze")
    .toTable("lakehouse.silver.daily_revenue")
)

Цепочка через Iceberg - это мощный паттерн: каждый слой Medallion читает предыдущий через Incremental Read, не зная ничего о Kafka. Kafka-стриминг существует только в одном месте (Bronze инжестия), все остальные слои работают через ACID Iceberg-таблицы.


7. Анти-паттерны и мифы

7.1 Миф 1: «Kappa всегда лучше Lambda - это же современная архитектура»

Самый распространённый миф, порождённый конференционными докладами. Реальность:

  • Netflix использует гибридную архитектуру с активным batch-processing для ML-обучения
  • Airbnb перешла на Lambda-подобную систему для финансовой отчётности именно потому, что стриминговый Replay занимал слишком долго
  • Uber использует оба подхода в разных системах в зависимости от SLA

Правда: Lambda и Kappa решают разные задачи. Выбор зависит от конкретных требований, не от модности.

7.2 Миф 2: «В Kappa нет батча - значит нет исторической обработки»

Это неверное понимание Kappa. В Kappa батч заменяется Replay - это другая форма той же идеи, а не её отсутствие. Replay через trigger=availableNow - это функционально эквивалент batch ETL.

Ограничение: Replay ограничен временем хранения в Kafka. Если нужна история за 5 лет, а retention = 90 дней - за пределами retention Replay невозможен. Решение: экспортировать Kafka-лог в S3 (через Kafka S3 Sink Connector или MirrorMaker2) и делать Replay из S3 через тот же Streaming API.

7.3 Миф 3: «Dual pipeline в Lambda = двойные расходы»

Не обязательно. Batch-пайплайн в Lambda запускается ночью и занимает 3–4 часа. Speed Layer работает постоянно, но с небольшим кластером (только последние часы данных). Суммарная стоимость часто оказывается сопоставима с Kappa-кластером 24/7 или даже ниже, если объём данных большой.

«Двойные расходы» - это расходы на разработку и поддержку: два репозитория, два CI/CD пайплайна, две системы мониторинга. Это реальная, но не инфраструктурная проблема.

7.4 Анти-паттерн: гибрид без дисциплины

Хуже Lambda или Kappa - неструктурированный гибрид: часть данных течёт через Kafka, часть через HDFS, бизнес-логика размазана между тремя пайплайнами, никто не знает, какой источник «правда».

Правильный гибрид - это не «Lambda и Kappa смешались в хаос», а единый Iceberg Lakehouse, куда разные пайплайны пишут в соответствии с чёткими ролями: batch пишет точные исторические агрегаты, streaming пишет горячие инкременты, оба в одни и те же таблицы с ACID.


8. Engineering Decision Framework: алгоритм выбора

8.1 Многомерная матрица решений

8.2 Итоговая шпаргалка архитектора

Критерий Выбирай Lambda/Batch Выбирай Incremental Batch Выбирай Kappa
Задержка SLA > 1 часа 15 мин – 1 час < 15 минут
Объём данных > 10TB/сутки 100GB – 10TB < 100GB
История для Replay > 1 года 1–12 месяцев < 3 месяцев
Kafka retention Нет / 7 дней 30–90 дней 90+ дней
Сложность join/agg Тяжёлые исторические Средние Лёгкие / windowed
Зрелость команды Batch experts Смешанная Streaming experts
Частота смены логики Редко Умеренно Часто
Стоимость инфра Низкая (S3) Средняя Высокая (Kafka долго)
Observability Проще Средне Требует зрелого стека

9. Production кейс-стади: три реальных сценария

9.1 Кейс 1: Антифрод-система в финтех банке

Контекст: 5 млн транзакций в сутки, 150 транзакций/сек в пике. Требование: решение о блокировке транзакции в течение 300 миллисекунд от момента инициации.

Анализ:

  • SLA 300ms → только continuous streaming, никакого batching
  • Объём 150 tx/сек × 86400 = ~13M событий/сутки ≈ 5GB - небольшой объём
  • Логика: stateful (история транзакций пользователя за 24 часа), pattern-matching по IP, device fingerprint
  • Команда: выделенная ML-платформа с опытом streaming

Решение: Kappa обязательна

Почему Lambda проиграла бы: даже «быстрый» Speed Layer Lambda не даёт 300ms. Batch Layer здесь вообще не применим - он нужен для исторической аналитики, а не для принятия решений в реальном времени. Kappa - единственный разумный выбор.

Kafka retention всего 7 дней: для антифрода не нужен Replay за год. Бизнес-логика модели меняется через MLOps-пайплайн (переобучение модели), а не через пересчёт исторических данных. Короткий retention экономит инфраструктуру.

9.2 Кейс 2: Корпоративный DWH для финансовой отчётности

Контекст: производственная компания, 50 филиалов, данные из SAP. Финансовая отчётность по МСФО требует аудиторской точности. 2TB новых данных в сутки, 500TB исторического архива. Отчёты нужны к 9:00 следующего рабочего дня.

Анализ:

  • SLA: к утру следующего дня - 8+ часов задержки приемлема
  • Объём: 500TB история → Replay займёт недели
  • Логика: сложные МСФО-стандарты, многоуровневые join SAP-таблиц, retention 10 лет по закону
  • Команда: BI-разработчики, DWH-архитекторы, без streaming expertise

Решение: Lambda / классический Batch

Почему Kappa проиграла бы: хранить 500TB в Kafka при retention на 10 лет = $50 000+/месяц только на storage. Lambda с 500TB на S3 IA = ~$6 250/месяц. Разница - $530 000 в год только на хранении данных. Плюс: МСФО требует auditable, воспроизводимых расчётов - batch с детерминированными джобами идеально соответствует этому требованию.

9.3 Кейс 3: Кликстрим маркетплейса - победа Incremental Batch

Контекст: e-commerce маркетплейс, 10M пользователей, 5 000 событий/сек в пике. Требование: обновление дашборда продавцов с задержкой до 15 минут. Команда из 3 Data Engineers, опыт batch, изучают streaming.

Анализ:

  • SLA: 15 минут - не требует continuous streaming
  • Объём: 5K событий/сек × 86400 = 432M событий/сутки ≈ 200GB - средний объём
  • Логика: воронка конверсии, сессии, выручка по категориям - несложно
  • Команда: batch experts, ограниченный опыт streaming

Решение: Incremental Batch (trigger=availableNow каждые 10 минут)

# DAG в Apache Airflow: запускает Spark каждые 10 минут
# Spark читает из Kafka всё накопившееся, обрабатывает, завершается

# Почему это дешевле Kappa 24/7:
# - Среднее время обработки 10-минутного инкремента: 2-3 минуты
# - Кластер работает 2-3 из каждых 10 минут = 20-30% времени
# - Kappa 24/7 = 100% времени
# - Экономия на compute: 70-80% при тех же данных

# Пример конфигурации:
AIRFLOW_SCHEDULE = "*/10 * * * *"   # Каждые 10 минут

# Spark-джоб (тот же run_incremental_batch() из примера выше)
# Checkpoint хранит прогресс между запусками
# failOnDataLoss=false: если Kafka удалил старые сегменты между запусками

# Результат:
# - Дашборды обновляются с задержкой 10-15 минут
# - Кластер работает 20-30% времени (vs 100% при Kappa)
# - Экономия: ~$15 000/год при стоимости кластера $2.50/час
# - Команда использует знакомый batch-код с минимальным streaming overhead

Ключевой инсайт кейса: при SLA 15 минут Incremental Batch экономит 70–80% compute-расходов по сравнению с Kappa 24/7 - при тех же функциональных результатах. Это деньги, которые можно потратить на развитие продукта.


10. Эволюция платформы: от Lambda к Unified Streaming

10.1 Типичная траектория зрелой платформы

Большинство production-систем не рождаются «чистыми» Lambda или Kappa - они эволюционируют:

Стадия 3 (Lakehouse Lambda) - где находятся большинство зрелых платформ сегодня. Iceberg/Delta устранили главные технические болячки классической Lambda (параллельная запись, consistence Serving Layer), но сохранили батч как основной инструмент исторической обработки.

Стадия 4 (Hybrid/Incremental) - «высший пилотаж» Modern Data Stack: один код, три режима исполнения. Большинство платформ к этому придут в 2025–2027 годах по мере зрелости Spark Structured Streaming.

Стадия 5 (Streaming-first / полная Kappa) - реальна для платформ с умеренными объёмами (< 1TB/сутки), коротким горизонтом хранения (< 90 дней) и зрелой streaming-командой. Для большинства enterprise-платформ на этой стадии останавливаться неоправданно дорого.

10.2 Финальная рамка: пять вопросов перед архитектурным выбором

Пять вопросов - это структура архитектурного design review. Ответ на каждый даёт вектор в сторону Lambda, Kappa или Hybrid. Суммируя векторы - получаем обоснованный выбор.

Главный вывод всего блока уроков: Lambda и Kappa - не «старое» и «новое». Это два инженерных инструмента с разными trade-offs. Lakehouse на Apache Iceberg и Spark Structured Streaming с 2022 года дают нам третий путь - адаптивный пайплайн, который меняет режим исполнения без изменения кода. Именно туда движется отрасль, и именно это разделяет Middle Data Engineer, выбирающего «что знаю», от Senior Data Engineer, выбирающего «что подходит».