Потоковая обработка: 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
Общая схема потоковой платформы
(events)"] A2["OLTP-база
CDC / Debezium"] A3["Клиенты
web / mobile"] end subgraph bus["Шина: Kafka"] T1["raw.*
сырые события"] T2["clean.*
валидные, со схемой"] DLQ["dlq.*
битые сообщения"] end subgraph proc["Обработка"] F1["Flink job:
парсинг + валидация"] F2["Flink job:
окна, join, агрегаты"] end subgraph sink["Приёмники"] S1["OLAP
ClickHouse / Pinot"] S2["Lakehouse
Iceberg / Delta"] S3["KV-сервинг
Redis / Cassandra"] S4["Обратно в Kafka
для других команд"] end A1 --> T1 A2 --> T1 A3 --> T1 T1 --> F1 F1 -->|валидные| T2 F1 -->|ошибки схемы| DLQ T2 --> F2 F2 --> S1 F2 --> S2 F2 --> S3 F2 --> S4 SR["Schema Registry
Avro / Protobuf"] -.контракт.-> F1 SR -.контракт.-> A1
Три вещи, которые видно на схеме и которые обычно забывают в первой версии пайплайна: отдельный слой сырых данных (без него нельзя переиграть историю после исправления парсера), 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».
Жизненный цикл окна
событие ключа Создано --> Накапливает: добавляем в агрегат Накапливает --> Накапливает: новое событие
внутри границ Накапливает --> Слито: session-окна
merge при перекрытии Слито --> Накапливает Накапливает --> Сработало: watermark >
конец окна Сработало --> Уточняется: опоздавшее событие
в пределах allowed lateness Уточняется --> Сработало: повторный emit
(update / retract) Сработало --> Удалено: конец окна +
allowed lateness < watermark Уточняется --> Удалено Удалено --> [*] Накапливает --> Отброшено: событие старше
watermark и вне lateness Отброшено --> [*]: side output / метрика
Обратите внимание на переход «Отброшено». В проде этот путь обязан быть измеряемым: счётчик 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%.
Семантика доставки: что означает каждый уровень
коммитятся атомарно K-->>D: read_committed видит только
закоммиченные записи
Разберём точно, что где происходит.
At-most-once. Оффсет коммитится до обработки; при падении сообщение теряется. Легально ровно в одном случае — когда данные заведомо избыточны (сэмплированная телеметрия), а потеря дешевле сложности.
At-least-once. Оффсет коммитится после успешной обработки. Дубли возможны при любом падении между записью в приёмник и коммитом оффсета. Это дефолтный режим 90% продакшн-пайплайнов, и он совершенно нормален — при условии, что приёмник идемпотентен.
Exactly-once. Здесь начинается терминологическая ложь. Никакая распределённая система не может гарантировать, что сообщение будет доставлено ровно один раз — это доказанный факт (проблема двух генералов). Что реально гарантируется — exactly-once processing semantics: эффект на состояние и на выход выглядит так, как будто каждое сообщение обработано один раз. Достигается это двумя механизмами:
- Идемпотентный producer (
enable.idempotence=true): продюсеру выдаётся PID, каждой записи — sequence number на партицию; брокер отбрасывает дубли ретраев. Защищает от дублей на участке producer → брокер. Включён по умолчанию начиная с Kafka 3.0. - Транзакции (
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
Состояние и чекпоинты: как Flink переживает падения
Оконные агрегаты, дедупликация, join’ы, паттерны CEP — всё это состояние, которое живёт внутри работающего джоба и может достигать десятков и сотен гигабайт. Задача движка: пережить падение любой машины, не потеряв состояние и не посчитав ничего дважды.
Flink решает её алгоритмом асинхронных барьеров (asynchronous barrier snapshotting) — вариантом классического Chandy–Lamport. Координатор просит источники вставить в поток специальную запись — барьер. Барьер течёт вместе с данными; оператор, получивший барьер по всем входам, сохраняет своё состояние и пропускает барьер дальше.
Практические следствия, которые определяют настройку джоба:
- 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 превращает молчаливую потерю в наблюдаемый поток, который можно посчитать и по которому можно завести алерт.
Flink SQL: то же самое, но декларативно
Большая часть продовых стриминг-задач сегодня пишется не на 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.
Типичные ошибки
- Ключ партиционирования выбран без учёта порядка. События одной сущности разлетаются по партициям — и «создан» приходит после «отменён». Ключ должен совпадать с сущностью, для которой важен порядок.
- Горячий ключ. 40% трафика — один
merchant_id; одна партиция и один воркер перегружены, остальные простаивают. Лечится двухфазной агрегацией с солью ключа (key + random(0..N)на первом этапе, свёртка на втором). - Состояние без TTL. Дедупликация, join’ы, session-окна растут вечно. Джоб живёт три недели и умирает по диску.
- Processing time вместо event time. Пайплайн становится недетерминированным; результаты нельзя ни воспроизвести, ни сверить с батчем.
- Приёмник не идемпотентен. At-least-once + обычный
INSERT= дубли после каждого рестарта. Плюс каждое обновление окна по allowed lateness создаёт лишнюю строку. - Слишком много партиций. Тысячи партиций на топик увеличивают время ребалансировки, нагрузку на контроллер и объём метаданных. Начинайте с числа, кратного планируемому параллелизму, а не с «на вырост в 100 раз».
- Нет DLQ. Одно сообщение с битым JSON останавливает потребителя, тот падает в бесконечный рестарт (poison pill), лаг растёт, дежурный ночью удаляет топик.
- Ребалансировка на каждом деплое.
max.poll.interval.msменьше времени обработки батча → группа считает потребителя мёртвым → stop-the-world ребаланс. Используйте cooperative-sticky-assignor и static membership (group.instance.id). - Схема без реестра. Продюсер добавил обязательное поле — все потребители легли. Контракт должен быть машинно-проверяемым, с политикой совместимости
BACKWARD. - Мониторинг только по 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.
Источники
- Tyler Akidau et al. «Streaming Systems» (O’Reilly, 2018) — лучшая книга по event time, окнам и водяным знакам: https://www.oreilly.com/library/view/streaming-systems/9781491983867/
- Tyler Akidau, «The Dataflow Model» (VLDB 2015) — исходная статья, задавшая современную модель: https://research.google/pubs/pub43864/
- Martin Kleppmann, «Designing Data-Intensive Applications», гл. 11 «Stream Processing»: https://dataintensive.net/
- Jay Kreps, «The Log»: https://engineering.linkedin.com/distributed-systems/log-what-every-software-engineer-should-know-about-real-time-datas-unifying
- Carbone et al., «Lightweight Asynchronous Snapshots for Distributed Dataflows» (arXiv, 2015): https://arxiv.org/abs/1506.08603
- Chandy & Lamport, «Distributed Snapshots: Determining Global States of Distributed Systems»: https://lamport.azurewebsites.net/pubs/chandy.pdf
- Официальная документация Apache Flink (event time, state, checkpointing): https://nightlies.apache.org/flink/flink-docs-stable/
- Документация Apache Kafka — дизайн, транзакции, конфигурация: https://kafka.apache.org/documentation/
- Confluent, «Transactions in Apache Kafka»: https://www.confluent.io/blog/transactions-apache-kafka/
- Jay Kreps, «Questioning the Lambda Architecture» (о каппа-архитектуре): https://www.oreilly.com/radar/questioning-the-lambda-architecture/
Что дальше
Мы научились считать на лету. Но и потоковые джобы, и батчевые расчёты нужно кем-то запускать, перезапускать, связывать зависимостями и переигрывать за прошлые даты. Об этом — следующая статья: Оркестрация пайплайнов: Airflow, dbt, идемпотентность и backfill.