← назад к разделу

В производных данных мы договорились: писать надо в источник правды, а кеши, индексы и витрины — выводить из него единым потоком изменений. Эта статья — про то, как этот поток устроен на практике.

Потоковая обработка — это, по сути, пакетная обработка, которая никогда не заканчивается: те же преобразования, но не над конечным набором данных, а над бесконечным потоком событий, и с низкой задержкой. Разберём две вещи, на которых спотыкаются чаще всего: как именно изменения из базы попадают в производные системы (CDC) и почему аналитика по потоку начинает врать, если перепутать время события и время обработки.

Обязательно

CDC: база как источник потока изменений

Как заставить поисковый индекс всегда соответствовать базе и при этом не скатиться в двойную запись (два отдельных write — в базу и в индекс, — которые однажды молча разойдутся)? Ответ — перехват изменений данных (change data capture, CDC). Идея: назначить базу ведущим узлом, читать её журнал изменений и применять эти изменения к производным системам ровно в том же порядке, в каком они легли в базу.

Как это работает. У базы уже есть журнал упреждающей записи (WAL) — тот самый, по которому идёт репликация. CDC-инструмент подключается к этому журналу как ещё одна «реплика» и раскодирует поток изменений. Например, Debezium умеет читать logical decoding в PostgreSQL, бинарный лог MySQL, поток изменений (change stream) MongoDB — и публикует все изменения в лог-брокер вроде Kafka. Дальше поисковый индекс, кеш, аналитическое хранилище — это просто потребители одного и того же потока; и раз поток сохраняет порядок изменений, все они сходятся к состоянию базы.

база брокер Kafka журнал WAL, Debezium поисковый индекс переиндексирует кеш обновляет ключи витрина догружает

Пишут только в базу, а её журнал читает Debezium и кладёт изменения в брокер: поисковый индекс, кеш и витрина просто три потребителя одного потока, и каждый применяет изменения в том же порядке, в каком они легли в базу.

Две детали, без которых картина CDC неполна:

  • Начальный снимок. Журнал хранит не всю историю (старые куски удаляются), а новому индексу нужны и давно не менявшиеся записи тоже. Поэтому старт устроен так: сначала снимок всей базы на известной позиции журнала, а затем поток изменений начиная с этой позиции.
  • Уплотнение журнала (log compaction). Чтобы не хранить журнал вечно, брокер выбрасывает старые записи по ключам, которые уже перезаписаны, оставляя по каждому ключу только последнее значение. Тогда журнал превращается в полную копию текущего состояния — и производную систему можно пересобрать с нуля из него одного. Само собой это не включается: у сообщений должен быть ключ, на топике должен стоять режим уплотнения, а удаление записи приходится изображать пустым сообщением с тем же ключом — его называют надгробием (tombstone). Надгробия тоже со временем убирают, иначе они копились бы вечно.

CDC асинхронен: производные системы отстают от базы на задержку репликации, как обычные реплики. Зато источник правды почти не чувствует медленных потребителей — они его не тормозят. Не чувствует, но не бесплатно: в PostgreSQL потребитель держит слот репликации, и база не выбрасывает журнал, пока он не подтвердит прочитанное. Потребитель, который встал на сутки, спокойно забивает диск базы — за отставанием слотов следят отдельно.

журнал цена 100 цена 120 остаток 5 цена 180 остаток 0 после уплотнения цена 180 остаток 0

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

Почему Debezium отстаёт и что с этим делать

Частый сюжет: Debezium настроен на две таблицы, в них пишут редко, а изменения приезжают в Kafka с задержкой в минуты. Причина почти всегда в том, как устроено логическое декодирование на стороне PostgreSQL, и лечится она там, а не в настройках коннектора.

Раскодируется весь журнал, фильтр стоит после. Слот репликации обслуживает один процесс walsender: он читает WAL базы подряд и раскодирует каждое изменение каждой таблицы. Решение «этой таблицы нет в публикации» принимает плагин вывода pgoutput, а к нему изменение попадает уже раскодированным. Поэтому чужая горячая таблица — журнал аудита, счётчики — стоит слоту процессорного времени, хотя в поток не попадает; распараллелить слот нельзя, это один процесс на одном ядре. Первое, что смотрят при отставании, — кто в базе пишет больше всех.

Транзакция отдаётся целиком после фиксации. Плагин получает изменения транзакции только на её COMMIT; до этого они копятся в памяти слота (logical_decoding_work_mem, по умолчанию 64 МБ) и сбрасываются на диск. Пакетное обновление на десять миллионов строк сорок минут не даёт потребителю ничего, а после фиксации приезжает одним куском, и всё зафиксированное позже ждёт очереди. С 14-й версии PostgreSQL умеет отдавать незавершённую транзакцию частями (режим streaming протокола pgoutput), но Debezium получает транзакции после фиксации, так что лечится это со стороны записи: большие обновления режут на пакеты и фиксируют по частям. Сколько транзакций ушло на диск, показывает pg_stat_replication_slots (spill_txns, spill_bytes).

Отставание слота — не всегда отставание данных. Слот хранит confirmed_flush_lsn — позицию, до которой потребитель подтвердил чтение, и база не удаляет WAL старше неё:

живой пример

SELECT slot_name, active, wal_status,
       pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), confirmed_flush_lsn)) AS lag
FROM pg_replication_slots;
Запустить

Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →

Debezium подтверждает позицию, когда событие ушло в Kafka. Если отслеживаемые таблицы молчат, а соседние пишут гигабайты, подтверждать нечего: lag растёт вместе с диском под WAL, хотя ни одно нужное изменение не потеряно. Для этого есть heartbeat.interval.ms и heartbeat.action.query: коннектор сам раз в интервал пишет строку в служебную таблицу, входящую в публикацию, получает её обратно как событие и подтверждает позицию.

Снимок большой таблицы. Стартовый снимок читает таблицу целиком до начала потока, а слот уже создан, и всё записанное за это время копится в WAL; после снимка коннектор догоняет накопленное, и первые часы выглядят как отставание. Большие таблицы снимают инкрементально, через сигнальную таблицу: снимок идёт кусками (incremental.snapshot.chunk.size) вперемешку с потоком и хвоста не оставляет.

TOAST. Длинные значения — текст, JSON — PostgreSQL хранит вне строки, и при REPLICA IDENTITY DEFAULT неизменившаяся такая колонка в событии UPDATE не приходит вовсе: вместо неё заглушка __debezium_unavailable_value, и потребитель обязан помнить прошлое значение сам. REPLICA IDENTITY FULL чинит это ценой полной старой строки в WAL на каждое обновление — больше декодирования и больше отставания.

Самописная репликация через свой сервис — обычно не CDC, а таблица исходящих (Java, Go, Node, Python): приложение в той же транзакции пишет событие в outbox, отдельный поток отправляет их в брокер. Чем лучше: в потоке только нужные события и с бизнес-смыслом, слота и декодирования чужих таблиц нет, диск под WAL не растёт. Чем хуже: событие появляется только там, где его написали, — прямой UPDATE из миграции в поток не попадёт; порядок гарантирован в пределах агрегата, а не всей базы; таблицу нужно чистить, доставка идёт как минимум один раз, и потребитель обязан быть идемпотентным. CDC берут, когда нужен полный след всех изменений базы; outbox — когда нужны события приложения.

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

Вторая ловушка тоньше. Потоковый процессор часто считает что-нибудь «за последние 5 минут». И тут внезапно оказывается, что «время» — это не одно понятие, а два: когда действие реально произошло и когда процессор до этого события наконец добрался. Между ними есть разрыв: события задерживаются в очередях, при сетевых сбоях, а особенно — при перезапуске потребителя (который потом залпом доедает всё накопившееся за простой). И вот тут группировка по времени обработки рисует фантом. Процессор на минуту завис, после перезапуска разом обработал накопленные события — и на графике «частота запросов» вырастает всплеск, хотя реальная частота ни на секунду не менялась. Считать надо по времени события. Первую шкалу называют временем события (event-time), вторую временем обработки (processing-time), и окна считают по первой.

Разницу видно на маленькой программе — один и тот же поток, два счёта по минутам:

живой пример

import java.util.Map;
import java.util.TreeMap;

public class WindowDemo {
    record Event(int eventSec, int arrivedSec) {}

    public static void main(String[] args) {
        Event[] stream = {
                new Event(65, 66), new Event(85, 86),
                new Event(125, 185), new Event(145, 186),
                new Event(185, 187), new Event(205, 206)
        };
        Map<Integer, Integer> byEvent = new TreeMap<>();
        Map<Integer, Integer> byArrival = new TreeMap<>();
        for (Event e : stream) {
            byEvent.merge(e.eventSec() / 60, 1, Integer::sum);
            byArrival.merge(e.arrivedSec() / 60, 1, Integer::sum);
        }
        System.out.println("по времени события:   " + byEvent);
        System.out.println("по времени обработки: " + byArrival);
    }
}
Запустить

Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →

живой пример

package main

import (
	"fmt"
	"maps"
	"slices"
	"strings"
)

type event struct{ eventSec, arrivedSec int }

func show(counts map[int]int) string {
	var parts []string
	for _, minute := range slices.Sorted(maps.Keys(counts)) {
		parts = append(parts, fmt.Sprintf("%d=%d", minute, counts[minute]))
	}
	return "{" + strings.Join(parts, ", ") + "}"
}

func main() {
	stream := []event{{65, 66}, {85, 86}, {125, 185}, {145, 186}, {185, 187}, {205, 206}}
	byEvent, byArrival := map[int]int{}, map[int]int{}
	for _, e := range stream {
		byEvent[e.eventSec/60]++
		byArrival[e.arrivedSec/60]++
	}
	fmt.Println("по времени события:  ", show(byEvent))
	fmt.Println("по времени обработки:", show(byArrival))
}
Запустить

Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →

живой пример

const stream = [[65, 66], [85, 86], [125, 185], [145, 186], [185, 187], [205, 206]];
const byEvent = new Map(), byArrival = new Map();
for (const [eventSec, arrivedSec] of stream) {
  byEvent.set(Math.floor(eventSec / 60), (byEvent.get(Math.floor(eventSec / 60)) ?? 0) + 1);
  byArrival.set(Math.floor(arrivedSec / 60), (byArrival.get(Math.floor(arrivedSec / 60)) ?? 0) + 1);
}
const show = (m) => '{' + [...m.keys()].sort((a, b) => a - b).map((k) => `${k}=${m.get(k)}`).join(', ') + '}';
console.log('по времени события:   ' + show(byEvent));
console.log('по времени обработки: ' + show(byArrival));
Запустить

Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →

живой пример

from collections import Counter

stream = [(65, 66), (85, 86), (125, 185), (145, 186), (185, 187), (205, 206)]
by_event = Counter(event_sec // 60 for event_sec, _ in stream)
by_arrival = Counter(arrived_sec // 60 for _, arrived_sec in stream)


def show(counts: Counter) -> str:
    return "{" + ", ".join(f"{k}={counts[k]}" for k in sorted(counts)) + "}"


print("по времени события:   " + show(by_event))
print("по времени обработки: " + show(by_arrival))
Запустить

Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →

В первом ряду частота ровная — по два события в минуту; во втором вторая минута пустая, а третья даёт всплеск из четырёх.

Аналогия: «Звёздные войны» вышли не в порядке эпизодов — сначала IV, V, VI, потом I, II, III. Дата события здесь — год по сюжету, дата обработки — когда фильм посмотрели. В потоке порядок сбивается так же.

Отдельная боль — отставшие события (stragglers). Окно 37-й минуты вы вроде уже закрыли и посчитали, а тут прилетает событие с меткой 37:59, которое застряло в сети. Что с ним делать? Вариантов два: либо игнорировать опоздавших (и тогда обязательно следить, какую долю вы отбрасываете), либо публиковать поправку к уже выданному окну. «Правильного» ответа нет — это осознанный выбор.

окно выбирается по времени события, а не по моменту прихода поток: события приходят вперемешку 10:01пришло 10:01 10:03пришло 10:03 10:02пришло 10:04 10:07пришло 10:07 10:04пришло 10:12 окно 10:00–10:04 окно 10:05–10:09 событий: событий: 0 1 2 3 4 0 1 набор идёт закрыто, итог 3 поправка: было 3, стало 4 набор идёт

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

Типы окон

Раз мы считаем по времени события, надо решить, как нарезать время на окна:

  • Падающее (tumbling) — окна фиксированной длины, и каждое событие попадает ровно в одно: минутные окна 10:03:00–10:03:59, 10:04:00–10:04:59, и так далее. Самый простой вариант.
  • Прыгающее (hopping) — окна тоже фиксированной длины, но перекрываются для сглаживания: например, 5-минутное окно с шагом в 1 минуту.
  • Скользящее (sliding) — берёт все события, оказавшиеся в пределах заданной ширины друг от друга; фиксированных границ у него нет. Тут будьте внимательны с названиями: в Flink классом SlidingEventTimeWindows названо как раз прыгающее окно из пункта выше, а не это. Смотрите на определение в документации, а не на слово.
  • Сессионное (session) — без фиксированной длины вообще: группирует события одного пользователя, идущие близко во времени, и закрывается, когда пользователь замолкает (скажем, 30 минут тишины). Классика веб-аналитики.

Соединения в потоке

Подсчёты по окну — только половина работы. Вторая половина — соединить поток с чем-то ещё, и это задача с двумя разными ответами.

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

Ключевая деталь, о которой забывают: справочник тоже имеет версии во времени. Заказ, оформленный в марте, должен считаться по цене марта, а не по сегодняшней, — поэтому справочник хранят с интервалами действия и соединяют по времени события, а не по текущему состоянию.

Чего делать нельзя: ходить в базу за справочником на каждое событие. На потоке в тысячи событий в секунду это тысячи запросов, и обработчик становится медленнее источника.

Поток плюс поток. Соединить «заказ создан» и «оплата прошла» сложнее: события приходят в разное время и в любом порядке. Такое соединение обязано быть по окну: «оплата в пределах получаса после заказа». Обработчик держит в состоянии заказы, ждущие оплату, и оплаты, пришедшие раньше заказа, а по истечении окна выбрасывает несовпавшие (и обычно отправляет их отдельно — это и есть «заказы без оплаты»).

Отсюда два практических следствия. Размер окна определяет объём состояния: окно в сутки на потоке в миллион событий — это миллион записей в памяти обработчика. И несовпавшие события — не ошибка, а результат: их надо куда-то направить, иначе они просто исчезнут.

Когда поток не нужен

Самая частая ошибка в этой теме — строить поток там, где хватает расписания. Поток дороже пакетной обработки по всем статьям: состояние обработчика, порядок событий, повторы, догон после сбоя, дежурство.

Пакетный пересчёт по расписанию даёт взамен три вещи, которых у потока нет: воспроизводимость (одинаковый вход даёт одинаковый выход, и это можно проверить), простота перезапуска (упало — запустили заново) и отсутствие состояния между запусками. Поэтому правило звучит так: если задержка в час или в сутки устраивает — считайте пакетом.

Признаки, по которым поток действительно нужен: свежесть измеряется секундами и это требование продукта (антифрод, уведомления, остатки); объём данных такой, что полный пересчёт не укладывается в окно; или события в принципе не хранятся так, чтобы их можно было пересчитать.

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

Начальный снимок и стык с потоком

Про снимок сказано, что он нужен, — но вся сложность в стыке: как сделать так, чтобы между снимком и потоком не потерялось и не удвоилось.

Порядок такой. Сначала фиксируется позиция в журнале базы — точка, с которой поток начнёт читать. Затем снимается снимок данных, согласованный с этой позицией (в PostgreSQL это делается в одной транзакции с созданием слота репликации, в MySQL — с фиксацией позиции журнала). Снимок грузится в приёмник. Потом обработчик включается с зафиксированной позиции и применяет всё, что случилось после неё.

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

На больших таблицах снимок делают по частям (по диапазонам ключа), и тогда поток может идти параллельно снимку: современные инструменты захвата изменений умеют такой режим, отслеживая, какие диапазоны уже перенесены. Проверяют результат так же, как любой переезд: сверкой счётчиков и контрольных сумм по диапазонам.

Где это применяется

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

Где спотыкаются начинающие:

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

Глубже: водяные знаки, смещения и повторная обработкарасширенное

Политика для опоздавших выше отвечает, что делать с событием, пришедшим после закрытия окна. Осталось два вопроса, которые определяют всё поведение потоковой системы: когда именно окно закрывается и что будет после перезапуска.

Водяной знак (watermark) это утверждение обработчика «все события со временем раньше T я уже видел». Считается он по данным: обычно максимальное время события среди увиденных минус допустимое опоздание, скажем, пять минут. Окно 10:00–10:04 закрывается, когда водяной знак переходит через 10:05, то есть когда пришло событие с временем 10:10 и старше. Пока поток течёт, знаки движутся; если источник замолчал, знак стоит, и окно не закроется никогда. Поэтому тихим источникам задают простой: не было событий минуту, считаем, что источник опустел, и двигаем знак по часам. У Kafka Streams и Flink это настройка задержки водяного знака, а не что-то, что нужно писать самим, но выбирать её значение придётся: маленькое опоздание закрывает окна быстро и теряет больше событий, большое честнее, но результат появляется позже.

Смещения (offsets). Потребитель Kafka запоминает, докуда дочитал, и делает это после обработки, иначе при падении между чтением и записью результат потеряется. «После обработки» означает, что при падении между записью результата и подтверждением смещения событие обработается дважды. Это доставка хотя бы один раз, и она лежит в основе всего: результат должен переживать повтор. Либо запись идемпотентна (обновление по ключу, INSERT ... ON CONFLICT), либо смещение хранят вместе с результатом в одной транзакции той базы, куда пишут, либо используют транзакции Kafka, когда результат тоже уходит в Kafka: там запись выхода и подтверждение входа фиксируются вместе, и это называют «ровно один раз», хотя относится оно только к цепочке Kafka-в-Kafka. Flink делает то же через контрольные точки состояния и двухфазную фиксацию во внешние хранилища.

Повторная обработка. Главное преимущество потока, о котором забывают: журнал можно прочитать заново. Нашли ошибку в расчёте, поправили код, сдвинули смещения группы потребителей на неделю назад (kafka-consumer-groups --reset-offsets --to-datetime), и производные данные пересчитались. Для этого журнал хранят достаточно долго (срок хранения топика это осознанное решение, а не умолчание в неделю), результат пишут так, чтобы повтор перезаписывал, а не дописывал, и состояние обработчика (окна, счётчики) при повторе сбрасывают, иначе старые окна сложатся с новыми. Повтор с водяными знаками тоже работает: при перечитывании недели знаки движутся по временам событий, и окна закрываются в том же порядке, что и в первый раз.

Коротко

  • CDC делает базу ведущей: производные читают её журнал и применяют изменения в том же порядке — вместо двойной записи, которая расходится.
  • Старт потребителя — снимок базы на позиции журнала и поток изменений с неё; уплотнённый журнал равен текущему состоянию.
  • CDC асинхронен: производные отстают на задержку репликации, зато медленный потребитель не тормозит источник правды.
  • Аналитику по времени считают по времени события: по времени обработки отставший процессор рисует всплеск, которого не было.
  • Окна бывают падающие, прыгающие, скользящие и сессионные; для опоздавших нужна политика — отбросить или выпустить поправку.
  • Водяной знак закрывает окна («всё раньше T уже видел»), смещения подтверждают после обработки, поэтому повтор нормален и результат обязан его переживать; журнал перечитывают заново после исправления, сбрасывая состояние.
  • Соединение со справочником делают локально (таблица рядом с обработчиком) и по времени события, а не запросом в базу на каждое событие.
  • Соединение двух потоков обязано иметь окно: его размер задаёт объём состояния, а несовпавшие события — результат, который надо куда-то направить.
  • Если задержка в час устраивает, пакетный пересчёт дешевле и воспроизводимее потока; промежуточный вариант — частый инкрементальный пакет.
  • Позицию журнала фиксируют до снимка, а не после: перекрытие лечится идемпотентностью записи, потеря — ничем.

Что почитать дальше

  • Производные данные — зачем вообще нужен единый поток изменений.
  • AMQP и Kafka — лог-брокер, на котором держатся CDC и повторное проигрывание.
  • Event sourcing — поток событий как источник правды.
  • Таблица исходящих (Java, Go, Node, Python) — второй способ не скатиться в двойную запись, когда CDC не подходит.