В производных данных мы договорились: писать надо в источник правды, а кеши, индексы и витрины — выводить из него единым потоком изменений. Эта статья — про то, как этот поток устроен на практике.
Потоковая обработка — это, по сути, пакетная обработка, которая никогда не заканчивается: те же преобразования, но не над готовым конечным набором данных, а над бесконечным потоком событий, и с низкой задержкой. Разберём две вещи, на которых спотыкаются чаще всего: как именно изменения из базы попадают в производные системы (CDC) и почему аналитика по потоку начинает врать, если перепутать время события и время обработки.
CDC: база как источник потока изменений
Как заставить поисковый индекс всегда соответствовать базе и при этом не скатиться в двойную запись (два отдельных write, которые молча расходятся)? Ответ — перехват изменений данных (change data capture, CDC). Идея: назначить базу ведущим узлом, читать её журнал изменений и применять эти изменения к производным системам ровно в том же порядке, в каком они легли в базу.
Как это работает. У базы уже есть журнал упреждающей записи (WAL) — тот самый, по которому идёт репликация. CDC-инструмент подключается к этому журналу как ещё одна «реплика» и раскодирует поток изменений. Например, Debezium умеет читать logical decoding в PostgreSQL, бинарный лог MySQL, oplog MongoDB — и публикует все изменения в лог-брокер вроде Kafka. Дальше поисковый индекс, кеш, аналитическое хранилище — это просто потребители одного и того же потока; и раз поток сохраняет порядок изменений, все они сходятся к состоянию базы.
Две детали, без которых картина CDC неполна:
- Начальный снимок. Журнал хранит не всю историю (старые куски удаляются), а новому индексу нужны и давно не менявшиеся записи тоже. Поэтому старт устроен так: сначала снимок всей базы на известной позиции журнала, а затем поток изменений начиная с этой позиции.
- Уплотнение журнала (log compaction). Чтобы не хранить журнал вечно, брокер выбрасывает старые записи по ключам, которые уже перезаписаны, оставляя по каждому ключу только последнее значение. Тогда журнал превращается в полную копию текущего состояния — и производную систему можно пересобрать с нуля из него одного.
CDC асинхронен: производные системы отстают от базы на задержку репликации, как обычные реплики. Зато источник правды почти не чувствует медленных потребителей — они его не тормозят.
Время события против времени обработки
Вторая ловушка тоньше и коварнее. Потоковый процессор часто считает что-нибудь «за последние 5 минут». И тут внезапно оказывается, что «время» — это не одно понятие, а два:
- время события (event-time) — когда действие реально произошло;
- время обработки (processing-time) — когда процессор до этого события наконец добрался.
Между ними есть разрыв: события задерживаются в очередях, при сетевых сбоях, а особенно — при перезапуске потребителя (который потом залпом доедает всё накопившееся за простой). И вот тут группировка по времени обработки рисует фантом. Процессор на минуту завис, после перезапуска разом обработал накопленные события — и на графике «частота запросов» вырастает всплеск, хотя реальная частота ни на секунду не менялась. Считать надо по времени события.
Аналогия. «Звёздные войны» вышли не в порядке эпизодов (сначала IV, V, VI, потом I, II, III). Если смотреть в порядке выхода, порядок сюжета нарушен. Дата события — это год по сюжету, а дата обработки — когда ты фильм посмотрел. В потоке порядок событий сбивается точно так же, и алгоритм обязан это учитывать.
Отдельная боль — отставшие события (stragglers). Окно 37-й минуты вы вроде уже закрыли и посчитали, а тут прилетает событие с меткой 37:59, которое застряло в сети. Что с ним делать? Вариантов два: либо игнорировать опоздавших (и тогда обязательно мониторить, какую долю вы отбрасываете), либо публиковать поправку к уже выданному окну. «Правильного» ответа нет — это осознанный выбор.
Типы окон
Раз мы считаем по времени события, надо решить, как нарезать время на окна:
- Падающее (tumbling) — окна фиксированной длины, и каждое событие попадает ровно в одно: минутные окна 10:03:00–10:03:59, 10:04:00–10:04:59, и так далее. Самый простой вариант.
- Прыгающее (hopping) — окна тоже фиксированной длины, но перекрываются для сглаживания: например, 5-минутное окно с шагом в 1 минуту.
- Скользящее (sliding) — берёт все события, оказавшиеся в пределах заданной ширины друг от друга; фиксированных границ у него нет.
- Сессионное (session) — без фиксированной длины вообще: группирует события одного пользователя, идущие близко во времени, и закрывается, когда пользователь замолкает (скажем, 30 минут тишины). Классика веб-аналитики.
Где это применяется
Как только производную систему (индекс, кеш, витрину) надо держать свежей почти в реальном времени или считать метрики по потоку событий — вы в потоковой обработке. Практическая рамка: синхронизируйте производные системы через CDC и лог-брокер, а не двойной записью; любую аналитику по времени считайте по времени события, а не обработки; а для окон честно решите заранее, что делаете с опоздавшими.
Где спотыкаются начинающие:
- Держат индекс/кеш в синхроне двойной записью — и они молча расходятся. CDC делает базу ведущим, а производные — потребителями её журнала.
- Считают частоту или среднее по времени обработки — и при отставании процессора получают фантомные всплески, которых в реальности не было.
- Считают окно закрытым по настенным часам — и теряют отставшие события. Нужна заранее выбранная политика: дропать (и мониторить долю) или публиковать поправку.
- Путают event-time и processing-time при перезапуске — и переигровка накопленного потока рисует аномалию на ровном месте.
Что почитать дальше: производные данные — зачем вообще нужен единый поток изменений; AMQP vs Kafka — лог-брокер, на котором держатся CDC и повторное проигрывание; event sourcing — поток событий как источник правды; двойная запись — антипаттерн, который CDC заменяет.