Data Engineering и ETL Потоковая обработка: Kafka, Flink, окна и семантика доставки
0%

Потоковая обработка: Kafka, Flink, окна и семантика доставки

Потоковая обработка: Kafka, Flink, окна и семантика доставки

В пакетной обработке мир устроен уютно: есть конечный набор файлов, есть момент «данные за вчера готовы», есть право пересчитать всё с нуля, если что-то сломалось. Потоковая обработка забирает эти удобства и взамен даёт секунды задержки вместо часов. Обмен получается неравноценным по сложности: почти вся тяжесть стриминга — это борьба с тем, что у потока нет конца, а значит, нет естественного момента «теперь можно считать».

Эта статья — про то, как индустрия научилась считать корректные агрегаты над бесконечными данными: про лог как фундаментальную абстракцию, про разницу между временем события и временем обработки, про окна и водяные знаки, про состояние и его снимки, и про то, почему фраза «у нас exactly-once» без уточнений почти всегда неправда.

Зачем вообще поток: три причины, а не одна

Первая причина очевидна — латентность. Антифрод, который узнаёт о мошеннической транзакции через шесть часов, бесполезен. Динамическое ценообразование, лента рекомендаций, алертинг по метрикам, отслеживание доставки на карте — всё это ломается на батче не потому, что батч «медленный», а потому, что ценность данных падает быстрее, чем их успевают посчитать.

Вторая причина тоньше — сглаживание нагрузки. Ночной батч на 200 машин создаёт пик потребления и держит железо простаивающим весь день. Стриминг размазывает ту же работу равномерно: инкремент маленький, состояние живёт в памяти, повторной свёртки терабайтов не происходит. Часто поток дешевле батча при той же общей пропускной способности, просто дешевизна прячется в утилизации.

Третья причина — архитектурная. Лог событий становится точкой интеграции между сервисами, и это меняет способ, которым системы обмениваются данными. Вместо десятков ночных выгрузок «система A читает базу системы B» появляется один контракт: сервис публикует факты о том, что у него произошло, а все заинтересованные читают их независимо, в своём темпе. Это ровно та же идея, что и в ETL vs ELT, только источник истины — не таблица, а упорядоченная последовательность изменений.

Ключевой текст, который стоит прочитать целиком: Jay Kreps, «The Log: What every software engineer should know about real-time data’s unifying abstraction» — https://engineering.linkedin.com/distributed-systems/log-what-every-software-engineer-should-know-about-real-time-datas-unifying. Из него выросли и Kafka, и половина современной дата-архитектуры.

Лог: абстракция, на которой всё держится

Лог — это append-only последовательность записей с монотонно возрастающими номерами. Всё. Никаких обновлений на месте, никакого удаления из середины. Из этой скудости растут все полезные свойства:

  • Порядок объективен. Внутри лога «раньше» и «позже» определены однозначно, без синхронизации часов между машинами.
  • Читатели независимы. Позиция чтения (offset) — состояние потребителя, а не лога. Десять потребителей читают один лог, не мешая друг другу и не копируя данные.
  • Воспроизводимость. Пока данные не удалены по retention, любой потребитель может вернуться назад и перечитать историю. Это и есть основа восстановления после сбоя и backfill’а новой логики.
  • Лог = материализованная таблица. Свернув лог изменений по ключу, получаем текущее состояние; развернув таблицу в поток изменений, получаем лог. Эта двойственность (stream–table duality) — центральная идея Kafka Streams и Flink SQL.

Устройство Kafka в объёме, который реально нужен инженеру

Топик разбит на партиции; партиция — это и есть физический лог на диске (сегменты файлов + индекс). Гарантия порядка даётся только внутри партиции, и это первое, обо что все спотыкаются. Ключ сообщения определяет партицию (hash(key) % partitions в дефолтном партиционере), поэтому все события одного пользователя/заказа/устройства попадают в одну партицию и обрабатываются по порядку — при условии, что ключ выбран правильно.

Партиция реплицируется на N брокеров; один из них — лидер, остальные — фолловеры. Множество реплик, догнавших лидера, называется ISR (in-sync replicas). Продюсер с acks=all получает подтверждение только после того, как запись легла во все ISR, а параметр min.insync.replicas задаёт, сколько реплик обязано быть живыми, чтобы запись вообще принималась. Классическая продакшн-комбинация: replication.factor=3, min.insync.replicas=2, acks=all — переживаем потерю одного брокера без потери данных и без остановки записи.

Потребители объединяются в consumer group: каждая партиция в момент времени назначена ровно одному потребителю группы. Отсюда жёсткий потолок параллелизма: потребителей в группе больше, чем партиций, — лишние простаивают. Количество партиций поэтому проектируется под будущий параллелизм (увеличить можно, уменьшить — нет, и увеличение ломает соответствие «ключ → партиция» для уже записанных данных).

Retention бывает двух видов. Обычный — по времени/размеру (retention.ms, retention.bytes): старое просто удаляется. Compaction (cleanup.policy=compact) — хранить последнее значение для каждого ключа бесконечно; так топик превращается в реплицируемую key-value таблицу, из которой можно восстановить состояние с нуля. На compacted-топиках строятся changelog’и Kafka Streams и справочники для join’ов.

# создаём топик под события заказов: 24 партиции, RF=3, неделя хранения
kafka-topics.sh --bootstrap-server kafka:9092 --create \
  --topic orders.events.v1 --partitions 24 --replication-factor 3 \
  --config min.insync.replicas=2 \
  --config retention.ms=604800000 \
  --config compression.type=zstd

# главная операционная метрика стриминга — лаг группы (в сообщениях на партицию)
kafka-consumer-groups.sh --bootstrap-server kafka:9092 \
  --describe --group orders-enricher

# перемотка группы на 2 часа назад: аварийный реплей после багфикса
kafka-consumer-groups.sh --bootstrap-server kafka:9092 \
  --group orders-enricher --topic orders.events.v1 \
  --reset-offsets --by-duration PT2H --execute

Общая схема потоковой платформы

Три вещи, которые видно на схеме и которые обычно забывают в первой версии пайплайна: отдельный слой сырых данных (без него нельзя переиграть историю после исправления парсера), DLQ для сообщений, которые не разобрались, и реестр схем как единый контракт. Про контракты и их эволюцию подробно — в статье про качество данных и governance.

Время: главный источник боли

У каждого события есть минимум два времени:

  • event time — когда событие произошло в реальности (часы устройства или сервиса-источника);
  • processing time — когда его увидел оператор потока;
  • (и часто третье — ingestion time, когда запись попала в Kafka; его ставит брокер, и оно монотонно, но не отражает реальность на клиенте).

Считать по processing time просто и бесполезно: результат зависит от того, как сегодня работала сеть и не перезапускался ли джоб. Один и тот же входной поток даст разные ответы на разных прогонах — то есть пайплайн недетерминирован, а значит, его нельзя ни протестировать, ни воспроизвести, ни сравнить с батчевым расчётом. Считать по event time правильно, но тогда возникает вопрос: сколько ждать опоздавших?

Время события и время обработки

Водяные знаки (watermarks)

Водяной знак W(t) — это утверждение системы: «все события с event time ≤ W уже пришли». Это эвристика, а не факт. Типичная стратегия — «ограниченная неупорядоченность»: watermark = (максимальный увиденный event time) − допуск, например 5 секунд. Как только watermark пересекает правую границу окна, окно считается полным и результат эмитируется.

Отсюда фундаментальный компромисс, из которого нет выхода:

Допуск watermark Задержка результата Полнота результата Стоимость состояния
маленький (1 с) низкая часть событий не успевает маленькое
большой (10 мин) высокая почти все события учтены большое
+ allowed lateness низкая, но результат уточняется максимальная большое дольше

Третья строка — самая практичная: выдать предварительный результат рано, а затем обновить его, когда придут опоздавшие. Это работает, только если приёмник умеет апдейты (upsert-таблица, ClickHouse с ReplacingMergeTree, Iceberg с merge-on-read) — иначе получите дубли вместо коррекции.

Отдельная ловушка: watermark движется по минимуму среди всех входных партиций. Если у топика 24 партиции и в одну из них давно ничего не пишут (например, регион ночью), watermark всего джоба замирает, окна не закрываются, состояние растёт. Лечится «idle source» настройкой (withIdleness в Flink), которая исключает молчащие партиции из расчёта.

Окна: как нарезать бесконечность

Типы окон на оси времени события

  • Tumbling — фиксированный размер, без пересечений. «Сколько заказов в каждую минуту». Каждое событие в ровно одном окне.
  • Sliding (hopping) — фиксированный размер, шаг меньше размера. «Средний чек за последние 10 минут, обновляется каждую минуту». Каждое событие попадает в размер / шаг окон — это прямой множитель к объёму состояния и к работе на событие. Sliding 24 ч с шагом 1 мин = 1440 копий каждого события. Так делать нельзя; вместо этого считают инкрементальный агрегат.
  • Session — границы задаёт сам поток: окно закрывается, если в течение gap не было событий. Пользовательские сессии, эпизоды просмотра. Технически сложнее всего: приход события в разрыв между двумя сессиями требует слияния окон.
  • Global / custom triggers — окно на весь поток плюс собственный триггер (по счётчику, по признаку в данных). Так реализуют, например, «закрыть окно, когда пришло событие end_of_session».

Жизненный цикл окна

Обратите внимание на переход «Отброшено». В проде этот путь обязан быть измеряемым: счётчик numLateRecordsDropped — одна из немногих метрик, по которой сразу видно, что данные тихо теряются. Если он ненулевой и растёт — либо допуск watermark мал, либо у источника проблемы с часами.

Сложность и объём состояния

Пусть K — число активных ключей, W — число окон, живых одновременно на ключ, S — средний размер агрегата.

  • Память: O(K · W · S). Для инкрементальных агрегатов (count, sum, min/max, HyperLogLog) S — константа; для окон, хранящих список событий (медиана, топ-N по сырым данным), S растёт с трафиком, и это главная причина взрывного роста состояния.
  • Работа на событие: O(W) — событие раскладывается по всем накрывающим окнам. Для sliding это размер/шаг.
  • Срабатывание: O(K_expired) на продвижение watermark, обычно реализовано таймерами в приоритетной очереди — O(log n) на вставку/извлечение таймера.

Практическое следствие: перед оконным агрегатом всегда стоит спросить, можно ли заменить хранение сырых событий на инкрементальный аккумулятор или скетч (HyperLogLog для уникальных, t-digest для квантилей). Это разница между 200 МБ и 200 ГБ состояния при одинаковом ответе с точностью 1%.

Семантика доставки: что означает каждый уровень

Разберём точно, что где происходит.

At-most-once. Оффсет коммитится до обработки; при падении сообщение теряется. Легально ровно в одном случае — когда данные заведомо избыточны (сэмплированная телеметрия), а потеря дешевле сложности.

At-least-once. Оффсет коммитится после успешной обработки. Дубли возможны при любом падении между записью в приёмник и коммитом оффсета. Это дефолтный режим 90% продакшн-пайплайнов, и он совершенно нормален — при условии, что приёмник идемпотентен.

Exactly-once. Здесь начинается терминологическая ложь. Никакая распределённая система не может гарантировать, что сообщение будет доставлено ровно один раз — это доказанный факт (проблема двух генералов). Что реально гарантируется — exactly-once processing semantics: эффект на состояние и на выход выглядит так, как будто каждое сообщение обработано один раз. Достигается это двумя механизмами:

  1. Идемпотентный producer (enable.idempotence=true): продюсеру выдаётся PID, каждой записи — sequence number на партицию; брокер отбрасывает дубли ретраев. Защищает от дублей на участке producer → брокер. Включён по умолчанию начиная с Kafka 3.0.
  2. Транзакции (transactional.id): запись в несколько партиций + коммит оффсетов потребления оформляются как одна атомарная операция (двухфазный коммит через координатора транзакций). Читатель с isolation.level=read_committed не видит незакоммиченные данные.

Ограничение, которое важно понимать: транзакции Kafka охватывают только Kafka. Как только пайплайн пишет во внешнюю систему (Postgres, S3, HTTP API), гарантия распадается, если приёмник не участвует в 2PC или не идемпотентен. Именно поэтому в проде «exactly-once» почти всегда реализуется как at-least-once + дедупликация на приёмнике:

  • уникальный ключ и INSERT ... ON CONFLICT DO NOTHING / MERGE;
  • ReplacingMergeTree в ClickHouse (схлопывание по версии);
  • запись во временный файл + атомарный rename в объектном хранилище;
  • Iceberg/Delta с идемпотентным commit по txn-id.
"""
At-least-once потребитель с идемпотентной записью — рабочая лошадка продакшна.
Ключевая идея: коммитим оффсет ТОЛЬКО после того, как данные надёжно в БД,
а запись делаем так, чтобы повтор был безвреден.
"""
from confluent_kafka import Consumer, KafkaError
import json, psycopg2

consumer = Consumer({
    "bootstrap.servers": "kafka:9092",
    "group.id": "orders-enricher",
    "enable.auto.commit": False,        # ручной коммит — обязательное условие
    "auto.offset.reset": "earliest",    # при первом старте читаем всё, а не только новое
    "isolation.level": "read_committed",
    "max.poll.interval.ms": 300_000,    # запас на медленный батч; иначе группу перебалансирует
})
consumer.subscribe(["clean.orders.v1"])

db = psycopg2.connect("postgresql://user:pass@analytics:5432/marts")

UPSERT = """
INSERT INTO order_facts (event_id, order_id, user_id, amount, event_ts)
VALUES (%s, %s, %s, %s, %s)
ON CONFLICT (event_id) DO NOTHING
"""  # event_id — сквозной идентификатор события от источника; дубль просто не запишется

BATCH = 500
buffer = []

try:
    while True:
        msgs = consumer.consume(num_messages=BATCH, timeout=1.0)
        for msg in msgs or []:
            if msg.error():
                if msg.error().code() == KafkaError._PARTITION_EOF:
                    continue
                raise RuntimeError(msg.error())
            e = json.loads(msg.value())
            buffer.append((e["event_id"], e["order_id"], e["user_id"],
                           e["amount"], e["event_ts"]))

        if buffer:
            with db:                       # транзакция БД: всё или ничего
                with db.cursor() as cur:
                    cur.executemany(UPSERT, buffer)
            buffer.clear()
            consumer.commit(asynchronous=False)  # только теперь двигаем оффсет
finally:
    consumer.close()
    db.close()

Порядок операций здесь — не стилистика, а весь смысл. Поменяйте местами commit и запись в БД — получите at-most-once и тихую потерю данных при рестарте пода.

"""
Read-process-write внутри Kafka с транзакциями (EOS).
Работает только когда и вход, и выход — Kafka.
"""
from confluent_kafka import Consumer, Producer, TopicPartition
import json

producer = Producer({
    "bootstrap.servers": "kafka:9092",
    "enable.idempotence": True,
    "transactional.id": "enricher-tx-1",   # уникален на инстанс, стабилен между рестартами
    "acks": "all",
})
consumer = Consumer({
    "bootstrap.servers": "kafka:9092",
    "group.id": "enricher",
    "enable.auto.commit": False,
    "isolation.level": "read_committed",
})
consumer.subscribe(["clean.orders.v1"])
producer.init_transactions()

while True:
    msgs = consumer.consume(num_messages=1000, timeout=1.0)
    if not msgs:
        continue
    producer.begin_transaction()
    try:
        for m in msgs:
            if m.error():
                continue
            out = transform(json.loads(m.value()))
            producer.produce("enriched.orders.v1", key=m.key(),
                             value=json.dumps(out).encode())
        # оффсеты входа коммитятся ВНУТРИ транзакции — вместе с выходными записями
        offsets = [TopicPartition(m.topic(), m.partition(), m.offset() + 1) for m in msgs]
        producer.send_offsets_to_transaction(offsets, consumer.consumer_group_metadata())
        producer.commit_transaction()
    except Exception:
        producer.abort_transaction()
        raise

Оконные агрегаты, дедупликация, join’ы, паттерны CEP — всё это состояние, которое живёт внутри работающего джоба и может достигать десятков и сотен гигабайт. Задача движка: пережить падение любой машины, не потеряв состояние и не посчитав ничего дважды.

Flink решает её алгоритмом асинхронных барьеров (asynchronous barrier snapshotting) — вариантом классического Chandy–Lamport. Координатор просит источники вставить в поток специальную запись — барьер. Барьер течёт вместе с данными; оператор, получивший барьер по всем входам, сохраняет своё состояние и пропускает барьер дальше.

Барьеры контрольных точек в Flink

Практические следствия, которые определяют настройку джоба:

  • Backend состояния. hashmap — состояние в куче JVM, быстро, но ограничено памятью и грузит GC. rocksdb — состояние на локальном диске в LSM-дереве, медленнее на порядок, зато держит сотни ГБ и умеет инкрементальные чекпоинты (в снимок уходят только новые SST-файлы). Для состояния больше нескольких ГБ выбор один — RocksDB.
  • Выравнивание (alignment). При перекосе входов быстрый канал буферизуется, пока не придёт барьер по медленному. Под backpressure это может растянуть чекпоинт на минуты. Лечение — unaligned checkpoints: барьер обгоняет данные, а буферы in-flight записываются в снимок; чекпоинт становится быстрым, но толстым.
  • Checkpoint vs savepoint. Чекпоинт — служебный снимок для автоматического восстановления, живёт по политике retention и заточен под скорость. Savepoint — явный, портируемый снимок для обновления версии джоба, миграции состояния, смены параллелизма. Правило: любой деплой стриминг-джоба — это stop-with-savepoint → деплой → start-from-savepoint.
  • TTL состояния. Состояние без TTL растёт вечно. Для дедупликации по event_id за сутки нужен StateTtlConfig на 24–26 часов, иначе через месяц джоб упрётся в диск. Это ошибка №1 по частоте в продовых Flink-джобах.
# flink-conf.yaml: минимальный продакшн-набор для stateful-джоба
execution.checkpointing.interval: 60s
execution.checkpointing.min-pause: 30s          # не душим джоб чекпоинтами подряд
execution.checkpointing.timeout: 10min
execution.checkpointing.mode: EXACTLY_ONCE      # аligned barriers + 2PC в sink
execution.checkpointing.unaligned: true         # спасает при backpressure
execution.checkpointing.tolerable-failed-checkpoints: 3

state.backend.type: rocksdb
state.backend.incremental: true
state.checkpoints.dir: s3://flink-state/checkpoints
state.savepoints.dir: s3://flink-state/savepoints
state.checkpoints.num-retained: 5

restart-strategy.type: exponential-delay
restart-strategy.exponential-delay.initial-backoff: 10s
restart-strategy.exponential-delay.max-backoff: 2min
"""
PyFlink: оконный агрегат по event time с watermark и обработкой опоздавших.
Считаем выручку по продавцу в минутных tumbling-окнах.
"""
from pyflink.common import WatermarkStrategy, Duration, Types
from pyflink.common.watermark_strategy import TimestampAssigner
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.window import TumblingEventTimeWindows
from pyflink.common.time import Time

env = StreamExecutionEnvironment.get_execution_environment()
env.set_parallelism(24)   # обычно = числу партиций входного топика

class OrderTs(TimestampAssigner):
    def extract_timestamp(self, value, record_timestamp):
        return value["event_ts_ms"]        # миллисекунды event time из самого события

wm = (WatermarkStrategy
      .for_bounded_out_of_orderness(Duration.of_seconds(5))  # допуск неупорядоченности
      .with_timestamp_assigner(OrderTs())
      .with_idleness(Duration.of_minutes(1)))                # молчащие партиции не тормозят watermark

stream = source.assign_timestamps_and_watermarks(wm)

late_tag = OutputTag("late-orders", Types.PICKLED_BYTE_ARRAY())

result = (stream
          .key_by(lambda e: e["seller_id"])
          .window(TumblingEventTimeWindows.of(Time.minutes(1)))
          .allowed_lateness(Time.minutes(10))   # 10 минут окно ещё уточняется
          .side_output_late_data(late_tag)      # то, что опоздало сильнее — не теряем, а отводим
          .reduce(lambda a, b: {"seller_id": a["seller_id"],
                                "amount": a["amount"] + b["amount"]}))

result.sink_to(upsert_sink)                     # приёмник ОБЯЗАН уметь upsert: окно эмитится не раз
result.get_side_output(late_tag).sink_to(late_audit_sink)

Три строки здесь несут всю нагрузку: for_bounded_out_of_orderness задаёт компромисс «свежесть против полноты», allowed_lateness включает уточнение результата, side_output_late_data превращает молчаливую потерю в наблюдаемый поток, который можно посчитать и по которому можно завести алерт.

Большая часть продовых стриминг-задач сегодня пишется не на DataStream API, а на SQL. Это резко снижает порог входа и делает логику сравнимой с батчевой (см. моделирование данных).

-- Источник: топик Kafka со схемой Avro из Schema Registry
CREATE TABLE orders (
    order_id    STRING,
    seller_id   STRING,
    amount      DECIMAL(12, 2),
    event_ts    TIMESTAMP(3),
    -- водяной знак объявляется прямо в DDL: допуск 5 секунд
    WATERMARK FOR event_ts AS event_ts - INTERVAL '5' SECOND
) WITH (
    'connector' = 'kafka',
    'topic' = 'clean.orders.v1',
    'properties.bootstrap.servers' = 'kafka:9092',
    'properties.group.id' = 'sql-agg',
    'scan.startup.mode' = 'group-offsets',
    'format' = 'avro-confluent',
    'avro-confluent.url' = 'http://schema-registry:8081'
);

-- Приёмник с первичным ключом => upsert-семантика, окна можно переэмитить
CREATE TABLE seller_minute_revenue (
    seller_id    STRING,
    window_start TIMESTAMP(3),
    revenue      DECIMAL(18, 2),
    orders_cnt   BIGINT,
    PRIMARY KEY (seller_id, window_start) NOT ENFORCED
) WITH (
    'connector' = 'upsert-kafka',
    'topic' = 'agg.seller_minute_revenue.v1',
    'key.format' = 'json',
    'value.format' = 'json'
);

-- Оконные табличные функции (TVF) — современный синтаксис Flink SQL
INSERT INTO seller_minute_revenue
SELECT
    seller_id,
    window_start,
    SUM(amount)   AS revenue,
    COUNT(*)      AS orders_cnt
FROM TABLE(TUMBLE(TABLE orders, DESCRIPTOR(event_ts), INTERVAL '1' MINUTE))
GROUP BY seller_id, window_start, window_end;

Дедупликация — отдельный частый приём: берём первую строку по ключу в порядке времени.

-- Убираем дубли событий (последствие at-least-once на входе) за окно суток
SELECT order_id, seller_id, amount, event_ts
FROM (
    SELECT *,
           ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY event_ts ASC) AS rn
    FROM orders
) WHERE rn = 1;

Важно: состояние такой дедупликации хранит все увиденные event_id и требует TTL (table.exec.state.ttl), иначе растёт неограниченно.

Join’ы в потоке: три принципиально разных зверя

Джойн в батче — это просто. Джойн в потоке требует ответа на вопрос «сколько ждать вторую сторону», и ответ определяет объём состояния.

1. Stream–stream interval join. Соединяем два потока в пределах временного интервала: «клик и показ рекламы, если клик случился в течение 30 минут после показа». Состояние = обе стороны за ширину интервала; чем шире интервал, тем больше памяти.

SELECT i.impression_id, c.click_id, i.event_ts, c.event_ts
FROM impressions i
JOIN clicks c
  ON i.impression_id = c.impression_id
 AND c.event_ts BETWEEN i.event_ts AND i.event_ts + INTERVAL '30' MINUTE;

2. Stream–table (lookup / temporal join). Обогащаем поток справочником. Наивный вариант — синхронный запрос в БД на каждое событие: это убивает пропускную способность и создаёт зависимость от чужого SLA. Правильный — держать справочник как compacted-топик и материализовать его в локальное состояние джоба. Дополнительный нюанс — temporal join: соединять надо не с текущим состоянием справочника, а с тем, каким он был на момент события (курс валюты на момент оплаты, а не сегодняшний).

-- FOR SYSTEM_TIME AS OF: версионный джойн по времени события
SELECT o.order_id,
       o.amount * r.rate AS amount_usd
FROM orders AS o
JOIN currency_rates FOR SYSTEM_TIME AS OF o.event_ts AS r
  ON o.currency = r.currency;

3. Stream–stream «вечный» join по ключу. Например, заказ и его доставка, между которыми могут пройти недели. Состояние здесь неограниченно по природе задачи, и его нужно ограничивать явно: TTL + вывод «не дождались» как отдельное бизнес-событие. Попытка сделать такой джойн «как в батче» — типичный способ уронить кластер через месяц после запуска.

// Kafka Streams: KStream × KTable — обогащение потока материализованным справочником.
// KTable читается из compacted-топика и живёт в локальном RocksDB инстанса.
KTable<String, Seller> sellers = builder.table(
        "ref.sellers.v1",
        Consumed.with(Serdes.String(), sellerSerde),
        Materialized.as("sellers-store"));

builder.stream("clean.orders.v1", Consumed.with(Serdes.String(), orderSerde))
       .selectKey((k, order) -> order.getSellerId())   // ко-партиционирование по ключу джойна
       .join(sellers, (order, seller) -> enrich(order, seller))
       .to("enriched.orders.v1", Produced.with(Serdes.String(), enrichedSerde));

Ключевое требование Kafka Streams, о которое спотыкаются все: ко-партиционирование. Оба топика должны иметь одинаковое число партиций и одинаковую схему ключей, иначе join физически невозможен (данные лежат на разных инстансах) — и selectKey вызовет неявный repartition-топик со своей стоимостью.

Выбор движка

Короткая шпаргалка по выбору:

  • Kafka Streams / ksqlDB — если и вход, и выход в Kafka, а команда пишет на JVM и не хочет отдельного кластера. Это библиотека, а не сервис: масштабируется добавлением инстансов приложения.
  • Flink — если нужны продвинутые окна, большое состояние (сотни ГБ), точная работа с event time, CEP или сложные join’ы. Цена — отдельный кластер и реальная операционная экспертиза.
  • Spark Structured Streaming — если у вас уже есть Spark-платформа и приемлема задержка в единицы секунд (micro-batch). Continuous processing до сих пор экспериментален. Логика переиспользуется с батчем — сильный аргумент, см. пакетную обработку.
  • Streaming-базы (Materialize, RisingWave, ClickHouse с materialized views) — если задача формулируется как «инкрементально поддерживать результат SQL-запроса», а не как произвольный dataflow.

Типичные ошибки

  1. Ключ партиционирования выбран без учёта порядка. События одной сущности разлетаются по партициям — и «создан» приходит после «отменён». Ключ должен совпадать с сущностью, для которой важен порядок.
  2. Горячий ключ. 40% трафика — один merchant_id; одна партиция и один воркер перегружены, остальные простаивают. Лечится двухфазной агрегацией с солью ключа (key + random(0..N) на первом этапе, свёртка на втором).
  3. Состояние без TTL. Дедупликация, join’ы, session-окна растут вечно. Джоб живёт три недели и умирает по диску.
  4. Processing time вместо event time. Пайплайн становится недетерминированным; результаты нельзя ни воспроизвести, ни сверить с батчем.
  5. Приёмник не идемпотентен. At-least-once + обычный INSERT = дубли после каждого рестарта. Плюс каждое обновление окна по allowed lateness создаёт лишнюю строку.
  6. Слишком много партиций. Тысячи партиций на топик увеличивают время ребалансировки, нагрузку на контроллер и объём метаданных. Начинайте с числа, кратного планируемому параллелизму, а не с «на вырост в 100 раз».
  7. Нет DLQ. Одно сообщение с битым JSON останавливает потребителя, тот падает в бесконечный рестарт (poison pill), лаг растёт, дежурный ночью удаляет топик.
  8. Ребалансировка на каждом деплое. max.poll.interval.ms меньше времени обработки батча → группа считает потребителя мёртвым → stop-the-world ребаланс. Используйте cooperative-sticky-assignor и static membership (group.instance.id).
  9. Схема без реестра. Продюсер добавил обязательное поле — все потребители легли. Контракт должен быть машинно-проверяемым, с политикой совместимости BACKWARD.
  10. Мониторинг только по CPU. Единственная метрика, которая действительно говорит о здоровье стриминга, — consumer lag и его производная. Растущий лаг = система не справляется, независимо от загрузки CPU.

Как это выглядит в проде

Наблюдаемость. Обязательный минимум панелей: лаг по группам и партициям (и производная лага — «догоняем или отстаём»), длительность и размер чекпоинтов, backpressure по операторам, numLateRecordsDropped, размер состояния, частота рестартов, счётчик DLQ. Алерты ставятся на лаг в единицах времени (сколько минут отставания), а не в сообщениях — бизнес понимает минуты.

Деплой. stop-with-savepoint → выкатка новой версии → start-from-savepoint. При несовместимом изменении состояния (сменилась структура агрегата) — либо миграция через State Processor API, либо запуск с нуля с реплеем из Kafka. Именно ради второго варианта retention сырых топиков держат в 7–30 дней.

Реплей и backfill. Kafka retention конечен, поэтому историю дольше него держат в озере (форматы и хранилища). Отсюда две классические архитектуры: лямбда (батч-слой считает исторически точный результат, стриминг — быстрый приблизительный, витрина объединяет) и каппа (единственный код — стриминговый, история переигрывается через тот же движок из лога/озера). Каппа проще концептуально, лямбда чаще встречается в реальности, потому что переиграть 2 года данных через стриминговый джоб дорого. Практический компромисс: один и тот же Flink SQL, исполняемый в streaming-режиме на живых данных и в batch-режиме на исторических.

Тестирование. Стриминг-логику тестируют детерминированно: подают фиксированную последовательность событий с явными таймстампами и водяными знаками и сверяют вывод. Во Flink для этого есть MiniClusterWithClientResource и тест-харнессы операторов; в Kafka Streams — TopologyTestDriver, который выполняет топологию без брокера. Отдельно тестируются сценарии «событие опоздало», «событие вне lateness», «дубль», «рестарт из savepoint».

Стоимость. Основные статьи: диск и сеть Kafka (репликация утраивает трафик), память под состояние, S3-запросы за чекпоинты (частые чекпоинты маленького джоба могут стоить дороже, чем сам джоб). compression.type=zstd на топиках даёт 3–5x экономии почти бесплатно.

Мини-итог

  • Лог — фундаментальная абстракция: упорядоченный, воспроизводимый, разделяемый источник истины. Порядок гарантируется в партиции, значит ключ решает всё.
  • Считайте по event time; водяной знак — явно настраиваемая эвристика «сколько ждём опоздавших», и это компромисс между свежестью и полнотой.
  • Тип окна выбирает бизнес-вопрос. Sliding умножает состояние на размер/шаг, session требует слияния, любое окно требует стратегии для опоздавших.
  • Состояние — главный ресурс стриминга. RocksDB + инкрементальные чекпоинты + обязательный TTL.
  • «Exactly-once» существует как семантика обработки, а не доставки. За пределами Kafka она достигается связкой at-least-once + идемпотентный приёмник.
  • Здоровье потоковой системы измеряется лагом, длительностью чекпоинтов и количеством отброшенных опоздавших — не загрузкой CPU.

Источники

Что дальше

Мы научились считать на лету. Но и потоковые джобы, и батчевые расчёты нужно кем-то запускать, перезапускать, связывать зависимостями и переигрывать за прошлые даты. Об этом — следующая статья: Оркестрация пайплайнов: Airflow, dbt, идемпотентность и backfill.

Нашли неточность? Выделите фрагмент текста — рядом появится жучок.

Нужен разбор именно вашей ситуации?

Статья описывает общий случай. Если у вас частный — можно разобрать его отдельно, платно. А если не хватает целого материала, предложите тему: её оплачивают вскладчину, и она выходит открытой для всех.

Доска запросов