Batch и streams: время, состояние, replay и сверка
Сравниваем полный и incremental расчёт, определяем late-policy, проверяем joins, retention и backpressure в бесплатной локальной модели.
Отчёт Клуба показывает регистрации по минутам. Телефон участника час был offline, затем отправил старое событие. Оперативный график уже закрыл ту минуту, а ночной пересчёт добавил регистрацию. Значения различаются, хотя оба обработчика правильно сложили числа. Не определён договор о времени и полноте.
Batch вычисляет ограниченный набор входов с выбранной границей. Stream обрабатывает поступающие изменения без заранее известного конца. Сначала выберем входные факты и инвариант: каждому уникальному ID соответствует один факт; одинаковый диапазон принятых фактов и одна версия алгоритма должны давать одинаковый агрегат. Внешние уведомления не входят в пересчёт аналитики.
Полный результат и производное состояние
Полный расчёт можно выразить как GROUP BY уникальных событий. Incremental обработчик хранит state и изменяет его по каждой записи. Повтор доставки без дедупликации удвоит сумму. Изменение записи также нельзя всегда считать новым фактом: CDC UPDATE может требовать вычесть старое и прибавить новое. Tombstone должен удалить вклад, а не просто исчезнуть из входа.
Stream-table duality полезна как точная связь: журнал изменений позволяет восстановить таблицу по префиксу, а таблица представляет состояние после применения. Она требует правил ключей, версий, deletes и удержания истории. Снимок в произвольный момент плюс хвост журнала без согласованной границы может пропустить или дважды применить изменение.
Materialized view быстрее отвечает на частый вопрос, но добавляет источник отставания и ошибок. У поискового индекса, отчёта и кеша могут быть разные позиции. Сверка должна сравнивать одинаковый префикс или watermark/snapshot boundary; сравнение активно меняющихся источника и копии без границы создаёт ложные расхождения.
Event time и закрытие окна
Рассмотрим события минут 0,2,1, повтор 0,5, затем новое 0. Event time приходит от источника, ingestion — от приёма, processing — от работы worker. Если час offline меняет минуту отчёта на текущую, это другой отчёт. Идентификатор записи и её timestamp тоже не взаимозаменяемы.
Для локального опыта окно имеет длину 1 минуту. Определим учебную границу: окно закрыто, когда window_end + grace <= max_seen_event_time. Grace равен 3. При поступлении минуты 5 окно минуты 0 закрыто: 1+3 <= 5. Новое событие минуты 0 считается поздним. Повтор прежнего ID подавляется до этого решения. Это явная модель; Kafka Streams использует свой data-driven stream time, а другие движки свои watermarks и policies. Kafka Streams: Time.
Watermark описывает ожидаемое продвижение в рамках правил источника и обработки. Он не запрещает прислать любое старое событие. Молчащая partition может остановить общий минимум; исключение idle источника позволяет продвигаться, но требует разбора его возвращения. Late-policy выбирает reject, correction или reopen/recompute. Точность, задержка публикации и размер state связаны с этим выбором.
Если приходит 10 тысяч событий/с и окно держится 10 минут, это до 6 миллионов входных событий до дедупликации/агрегации. При среднем state record 100 байт грубая оценка — 600 MB без индексов и runtime overhead. Хранение только агрегата может быть дешевле, но correction/delete и joins потребуют дополнительных сведений. Не объявляйте память по размеру одного счётчика, если обещаете исправлять любой старый факт.
Join требует временного договора
Stream-stream join сопоставляет два потока. Покупка и подтверждение оплаты могут прийти в любом порядке. Нужно определить ключ, допустимый интервал между событиями, позднее подтверждение и срок удержания несопоставленной стороны. Немедленное отсутствие оплаты не доказывает окончательный отказ.
Stream-table join добавляет справочник к событию. Если событие 12:00 обрабатывается в 14:00, а цена менялась в 13:00, какая цена нужна? Latest-table lookup и as-of lookup дают разные ответы. Для воспроизводимого исторического результата надо хранить факт цены в событии или иметь версионный справочник с нужным временем. Replay поверх нынешнего справочника может изменить прошлое.
Repartition по ключу join добавляет сеть, хранение и возможность hot key. Запрос по одной популярной встрече не распределится от общего большого числа ключей. Join state очищается по договору времени; без границы медленная сторона создаёт неограниченный backlog.
Retention и восстановление
Consumer сохраняет позицию. Restart читает после неё в обычном пути, либо выбранный диапазон при replay. Если retention уже удалил необходимое начало, нужна согласованная исходная таблица и хвост. Log compaction, если используется, сохраняет определённую историю по ключу, а не автоматически полный audit; для event sourcing этого может быть недостаточно.
Recovery состояния требует согласованности результата и позиции. Позиция после записи 45 при незавершённой 44 создаёт пропуск. SQL consumer может хранить inbox+effect атомарно и только затем подтвердить позицию; Kafka-to-Kafka pipeline использует собственный транзакционный механизм. Обычный HTTP-эффект не попадает в него. Эти границы разобраны в Kafka Design и доставке.
Replay в новую версию результата отделяет вычисление от текущих читателей. Путь: зафиксировать входную границу → пересчитать новой логикой → догнать хвост → сверить → переключить чтение → сохранить безопасный откат. Если старая версия уже не получает новые изменения, откат к ней после cutover требует новой синхронизации. Такой же принцип работает для secondary index.
Backpressure и нагрузка
При producer 100/с и consumer 60/с backlog растёт на 40/с. Буфер ограничивает память и время до отказа. Приёмник выбирает pause, rejection или durable buffering. Возврат «принято» имеет смысл только после хранения по обещанному контракту. Retrying каждого rejected запроса немедленно делает перегрузку сильнее.
DLQ — отдельное состояние процесса, а не завершение работы. Нужны причина ошибки, исходный ID, схема, owner, срок и правило replay. Событие с плохой схемой может блокировать partition; пропуск требует явного решения о порядке и потерянном эффекте. Метрики показывают oldest pending age, lag, service rate, retry rate и размер state, отдельно по значимым группам.
Запускаем оригинальную модель
Скачайте бесплатный комплект лабораторий либо используйте tools/system-design-labs в репозитории. Нужны Python 3.9+ и стандартный sqlite3. Из каталога лабораторий:
python3 stream_lab.py
python3 stream_lab.py --testСначала предскажите effects, sum, next_offset для трёх вариантов crash. Затем предскажите окно 0 до и после позднего e. Только после этого сравните JSON. Модель создаёт новые временные базы и не принимает путь к рабочим данным.
Подсказка 1: ранний offset забывает ещё не выполненную работу. Подсказка 2: поздний offset повторяет уже сохранённую работу. Подсказка 3: inbox и эффект должны быть одной транзакцией, но код обработчика всё равно может выполниться больше одного раза.
Разбор: early после restart — два эффекта, сумма 5; late — четыре, сумма 7; safe — три, сумма 6 при четырёх доставках. Event-time результат окна 0 равен 2, full batch равен 13 из-за позднего уникального значения 11. Replay той же последовательности и политики совпадает со stream; replay с другой grace может дать иной результат законно. Один ID с разным содержанием отвергается как конфликт.
Во второй части producer предлагает пачку прежде service каждого секундного такта. На седьмом такте rejected_total равен 140, очередь после service — 140. Это дискретный порядок; непрерывная модель имеет другую траекторию. Измените consumer на 120 и затем на 0, объясните результат до запуска. Граница буфера и честный статус приёма должны сохраниться.
Критерии готовности: таблица crash windows, объяснение расхождения batch/stream, найденный предел retention, версия алгоритма/schema для replay, отдельный путь внешних уведомлений. Пересматривайте архитектуру при требованиях исторических joins, более длинном offline, праве удаления фактов или недопустимости неполного оперативного результата.
Первичные источники
Kafka Design, Kafka Streams: Core Concepts, MIT 6.5840: MapReduce и проверка отказов. Числа, late-policy и лаборатория — самостоятельные учебные конструкции.
Запишите ход рассуждений, расчёты и вопросы. Сохраните текст перед уходом со страницы. После входа в аккаунт ответ участвует в общей синхронизации прогресса. Автоматической оценки архитектуры здесь нет.