Сквозные гарантии и жизнь пайплайна

Каждое звено обещает своё, а отвечать приходится за цепочку целиком: как at-least-once складывается в effectively-once, где заканчивается транзакция Kafka, почему дедуп на стоке обязателен и как пайплайн переживает replay, эволюцию схемы и битые события

«У нас Kafka с exactly-once» — фраза, после которой стоит спросить: exactly-once где именно? Транзакционный producer Kafka честно даёт атомарность внутри себя — но пайплайн редко заканчивается на Kafka. У него в начале PostgreSQL, в середине Kafka Streams, в конце ClickHouse, и на каждом шве гарантия может незаметно ослабнуть с exactly-once до at-least-once. Дальше — либо дубли в витрине, либо осознанная идемпотентность на приёмнике. Эта статья — про то, где именно рвётся цепочка гарантий и чем это закрывают на практике: не абстрактно, а на живых числах одного стенда, где дубли, replay, эволюция схемы и битые события — не мысленный эксперимент, а то, что реально произошло и было измерено.

Это третья, заключительная статья серии «Стриминг и обработка данных» — она сшивает первую (CDC-пайплайн PostgreSQL → Debezium → Kafka → ClickHouse) и вторую (Kafka Streams без кластера) статьи: там уже был весь путь события от строки в PostgreSQL до строки в витрине, здесь — вопрос, что гарантировано на этом пути, а что — нет.

Ретрофутуристическая схема в стиле «Полдень. XXI век»: цепочка звеньев PostgreSQL → Debezium → Kafka → Kafka Streams → ClickHouse изображена как последовательность шлюзов на приборной панели, у каждого шлюза — своя лампа-гарантия (at-least-once, EOS, идемпотентность), последний шов перед ClickHouse отмечен разрывом пунктирной линии с подписью «здесь транзакция Kafka заканчивается»

В статье

Композиция гарантий: что обещает каждое звено

У пайплайна CDC → Kafka → Kafka Streams → ClickHouse нет одной сквозной гарантии — есть цепочка локальных, и итоговая надёжность равна самому слабому звену, а не самому сильному.

  • PostgreSQL — источник истины. Что записано в таблице orders, то и есть факт; всё остальное в пайплайне — производная от него, и в этой статье PG-значения многократно используются именно как эталон для сверки.
  • Debezium / Kafka Connect — читает журнал изменений и публикует в Kafka. Гарантия здесь — at-least-once: коннектор может повторно доставить событие после рестарта, но не потеряет его молча.
  • Kafka как транспорт — держит сообщение с гарантиями конкретного брокера: репликация, durability, поведение при подтверждении записи. Это отдельная тема, разобранная в «Брокеры и стриминг: транзакции в очереди и логе» — здесь она не повторяется.
  • Kafka Streams — обрабатывает поток и умеет exactly-once semantics (EOS) для операции внутри Kafka, через транзакционный producer. Механика этой транзакции — предмет отдельной статьи «Exactly-once и транзакции в Kafka»; здесь важна не механика, а граница действия — см. следующий раздел.
  • Go-консьюмер → ClickHouse — читает результат Kafka Streams и пишет в витрину. Это уже не Kafka: гарантия по умолчанию — at-least-once (коммит offset после успешной обработки), и именно здесь дубли становятся физически видимыми, если ничего не предпринять.

Это принципиально другая картина, чем у одного брокера с собственной сквозной гарантией «из коробки» — например, JetStream в NATS, где надёжность и производительность настраиваются внутри одной системы (см. «Надёжность и производительность NATS»). В композиции из нескольких независимых систем сквозной гарантии никто не выдаёт по умолчанию — её приходится собирать самому на каждом шве.

Где кончается транзакция Kafka

Транзакция Kafka Streams в режиме EOS покрывает шаблон «прочитал → обработал → записал» внутри Kafka: чтение из входного топика, запись в выходной топик и коммит offset входа — атомарны как единое целое. Если Streams упадёт посреди обработки, при рестарте либо весь эффект применится, либо не применится ни один его кусочек — половинчатого состояния внутри Kafka не бывает.

Но у этой транзакции есть чёткая граница, и она — на выходе из Kafka. Запись в ClickHouse, которую делает Go-консьюмер, не является частью этой транзакции: это отдельный шаг, отдельный клиент, отдельная сеть. К моменту, когда данные покидают Kafka, гарантия EOS уже отработала своё — дальше начинается обычный at-least-once: консьюмер может прочитать сообщение из customer.totals, записать его в ClickHouse и упасть до коммита offset — при рестарте то же сообщение придёт снова.

Граница EOS не заканчивается ровно там, где кажется на первый взгляд — есть ещё один шов, уже внутри чтения из Kafka, а не после него. customer.totals — транзакционный вывод Kafka Streams (PROCESSING_GUARANTEE_CONFIG=EXACTLY_ONCE_V2): запись физически попадает в топик ещё до того, как брокер решит, коммитить транзакцию или прервать (abort) её. Консьюмер с дефолтной изоляцией read_uncommitted (дефолт клиентской библиотеки franz-go) отдаёт код приложения именно физически записанное — в том числе содержимое транзакции, которая ещё не закоммичена или в итоге окажется прерванной. Если Kafka Streams убить (kill -9) посреди активной EOS-транзакции, read_uncommitted-консьюмер рискует прочитать и протащить в витрину данные, которых с точки зрения Kafka никогда не было. kill -9 этого процесса на стенде реально происходит (см. дальше разбор TaskCorruptedException в разделе про лаг), но там это проявилось иначе — как рассинхронизация чекпойнта RocksDB, другой отказ, чем прерванная транзакция продюсера (что именно измерялось в этом прогоне, а что нет, — оговорка ниже). Поэтому Go-консьюмер сконфигурирован на kgo.FetchIsolationLevel(kgo.ReadCommitted()): видимость записи только после того, как брокер зафиксировал commit-маркер транзакции — это и есть контрактная часть EOS для потребителей, а не опция для тюнинга.

Честно про эту правку и границы того, что доказано: прогон, из которого взяты цифры этой статьи, выполнен только с read_committed. Abort транзакции отдельно не форсировался, и разница между двумя уровнями изоляции в этом прогоне не измерялась — по итоговым числам нельзя даже сказать, были ли вообще скрытые записи прерванных транзакций (под read_committed они как раз и невидимы). Поэтому read_committed здесь — не «было хуже, стало лучше», а обязательная часть контракта EOS для потребителя транзакционного топика: настройка закрывает класс риска по построению, а не потому, что мы поймали конкретный abort. Форсировать реальный abort и сравнить оба режима — отдельная, более рискованная для стенда демонстрация, и в эту правку она не входит.

Отсюда и тезис, вокруг которого строится вся статья: exactly-once end-to-end не бывает подарком сервера — это идемпотентность плюс дедуп на стоке. Транзакционность Kafka снимает проблему дублей внутри Kafka; проблему дублей после Kafka снимает не Kafka, а то, как устроена запись в ClickHouse. Дальше — как это выглядит на цифрах.

Семантика чтения changelog: почему наивная сумма врёт

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

customer_totals/raw_totals заполняются из customer.totals — а это не поток независимых событий, а changelog накопительных снимков состояния клиента: каждое сообщение значит «клиент X теперь на orders=N, total=M», и по мере поступления новых заказов один и тот же клиент даёт в топике всё новые и новые снимки, каждый раз чуть больше предыдущего (в этом и состоит внутренняя механика state store Kafka Streams — она за пределами этой статьи, см. «Kafka Streams без кластера»). Если просто сложить total по всем строкам такой таблицы, результат окажется завышен уже в baseline, без единого дубля — потому что складываются не N независимых чисел, а ~N растущих частичных сумм одного и того же клиента.

Это проверено на стенде напрямую: два прогона одного и того же Go-консьюмера над одним и тем же топиком customer.totals, единственная переменная — флаг -dup (искусственно дублирует каждое сообщение перед вставкой).

Baseline, без -dup, полное перечитывание топика с начала, ClickHouse только что пересоздан пустым:

консьюмер: вставлено=51 битых=0 dlq_ошибок=0 ошибок_вставки=0 ошибок_fetch=0 (dup=false, from-start=true, mode=totals)

raw_totals count()                          51
customer_totals count() (no FINAL)          51
customer_totals FINAL count()/sum(total)    5   2831.46

PG-истина в этот момент (SELECT count(*), sum(amount) FROM orders): 51|2831.46 — совпадает с customer_totals FINAL до цента.

Тот же топик, тот же сброс offset на earliest, единственное отличие — флаг -dup:

консьюмер: вставлено=102 битых=0 dlq_ошибок=0 ошибок_вставки=0 ошибок_fetch=0 (dup=true, from-start=true, mode=totals)

raw_totals count()                          153
customer_totals count() (no FINAL)          153
customer_totals FINAL count()/sum(total)    5   2831.46

raw_totals выросла с 51 до 153, то есть на 102 строки — но это не «эффект флага -dup» в чистом виде. Baseline-прогон уже перечитывает весь топик с начала и вставляет 51 строку сам по себе (та же таблица, продолженная); прогон с -dup делает то же самое перечитывание, но дублирует каждое сообщение, добавляя ещё 102. Дальше — то, ради чего этот прогон вообще ставился:

customer_totals FINAL: sum=2831.46                                       -- = PG, при 153 физических строках
raw_totals наивная sum(total) (без дедупа): 49779.99                     -- врёт в 17.58 раза
raw_totals argMax(total, version) по customer_id (то, что делает FINAL): 2831.46   -- = PG

Наивная sum(total) по raw_totals — 49779.99, что расходится с PG-истиной (2831.46) в 17.58 раза. Ключевой момент, который в предыдущей редакции статьи был упущен: эта дельта не измеряет эффект дублей. Она была бы такой же огромной и при -dup=false, если бы топик перечитали дважды подряд без единого искусственного дубля — потому что naive-sum по changelog снимков растёт от самого факта повторного чтения снимков, а не от факта их дублирования. То, что она действительно измеряет, — это то, ЧТО значит правильно прочитать changelog накопительных состояний: не «сложить все строки», а «взять последнее состояние каждого клиента» — argMax(total, version), та же операция, которую под капотом делает модификатор FINAL, и она даёт 2831.46, byte-в-byte совпадение с PG.

А инвариантность самой customer_totals FINAL (2831.46 что при 51 строке, что при 153) — это тоже не то доказательство дедупа, которым может показаться на первый взгляд: снимок состояния клиента идемпотентен по своей природе (повторная доставка того же снимка — это просто повтор той же информации, а не новая), поэтому устойчивость FINAL здесь показывает устойчивость чтения именно ЭТОГО changelog к дублям снимков, а не изолированный эффект дублей на независимых данных. Для чистого доказательства нужен другой эксперимент — он ниже.

Изолированный дедуп-эксперимент: 2N против N

Чтобы измерить эффект дублей отдельно от семантики changelog, стенд читает другой источник — orders.events, топик Debezium-outbox событий (один заказ — одно сообщение), а не снимков состояния:

{"amount":59.68,"order_id":1,"customer_id":2}
{"amount":80.33,"order_id":5,"customer_id":1}
{"amount":87.63,"order_id":7,"customer_id":3}

Сток — две новые таблицы ClickHouse: order_counts_raw (SummingMergeTree, cnt=1 на каждое событие, без дедупа) и order_dedup (ReplacingMergeTree(version) по order_id, FINAL оставляет одну строку на уникальный заказ). Заполняет их новый режим консьюмера -mode dedup-demo — тот же клиент, та же обвязка read_committed/DLQ/at-least-once, другой вход и сток.

Baseline — первое прочтение orders.events с начала:

консьюмер: вставлено=51 битых=0 dlq_ошибок=0 ошибок_вставки=0 ошибок_fetch=0 (dup=false, from-start=true, mode=dedup-demo)

order_counts_raw sum(cnt)                    51
order_counts_raw count() (физические строки) 51
order_dedup FINAL count()                    51
order_dedup count() (физические строки, no FINAL)   51
PG count(*) FROM orders                      51

N=51, совпадает с PG. Дублей ещё не было — обе таблицы физически совпадают с N. Replay — тот же -from-start, топик перечитан целиком ещё раз, без единого изменения в источнике (чистая имитация повторной доставки/backfill):

консьюмер: вставлено=51 битых=0 dlq_ошибок=0 ошибок_вставки=0 ошибок_fetch=0 (dup=false, from-start=true, mode=dedup-demo)

order_counts_raw sum(cnt)                    102
order_counts_raw count() (физические строки) 102
order_dedup FINAL count()                    51
order_dedup count() (физические строки, no FINAL)   102

Это и есть чистое доказательство эффекта дубля, без искажения семантикой changelog. order_counts_raw+1 на каждое событие, без дедупа — sum(cnt) буквально удвоилась: 51 → 102, 2N. Каждый из 51 реального заказа физически посчитан дважды, и это видно напрямую: SummingMergeTree консолидирует строки одного customer_id фоновым merge, но не меняет общую сумму — sum(cnt) инвариантна к тому, прошёл фоновый merge или нет. order_dedup — дедуп по order_id — дал FINAL count()=51, N, несмотря на 102 физические строки: повторная доставка того же order_id гасится при merge/FINAL.

2N против N на одном и том же входе, единственная переменная — повторное чтение того же топика событий, — это и измеряет тезис «дедуп на стоке обязателен»: не «наивная сумма врёт в 17.58 раза» (это про changelog-семантику, см. раздел выше), а «без дедупа число реальных сущностей удвоилось, с дедупом — нет». Вместе с инвариантностью FINAL при полном replay в разделе про backfill ниже это и есть две ноги доказательства, а не одна слабая.

Идемпотентное состояние на стоке и ReplacingMergeTree

Слово «идемпотентная» здесь стоит уточнить сразу: сама вставка не идемпотентна — повторная доставка физически добавляет строки в таблицу (это и показал изолированный эксперимент выше). Идемпотентны итоговое состояние и чтение через FINAL/argMax: сколько бы повторов ни легло физически, последнее состояние каждого ключа читается одинаковым. Именно это свойство и обеспечивает ReplacingMergeTree.

Консьюмер льёт каждое событие сразу в две таблицы: raw_totals (обычный MergeTree, хранит всё физически, включая дубли) и customer_totals (ReplacingMergeTree, движок с версионным ключом, предназначенный гасить повторные версии одной и той же строки). В момент замера §3 (153 строки) обе таблицы физически одинаковы — консьюмер пишет в них одно и то же сообщение. Но это состояние конкретного момента, а не постоянное свойство: ReplacingMergeTree гасит дубли не на вставке, а асинхронно, фоновым merge (с оговоркой про настройку optimize_on_insert — ниже). Пока merge не отработал, в самой таблице лежат все физические строки, включая повторы — читать её нужно всегда через FINAL или явный argMax(total, version), а не полагаться на голый count().

Насколько это не постоянное свойство, видно на самой таблице: к моменту, когда весь прогон FIXTURES отработал (после §7), raw_totals содержит 306 строк (все физические вставки), а customer_totals — уже 7 (смёрджено фоном). «Физически одинаковы» здесь больше не выполняется, паритет количества строк без FINAL был случайностью момента замера, а не гарантией.

К raw_totals модификатор FINAL вообще неприменим — движок MergeTree его не поддерживает, и это не абстракция, а буквальная ошибка сервера:

$ clickhouse-client -q "SELECT count() FROM raw_totals FINAL"
Received exception from server (version 26.6.1):
Code: 181. DB::Exception: Received from localhost:9000. DB::Exception: Storage MergeTree doesn't support FINAL. (ILLEGAL_FINAL)
(query: SELECT count() FROM raw_totals FINAL)

Значит ли это, что MergeTree не мёрджится вовсе, в отличие от ReplacingMergeTree? Нет — мёрджится точно так же, просто не гасит строки. Это подтверждено контролируемым экспериментом из system.part_log: обе таблицы получили одни и те же пять кусков-источников и слились в один и тот же по имени результирующий кусок, в одну и ту же секунду:

event_time            table              part_name    rows  event_type   merged_from
2026-07-17 18:37:57   customer_totals    all_1_5_1    5     MergeParts   [all_1_1_0..all_5_5_0]
2026-07-17 18:37:58   raw_totals         all_1_5_1    304   MergeParts   [all_1_1_0..all_5_5_0]

Единственная переменная между этими двумя строками — движок таблицы. ReplacingMergeTree (customer_totals) на выходе дал кусок из 5 строк — дубли схлопнуты прямо во время merge. MergeTree (raw_totals) на тех же пяти кусках, в ту же секунду, дал кусок из 304 строк — ни одна не погашена, merge физически объединил данные, но количество не тронул. Вывод: разницу даёт не «мёрджится / не мёрджится», а гасит ли конкретный движок повторные версии строки во время merge. (Число 304 здесь — состояние уже позже по прогону, после досева и рестарта Kafka Streams из раздела про лаг ниже, а не 153 строки, на которых стоит доказательство выше; таблица продолжает расти на протяжении всего прогона.)

Отдельная оговорка про демо-настройки стенда: вставка в ClickHouse на стенде идёт с явным optimize_on_insert=0. На дефолтной конфигурации ClickHouse (optimize_on_insert=1) ReplacingMergeTree может погасить дубли уже на вставке — но только если обе копии строки попали в один и тот же батч. Стенд отключает эту настройку намеренно, потому что флаг -dup кладёт обе копии сообщения в один батч, и с дефолтной настройкой демонстрация дедупа через merge (а не через вставку) была бы не видна. Реальная пересдача at-least-once в проде почти всегда кросс-батчевая — повтор приходит отдельным сообщением спустя время, в другом батче — и в этом случае гашение на вставке не спасает независимо от optimize_on_insert: дедуп на чтении через FINAL/argMax нужен всегда, а не только при выключенной оптимизации.

Replay безопасен именно из-за дедупа

Дедуп на стоке — это не страховка на редкий случай, а то, что делает возможным осмысленный replay, и это вторая нога доказательства из раздела про изолированный дедуп-эксперимент выше. Проверено последовательно, продолжая состояние после customer.totals-прогонов и изолированного дедуп-эксперимента (raw_totals=153, customer_totals FINAL=5/2831.46, PG=51 заказов): досеяно 20 новых заказов.

Инкрементальный подхват (без -from-start, продолжение с committed offset):

до:   raw_totals=153, FINAL=5/2831.46
$ go run . -for 15s
консьюмер: вставлено=20 битых=0 dlq_ошибок=0 ошибок_вставки=0 ошибок_fetch=0 (dup=false, from-start=false, mode=totals)
после: raw_totals=173, FINAL=5/3757.82

Подхвачены ровно новые 20 записей — витрина выросла на сумму этих заказов. PG в этот момент: 71|3757.82 — совпадает с FINAL до цента. Дальше — сброс committed offset на earliest и полное перечитывание топика с начала, то есть replay:

до:   raw_totals=173, FINAL=5/3757.82
$ go run . -from-start -for 15s
консьюмер: вставлено=71 битых=0 dlq_ошибок=0 ошибок_вставки=0 ошибок_fetch=0 (dup=false, from-start=true, mode=totals)
после: raw_totals=244, FINAL=5/3757.82

raw_totals физически выросла ещё на 71 строку — весь топик (51 из baseline + 20 новых) перечитан заново и вставлен второй раз (173+71=244). А вот customer_totals FINAL не изменилась: 5/3757.82 — байт-в-байт та же сумма, что и до replay, несмотря на то что физически витрина только что получила ещё 71 повторную вставку.

Это и есть первая нога доказательства из раздела про изолированный дедуп-эксперимент: replay безопасен не потому, что Kafka что-то дедуплицировала (она не занимается этим), и не потому, что консьюмер стал умнее — а потому, что на чтении витрина всегда проходит через FINAL/argMax. Без дедупа на стоке та же операция replay удвоила бы сумму в отчёте; с ним — она не сдвинулась ни на копейку. Вместе с 2N против N на orders.events выше это и закрывает тезис-якорь статьи двумя независимыми проверками, а не одной.

Эволюция схемы сквозь звенья

Третья проверка сквозной цепочки — что произойдёт, если в событие добавить поле, которого раньше не было, без специальной подготовки на приёмной стороне. На стенде отправлено событие с новым полем currency: order_id=132 customer_id=3 amount=42.00 currency=RUB.

PG-истина независимо от Kafka: у customer_id=3 стало orders=27 total=1112.08. Kafka Streams, посчитав агрегат независимо и до ClickHouse, дал тот же результат:

{"customer_id":3,"orders":27,"total":1112.08}
Streams независимо посчитал orders=27 total=1112.08 — совпадает с PG (27/1112.08) до цента.

Новое поле не сломало и не исказило агрегацию Kafka Streams. При этом currency физически не доезжает до customer_totals — не потому что кто-то его специально фильтрует, а потому что TotalsProcessor пересобирает JSON заново из трёх полей (customer_id/orders/total), и лишнее поле обрезается уже на этом шаге. Это факт конкретной топологии, а не общая устойчивость соседнего звена к неизвестным полям.

У самого Go-консьюмера здесь важно различить два разных свойства, которые легко спутать. Первое — устойчивость к незнакомым полям: консьюмер разбирает JSON обычным encoding/json.Unmarshal без DisallowUnknownFields(), поэтому поле вроде currency, дойди оно до консьюмера, было бы молча проигнорировано, а не уронило бы парсинг. Второе — валидация значений: то же сообщение с семантически недопустимым customer_id (вне диапазона реальных клиентов 1..5) или без обязательных полей (orders=0) консьюмер отвергает в DLQ — см. следующий раздел. Незнакомое поле терпится, недопустимое значение отбрасывается; это два независимых, оба намеренных поведения, и путать их не стоит.

Сквозной результат эволюции схемы виден по витрине: после того как событие с currency прошло всю цепочку, FINAL-сумма реальных клиентов сходится с PostgreSQL до цента.

CH customer_totals FINAL (реальные клиенты): count=5 sum=6534.23
PG (истина):                                 count=132 sum=6534.23
CH customer_id=3 (после эволюции):           orders=27 total=1112.08

Добавление поля не исказило агрегат: агрегация Kafka Streams осталась верной, а витрина реальных клиентов сошлась с источником до цента. Само поле currency при этом в выходной контракт не входит и намеренно отбрасывается — TotalsProcessor пересобирает JSON заново из трёх полей, так что до customer.totals оно не доезжает вовсе (это ожидаемое поведение топологии, а не потеря).

DLQ: битое событие не роняет пайплайн

Отдельно от «нового, но валидного» поля — случай действительно битых сообщений, и стенд проверяет два разных вида «битости».

Первый — синтаксически битое, не-JSON. На стенде в топик customer.totals инъецировано некорректное JSON-тело, и следом запущен консьюмер:

$ bash scripts/inject-bad.sh
битое событие отправлено в customer.totals

$ go run . -for 15s
битое событие offset=140: invalid character 'o' in literal null (expecting 'u')
консьюмер: вставлено=0 битых=1 dlq_ошибок=0 ошибок_вставки=0 ошибок_fetch=0 (dup=false, from-start=false, mode=totals)

Второй — семантически пустое: {} — синтаксически валидный JSON, encoding/json.Unmarshal разберёт его без единой ошибки и получит customer_id=0, orders=0, total=0 — которое без отдельной проверки тихо уехало бы в витрину как мусор. Консьюмер трактует customer_id=0 как битую запись (диапазон реальных клиентов стенда — 1..5, см. scripts/seed.sh) тем же путём, что и синтаксически невалидный JSON:

$ echo '{}' | docker exec -i ds-kafka /opt/kafka/bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic customer.totals

$ go run . -for 15s
битое событие offset=141: семантически невалидная запись: customer_id=0 (диапазон клиентов стенда — 1..5, см. scripts/seed.sh)
консьюмер: вставлено=0 битых=1 dlq_ошибок=0 ошибок_вставки=0 ошибок_fetch=0 (dup=false, from-start=false, mode=totals)

Проверка не сводится к одному customer_id, и здесь важна тонкость контракта данных. encoding/json.Unmarshal при разборе в обычное числовое поле-значение не отличает отсутствующее поле от явного нуля — {"customer_id":1} разберётся в orders=0, total=0 без единой ошибки, и без проверок такая запись уехала бы в витрину ложным агрегатом. (Отличить можно, но для этого поле должно быть указателем — к этому вернёмся ниже.) Реальный снимок customer.totals всегда имеет orders>=1, поэтому orders==0 отвергается:

$ echo '{"customer_id":1}' | docker exec -i ds-kafka /opt/kafka/bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic customer.totals

$ go run . -for 15s
битое событие offset=146: семантически невалидная запись: orders=0 у customer_id=1 (реальный агрегат customer.totals всегда orders>=1, см. TotalsProcessor.process())
консьюмер: вставлено=0 битых=1 dlq_ошибок=0 ошибок_вставки=0 ошибок_fetch=0 (dup=false, from-start=false, mode=totals)

Но orders>=1 не покрывает total: сообщение {"customer_id":1,"orders":1} прошло бы проверку orders, а total в нём отсутствует и превратился бы в 0. Отличить отсутствующий total от допустимого явного total:0 в Go можно только разбором в структуру с указателем (*float64nil означает «поля не было»); отдельная проверка отвергает отсутствие:

$ echo '{"customer_id":1,"orders":1}' | docker exec -i ds-kafka /opt/kafka/bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic customer.totals

$ go run . -for 15s
битое событие offset=147: семантически невалидная запись: total отсутствует у customer_id=1; обязательное поле реального агрегата (см. TotalsProcessor.process())
консьюмер: вставлено=0 битых=1 dlq_ошибок=0 ошибок_вставки=0 ошибок_fetch=0 (dup=false, from-start=false, mode=totals)

Во всех случаях пайплайн не упал: битых=1, ошибок_вставки=0, dlq_ошибок=0, процесс завершился штатно. Ни одно из сообщений не потерялось молча — все попали в customer.totals.dlq целиком, с контекстом:

dlq-reason:invalid character 'o' in literal null (expecting 'u')
dlq-origin-topic:customer.totals
dlq-origin-offset:140
тело: not-a-json{{{

dlq-reason:семантически невалидная запись: customer_id=0 (диапазон клиентов стенда — 1..5, см. scripts/seed.sh)
dlq-origin-topic:customer.totals
dlq-origin-offset:141
тело: {}

dlq-reason:семантически невалидная запись: orders=0 у customer_id=1 (реальный агрегат customer.totals всегда orders>=1, см. TotalsProcessor.process())
dlq-origin-topic:customer.totals
dlq-origin-offset:146
тело: {"customer_id":1}

dlq-reason:семантически невалидная запись: total отсутствует у customer_id=1; обязательное поле реального агрегата (см. TotalsProcessor.process())
dlq-origin-topic:customer.totals
dlq-origin-offset:147
тело: {"customer_id":1,"orders":1}

Здесь показаны четыре сценария, каждый из которых воспроизводится командой выше. Физический снимок топика на стенде содержит на одну запись больше — повторный {}, оффсет происхождения которого достоверно установить не удалось; он приведён в FIXTURES.md стенда как есть, без попытки объяснить задним числом.

Причина (dlq-reason), исходный топик и оффсет — в заголовках DLQ-записи, тело — как есть, без попытки его починить или отбросить. Это минимальный контракт DLQ на этом звене: пайплайн продолжает работать при следующем сообщении, а не блокируется на одном плохом, и ничего не теряется без следа — разобрать и переиграть битое событие можно по данным самой DLQ-записи.

Честная оговорка про границы этого контракта: DLQ покрывает сток — Go-консьюмер между customer.totals и ClickHouse. Битый вход на уровне Kafka Streams (невалидный JSON, нечисловой или отсутствующий customer_id/amount в orders.events) до DLQ не доезжает вовсе: топология обработки отбрасывает такое событие молча, с записью в stderr процесса Kafka Streams, а не в customer.totals.dlq. Если бы статья обещала «пайплайн целиком ничего не теряет», это было бы неточно — гарантия DLQ здесь узкая: не теряется то, что дошло до стока и не прошло его собственную проверку, а не всё, что вообще вошло в пайплайн.

Лаг как метрика здоровья

Лаг консьюмер-группы — стандартная метрика здоровья пайплайна, но здесь есть две ловушки, обе видны на живых числах стенда.

Первая — не всякое ненулевое число значит отставание. В покое (Kafka Streams и Go-консьюмер простаивают, событий не поступает):

=== группа ds-streams ===
ds-streams  ds-streams-customer-totals-repartition  0  33  33  0
ds-streams  ds-streams-customer-totals-repartition  1  17  17  0
ds-streams  ds-streams-customer-totals-repartition  2  32  32  0
ds-streams  orders.events                           0  33  33  0
ds-streams  orders.events                           1  25  25  0
ds-streams  orders.events                           2  13  13  0
=== группа ds-sink ===
ds-sink   customer.totals  0  76  77  1

ds-streams держит LAG=0 на всех партициях — здесь всё однозначно. А вот у ds-sink LAG=1 в этом конкретном замере, и это не обязательно отставание: это хвостовой control-маркер транзакционного продюсера EXACTLY_ONCE_V2, который Go-консьюмер физически не может «дочитать» — это не сообщение с данными. Важная оговорка: LAG=1 здесь — не постоянный инвариант, а наблюдение конкретного прогона, и зависит от того, была ли последняя запись в топике именно control-маркером. В другой момент, если последней физической записью оказались данные, а не маркер, ds-sink на этом же топике в покое мог бы показать и LAG=0. Заводить алерт «LAG > 0 у ds-sink — есть проблема» без учёта этой природы значит регулярно получать ложные срабатывания на пустом месте — но и полагаться на «LAG=1 всегда безопасен» тоже нельзя: это не гарантированное свойство, которое можно зашить в порог алерта, а факт, который надо смотреть на конкретном прогоне.

Вторая ловушка — растущий лаг нужно уметь отличать от «не успевает» и не путать с фактом, что звено попросту не запущено. Kafka Streams остановлен kill -9 без предварительной очистки локального state dir, orders.events продолжает наполняться:

$ bash scripts/seed.sh 30   # Streams остановлен
ds-streams  orders.events  0  33  46  13
ds-streams  orders.events  1  25  33  8
ds-streams  orders.events  2  13  22  9
                                       (итого 30)

$ bash scripts/seed.sh 30   # ещё 30, Streams всё ещё остановлен
ds-streams  orders.events  0  33  54  21
ds-streams  orders.events  1  25  41  16
ds-streams  orders.events  2  13  36  23
                                       (итого 60)

Лаг вырос с 30 до 60 — ровно на число заказов, добавленных за то время, пока звено остановлено. Рестарт после такой нечистой остановки (kill -9 без предварительной очистки локального состояния) сам по себе не тихий: локальный чекпойнт RocksDB разошёлся с состоянием брокера, и в логе появляется TaskCorruptedException. Это не сбой, требующий вмешательства оператора, а штатное самообнаружение: Kafka Streams переинициализировал таски, восстановил состояние из changelog-топика и нагнал весь бэклог сам — lag.sh после рестарта снова показал ds-streams LAG=0 на всех партициях, а ds-sink тем временем вырос до LAG=64 (Kafka-уровневый бэклог: Streams уже произвёл агрегаты в customer.totals, но они ещё не слиты в ClickHouse). Разовый прогон Go-консьюмера сразу после (вставлено=60) слил этот бэклог в витрину без разрывов — raw_totals выросла до 304 строк, customer_totals FINAL5/6492.23, совпадает с PG (131|6492.23) до цента, а ds-sink вернулся к LAG=1 (тот же control-маркер, см. оговорку выше).

Здесь стоит явно оговорить то, что легко упустить при чтении одной только метрики Kafka-лага: LAG=0 у консьюмер-группы ds-streams говорит про уровень Kafka — что все сообщения из orders.events вычитаны и обработаны Kafka Streams. Это не то же самое, что «весь бэклог доехал до ClickHouse»: агрегаты, которые Kafka Streams нагнал в customer.totals, всё ещё нужно забрать оттуда отдельным прогоном Go-консьюмера (тем же вызовом, что и инкрементальный подхват в разделе про replay) — это отдельная консьюмер-группа со своим лагом, что и подтвердил разрыв между LAG=64 и последующим сливом выше. Проверять здоровье пайплайна по лагу одного-единственного звена недостаточно: нужен лаг на каждой границе между звеньями, а не только на входе.

Карта выбора

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

  • Композиция независимых систем (CDC → Kafka → Streams → аналитическое хранилище) — оправдана, когда источник, транспорт и потребитель по природе разные системы: PostgreSQL как OLTP, Kafka как журнал, ClickHouse как аналитика. Цена — необходимость проектировать идемпотентность на каждом шве самостоятельно, как показано выше.
  • Один брокер со своей сквозной гарантией (например, JetStream в NATS, см. «Надёжность и производительность NATS») — проще эксплуатационно, если вся коммуникация умещается в одну систему и не требует отдельного аналитического хранилища; надёжность и производительность настраиваются внутри одной модели, а не собираются из нескольких.
  • Транзакции брокера (RabbitMQ vs Kafka) — если вопрос именно в гарантиях самого транспортного слоя, а не композиции нескольких систем, разбор в «Брокеры и стриминг: транзакции в очереди и логе».
  • Дедуп на чтении против дедупа на записиReplacingMergeTree/FINAL (эта статья) не единственный способ; для агрегатов, посчитанных сразу в момент вставки, есть материализованные представления ClickHouse — см. «Материализованные представления и real-time агрегации в ClickHouse».

Общая карта серий и статей о брокерах и обработке потоков — в хабе «Карта messaging: очереди, брокеры, стриминг».

Демо и версии

Все числа и логи в статье — из живых прогонов одного стенда, единым проходом после полной пересборки (docker compose down -v && up -d), без ничего досочинённого.

digital-cookbook/messaging/data-streaming
Компонент Версия
PostgreSQL 18.4
Kafka 4.3.1
Kafka Connect (framework) 4.3.0
Debezium connector 3.6.0.Final
Kafka Streams 4.3.1
Jackson (databind) 2.18.2
ClickHouse 26.6.1.1193
Go 1.26.3
franz-go / kadm v1.18.0 / v1.14.0
clickhouse-go/v2 v2.30.0
shopspring/decimal v1.4.0

Версии зафиксированы на момент прогона — поведение ReplacingMergeTree/FINAL, EXACTLY_ONCE_V2 в Kafka Streams и обработки DLQ чувствительно к минорным изменениям, при обновлении стоит перепроверять.

Документация

Обсуждение в Telegram

Присоединиться →

Комментарии