Вступление
Меня зовут Василий Миловидов, я системный архитектор в крупной компании. Мы строим обмен между сервисами на событиях: сервис меняет агрегат в PostgreSQL, в той же транзакции кладёт задачу в outbox, воркер на River, публикует событие в Kafka с ключом aggregate_id.
Разбирая соседнюю задачу, я обнаружил, что эта схема не даёт гарантии, которую мы в неё заложили. Строгого порядка событий одного агрегата в ней нет, и Kafka тут ни при чём: порядок теряется раньше, на своих же воркерах.
Статья написана для моей команды, чтобы у всех была одна картина проблемы и один язык для её обсуждения. Если она пригодится кому‑то ещё, тем лучше. Отдельно я старался, чтобы текст был понятен человеку, который только начинает работать с распределёнными системами: каждый механизм я объясняю с точки зрения того, что он делает и от чего защищает.
Всё, что касается поведения River, Kafka и PostgreSQL, проверено по документации, ссылки собраны в конце. Там, где утверждение является моим выводом, а не документированным поведением, я это отмечаю отдельно.
Словарь
Дальше по тексту эти слова используются в одном и том же значении. Если термины знакомы, раздел можно пропустить.
Термин | Что означает в этой статье |
|---|---|
Агрегат | Объект, изменения которого имеет смысл рассматривать как одну связную историю: заказ, посылка, счёт. У агрегата есть идентификатор, по которому его события собираются вместе. |
Событие | Запись о свершившемся факте: заказ оплачен, посылка собрана. После публикации событие не меняется. |
Версия агрегата | Порядковый номер события в потоке событий одного агрегата. Единственный атрибут, по которому определяется бизнес‑порядок. |
Outbox | Таблица исходящих сообщений в той же базе, что и бизнес‑данные. Позволяет одной транзакцией записать изменение и намерение отправить сообщение. |
Задача и воркер | Задача это строка в таблице river_job, описывающая единицу работы. Воркер это код, который её выполняет; воркеров несколько и работают они параллельно. |
Продюсер и потребитель | Продюсер пишет записи в Kafka, потребитель их читает. В нашей схеме продюсером выступает воркер River. |
Партиция и ключ | Топик Kafka разбит на партиции. Порядок записей гарантируется внутри одной партиции, а записи с одинаковым ключом попадают в одну и ту же партицию. |
Per‑key ordering | Строгий порядок внутри одного ключа при полной независимости разных ключей. В нашем случае ключ это идентификатор агрегата. |
Идемпотентность | Свойство обработчика, при котором повторная обработка того же самого сообщения не меняет результат. |
Доставка не менее одного раза (at‑least‑once) | Гарантия, при которой сообщение точно дойдёт, но может прийти повторно. Дубли лечатся идемпотентностью потребителя. |
Отравленное событие (poison event) | Событие, которое не удаётся опубликовать никогда: битый payload, отвергнутая схема, слишком большой размер. При строгом порядке оно блокирует все последующие события своего агрегата. |
Коротко
Проблема формулируется в четырёх пунктах.
Transactional outbox решает задачу атомарности: событие не потеряется, если транзакция закоммитилась. Задачу порядка он не решает вообще.
Ключ партиционирования Kafka сохраняет порядок поступления записей в партицию. Порядок поступления задаёт продюсер, то есть наши воркеры, работающие параллельно.
До Kafka порядок теряется в трёх независимых местах: при назначении позиции события, при параллельной обработке задач и внутри продюсера. После Kafka его добивает потребитель, раздающий записи в пул горутин.
Гарантию порядка нельзя купить одной настройкой. Её собирают из явной версии события, механизма сериализации по ключу и политики на случай события, которое не удаётся опубликовать никогда.
River Pro закрывает часть этой работы механизмом Sequences, но не всю, и об этом отдельный раздел.
Как устроен путь события
Схема, о которой идёт речь, выглядит так.

Воркер здесь выступает продюсером, то есть тем, кто пишет в Kafka; сервисы, которые читают, называются потребителями. River здесь выбран не случайно. Он хранит очередь задач в том же PostgreSQL, где лежат бизнес‑данные, поэтому вставку задачи можно сделать частью той же транзакции:
BEGIN; UPDATE orders SET status = 'paid' WHERE id = 123; INSERT INTO river_job (...) VALUES (...); COMMIT;
Это решает проблему двойной записи. Если бы мы меняли базу и отдельно писали в Kafka, между двумя операциями существовал бы момент, в который процесс может упасть: данные изменены, событие не отправлено, и никто об этом не знает. Внутри одной транзакции такого момента нет. Либо коммит прошёл и задача существует, либо был откат и задачи нет. River описывает это как одно из основных свойств своей архитектуры и специально предоставляет InsertTx для вставки задачи в чужой транзакции.
Дальше задачу забирает воркер, читает событие и публикует его в Kafka. Атомарность и порядок это разные свойства, и первое не влечёт второе.
Что мы называем порядком
Прежде чем что‑то гарантировать, полезно определить, что именно. Для одного агрегата существует как минимум четыре разных порядка, и они не обязаны совпадать.
Порядок | Чем задаётся | Наблюдаем ли он |
|---|---|---|
Порядок начала бизнес‑операций | временем прихода запросов в сервис | нет, зависит от балансировки и таймингов |
Порядок изменения данных в транзакции | моментом выполнения UPDATE | нет, транзакция ещё не видна другим |
Порядок коммитов транзакций | моментом COMMIT | да, это момент, когда факт стал существовать для системы |
Порядок записей в партиции Kafka | моментом приёма записи брокером | да, это то, что увидит потребитель |
Требование, которое имеет смысл, звучит так: порядок записей в партиции Kafka должен совпадать с порядком коммитов транзакций, изменивших агрегат.
Именно коммит, а не начало операции. Запрос A мог прийти раньше запроса B, но если транзакция B закоммитилась первой, то для системы факт B произошёл раньше. Восстанавливать «настоящий» порядок по времени начала операции бессмысленно: до коммита изменения не существует ни для кого, кроме своей транзакции.
Требование при этом должно быть локальным для агрегата. Разные заказы обязаны обрабатываться параллельно, иначе система не масштабируется. В литературе это называется per‑key ordering: строгий порядок внутри ключа, полная независимость между ключами.

Стрелки существуют только внутри рамок. Между рамками порядка нет и не требуется.
Кому именно нужен строгий порядок
Этот вопрос стоит задать до того, как выбирать механизм, потому что строгий порядок стоит дорого, а нужен он не всем потребителям. Потребители делятся на три группы, и цена гарантии для них разная.
Тип потребителя | Что ему нужно | Что ломается при нарушении порядка |
|---|---|---|
Восстанавливает текущее состояние (проекция, витрина, кэш) | последнее состояние | старое состояние затирает новое и остаётся навсегда |
Реагирует на переходы (уведомления, интеграции) | каждый факт, порядок важен | письмо «заказ собран» уходит раньше письма «заказ оплачен» |
Считает агрегаты и метрики | каждый факт, порядок не важен | ничего, если повторную обработку одного события он переживает без последствий |
Первой группе строгий порядок публикации не нужен. Ей достаточно версии в сообщении и правила «применяй, только если версия больше применённой». Событие, пришедшее с опозданием, просто отбрасывается: состояние всё равно сойдётся к правильному, потому что более новое событие содержит более новое состояние.

Это правило «побеждает последняя запись по версии». Оно дешевле строгой упорядоченной публикации: продюсер не сериализуется, воркеры остаются параллельными, отравленное событие, то есть событие, которое не удаётся опубликовать никогда, не блокирует поток. Платим мы тем, что промежуточные состояния могут не доехать до потребителя, а сообщение должно нести достаточно данных, чтобы применяться независимо от предыдущих.
Второй группе нужен именно порядок, и вся дальнейшая часть статьи про неё.
Практический вывод простой: прежде чем строить сериализацию, надо назвать потребителя, который ломается от нарушения порядка. Если такого потребителя нет, задача решается версией в сообщении и проверкой на стороне потребителя, и на этом можно остановиться.
Почему ключа Kafka недостаточно
Логика, из‑за которой возникает ложное чувство безопасности, выглядит убедительно. Kafka распределяет записи по партициям по хешу ключа, все события агрегата 123 идут в одну партицию, а внутри партиции порядок записей сохраняется. Значит, порядок есть.

Пропущен один шаг. Kafka сохраняет тот порядок, в котором записи в неё поступили. Поступают они от продюсера, а продюсер это наш воркер. Kafka не знает ни про транзакции PostgreSQL, ни про то, какое событие было раньше по бизнесу. Если два воркера отправили события в обратном порядке, Kafka честно запишет их в обратном порядке и ничего не нарушит.
Здесь же стоит зафиксировать два ограничения, о которых часто вспоминают поздно.
Изменение числа партиций у топика меняет отображение ключа в партицию. События одного агрегата, отправленные до и после расширения топика, могут оказаться в разных партициях, и порядок между ними исчезнет. Планировать число партиций надо с запасом на рост.
Порядок гарантируется внутри одной партиции, то есть внутри одного топика. Если события одного агрегата разложены по разным топикам, между ними порядка нет ни при какой конфигурации. Это важно для гранулярной схемы событий, где на каждый тип факта заводится свой топик.
Где именно ломается порядок
Первое: позиция события назначается не там, где кажется
Проще всего считать, что порядок задан идентификатором задачи River. Задачи получают идентификаторы из последовательности PostgreSQL, значит id = 100 вставлена раньше, чем id = 101.
Идентификатор действительно выдаётся раньше. Проблема в том, что выдаётся он в момент INSERT, то есть внутри транзакции, а видимой задача становится в момент COMMIT. Между этими моментами может пройти сколько угодно времени.

Порядок коммитов здесь обратный порядку идентификаторов. Воркер, забирающий доступные задачи, увидел сначала 101 и вполне мог её отработать до появления 100.
К этому добавляется вторая особенность последовательностей: значения из них не образуют непрерывный ряд. Значение, выданное транзакции, которая потом откатилась, не возвращается в последовательность, поэтому в идентификаторах остаются дыры. Из наблюдения id = 100, id = 101, id = 103 нельзя заключить ничего о том, что происходило с бизнес‑сущностью.
Вывод: технический идентификатор outbox не является бизнес‑порядком и не должен им становиться. Порядок нужно моделировать явно, об этом ниже.
Второе: параллельные воркеры
Это главная причина, и она не связана ни с какими особенностями хранения.
Пусть в базе последовательно закоммитились три факта одного заказа, и в River появились три задачи. Очередь обрабатывается пулом воркеров, размер пула задаётся MaxWorkers для очереди. Доступные задачи выбираются в порядке приоритета и scheduled_at, то есть выборка примерно соответствует порядку появления. Но выбирает их несколько воркеров одновременно, а дальше каждая задача живёт своей жизнью: у одной медленно ответил брокер, вторая попала на воркер, который в этот момент собирает мусор, третья прошла мгновенно.

Порядок выборки задач это не порядок их завершения. Никакой связи между задачами одного агрегата River не создаёт: для него это просто две независимые строки в таблице.
Отдельно стоит сказать про повторы. Если публикация упала с ошибкой, задача уходит на повтор с экспоненциальной задержкой. Политика по умолчанию даёт задержку перед следующей попыткой равную числу уже сделанных попыток в 4-й степени, в секундах с разбросом ±10%: 1 секунда, 16 секунд, 1 минута 21 секунда, 4 минуты 16 секунд и дальше. То есть после первой же ошибки окно рассинхронизации измеряется не миллисекундами, а минутами, и порядок нарушается уже почти наверняка.
Третье: продюсер Kafka
Даже если по одному агрегату публикует ровно один воркер за раз, остаётся сам продюсер. Он отправляет запросы брокеру асинхронно и может держать несколько запросов в полёте одновременно. Если первый запрос упал и ушёл на повтор, а второй прошёл, записи окажутся в партиции в обратном порядке.
Защита от этого называется идемпотентным продюсером: брокер нумерует записи от конкретного продюсера и отвергает нарушение нумерации, что заодно сохраняет порядок при повторах. Требования такие: acks=all, retries больше нуля и не более пяти запросов в полёте на соединение.
Дальше начинается место, на котором Go‑команды регулярно спотыкаются, потому что умолчания у клиентов разные.
Клиент | Идемпотентность по умолчанию | Что нужно проверить |
|---|---|---|
Java, Kafka 3.0 и новее | включена, если нет конфликтующих настроек | что настройки не выключили её неявно |
franz‑go (kgo) | включена | что не выставлен |
librdkafka, confluent‑kafka‑go | выключена | включить |
Sarama | выключена |
|
И главное ограничение, которое снимает половину иллюзий: идемпотентность действует в пределах одного экземпляра продюсера. Продюсер получает свой идентификатор при подключении, при перезапуске процесса он новый, а между разными подами это просто разные продюсеры. Никакого порядка между записями, отправленными двумя разными процессами, Kafka не обещает.
Отсюда требование к реализации, которое надо держать в голове: публикация события агрегата должна быть синхронной. Воркер отправляет запись, дожидается подтверждения и только после этого считает событие опубликованным и отдаёт очередь следующему событию. Отправка без ожидания подтверждения ломает порядок даже при одном воркере.
Четвёртое: параллельная обработка у потребителя
Порядок в партиции не спасает, если потребитель раздаёт записи в пул горутин. Сообщения читаются по порядку, а применяются как получится. Формально это уже не наша зона ответственности, но требование к потребителю стоит записать явно, иначе всю работу продюсера обесценит одна строчка в чужом сервисе.
Вся цепочка целиком

Гарантия существует только если каждая стрелка имеет собственный инвариант. Разорванная в любом месте цепочка не восстанавливается ниже по течению.
Это не дефект River
River делает ровно то, что обещает: атомарную постановку задачи, повторы, параллельную обработку, доставку не менее одного раза(at‑least‑once). Гарантию порядка внутри ключа он в открытой версии не заявляет, и претензий тут быть не может.
Свойство, о котором идёт речь, принадлежит не библиотеке, а шаблону. Любой outbox с пулом обработчиков ведёт себя так же, будь то River, Sidekiq, Celery или самописный поллер с SKIP LOCKED(поллер, разбирающий таблицу по строкам с пропуском занятых). Ошибка была на нашей стороне: мы вывели гарантию из свойства Kafka, которое к ней не относится.
Насколько это вероятно: считаем, а не угадываем
«Теоретически возможно» и «происходит каждый день» требуют разных решений. Прежде чем строить сериализацию, стоит измерить два числа.
Первое: сколько у нас пар соседних событий одного агрегата, попадающих в опасное окно. Опасное окно это примерно время публикации одного события, от момента, когда задача стала доступной, до подтверждения от Kafka. Если следующее событие того же агрегата коммитится внутри этого окна, две задачи будут выполняться одновременно и их порядок определит случай.
Запросы ниже предполагают таблицу событий с полями aggregate_id, aggregate_version и published_at. Почему схема именно такая, разбирается в разделе про версию событий, здесь она нужна только как источник чисел
WITH pairs AS ( SELECT aggregate_id, created_at, lag(created_at) OVER ( PARTITION BY aggregate_id ORDER BY aggregate_version ) AS prev_created_at FROM domain_event WHERE created_at >= now() - interval '7 days' ) SELECT count(*) FILTER (WHERE prev_created_at IS NOT NULL) AS total_pairs, count(*) FILTER ( WHERE created_at - prev_created_at < interval '100 milliseconds' ) AS risky_pairs FROM pairs;
Запрос берёт события за неделю, для каждого события находит предыдущее событие того же агрегата и считает, у скольких пар разрыв меньше ста миллисекунд. Оконная функция lag возвращает значение из предыдущей строки внутри группы, группа задана PARTITION BY aggregate_id, а порядок внутри группы задан версией.
Второе: сколько нарушений уже произошло. Если в таблице событий есть published_at, нарушение видно как пара, у которой предыдущая по версии запись опубликована позже текущей.
WITH ordered AS ( SELECT aggregate_id, aggregate_version, published_at, lag(published_at) OVER ( PARTITION BY aggregate_id ORDER BY aggregate_version ) AS prev_published_at FROM domain_event WHERE published_at IS NOT NULL AND published_at >= now() - interval '7 days' ) SELECT count(*) AS inversions FROM ordered WHERE prev_published_at > published_at;
У этого измерения есть две оговорки, и их надо озвучить вместе с результатом. published_at пишут разные поды, между их часами есть расхождение, поэтому на разрывах в единицы миллисекунд возможны и ложные срабатывания, и пропуски. И если таблицы событий нет, а считать приходится по river_job, завершённые задачи не хранятся вечно: их удаляет служебный процесс cleaner, по умолчанию через сутки после завершения. Окно наблюдения по river_job этим и ограничено.
Дальше арифметика. Числа ниже иллюстративные, подставьте свои.
Величина | Значение | Откуда |
|---|---|---|
Событий в сутки | 3 000 000 | метрика публикации |
Пар соседних событий одного агрегата | 1 800 000 | запрос выше |
Доля пар в окне 100 мс | 0,5% | запрос выше |
Пар в опасном окне за сутки | 9 000 | 1 800 000 умножить на 0,005 |
Доля пар, которые реально переставились | около половины | гонка симметрична |
Перестановок в сутки | около 4 500 | 9 000 разделить на 2 |
При таких числах перестановки идут постоянным фоном. Если же измерение даст двадцать пар в сутки, разговор совсем другой, и проверки на стороне потребителя может хватить.
Отдельно посчитайте вклад повторов. Каждая ошибка публикации отправляет задачу на паузу, а следующие события агрегата тем временем публикуются. При задержке в 16 секунд на второй попытке за это время может уехать вперёд весь дальнейший поток агрегата.
Что River OSS гарантирует, а что нет
Свойство | Есть в OSS | Комментарий |
|---|---|---|
Атомарная постановка задачи в чужой транзакции | да |
|
Повторы с экспоненциальной задержкой | да | по умолчанию 25 попыток, последняя примерно через три недели |
Параллельная обработка очереди | да |
|
Выборка задач в порядке приоритета и | да | это порядок выборки, не порядок завершения |
Отложить задачу без расхода попыток | да |
|
Выполнение не более одного раза | нет | гарантия не менее одного раза, дубли возможны |
Последовательное выполнение задач по ключу | нет | это Sequences в River Pro |
Две детали из этой таблицы стоит развернуть, потому что они влияют на расчёты дальше.
MaxWorkers ограничивает число воркеров одного экземпляра клиента, а не глобально. Три пода с MaxWorkers = 1 дадут три одновременно работающих воркера.
Гарантия «не менее одного раза» не теоретическая. В River есть служебный процесс rescuer, спасающий зависшие задачи: если задача выполняется дольше порога (по умолчанию час), она будет запланирована заново, возможно параллельно с уже идущим выполнением.. То есть одна и та же задача может выполняться дважды одновременно, и это документированное поведение, а не сбой.
River Pro Sequences: что покупается и что не покупается
Sequences это механизм River Pro, который выполняет задачи одной последовательности строго по одной за раз, при этом разные последовательности выполняются параллельно. Ключ последовательности вычисляется из вида задачи и её аргументов, то есть aggregate_id подходит на эту роль напрямую. Это ровно та модель, которая нужна.
Купив её, вы перестаёте писать собственную сериализацию. Но остаётся четыре вещи, о которых надо знать заранее.
Порядок внутри последовательности задаётся идентификатором задачи, по возрастанию. Это документировано. Отсюда мой вывод, который я отмечаю именно как вывод: механизм наследует проблему из первого раздела про идентификаторы. Задача, вставленная раньше, но закоммиченная позже, до коммита невидима, и последовательность может продвинуться без неё. Если две транзакции по одному агрегату идут параллельно, порядок выполнения может не совпасть с порядком коммитов. Значит, сериализовать назначение позиции на записи всё равно придётся, и покупка Pro от этой работы не освобождает. Проверять это утверждение надо тестом на своей версии, а не верить статье.
Отравленное событие остаётся на вашей стороне. По умолчанию последовательность останавливается, если задача отменена или отброшена, и продолжить её можно, либо повторив проблемную задачу, либо вручную запустив следующую в обход. Флаги ContinueOnCancelled и ContinueOnDiscarded меняют поведение на «идти дальше», то есть отказываются от строгой семантики. Выбор всё равно за вами, механизм только даёт форму этому выбору.
Первая задача в простаивающей последовательности ждёт дольше остальных: её переводит в доступное состояние периодический служебный процесс River Pro под названием служба сопровождения последовательностей (sequence maintenance). Последующие задачи, если они уже стоят в очереди, планируются сразу после завершения предыдущей. Для потока, где у большинства агрегатов одно‑два события, эта задержка платится почти на каждом событии, и её надо померить на своей нагрузке до принятия решения.
Ручное вмешательство ломает гарантию. Документация прямо говорит, что при ручном повторе задач в последовательности одновременное выполнение возможно. То есть кнопка «повторить» в River UI это операция, ломающая инвариант. Если такой интерфейс доступен дежурному, порядок держится на договорённости, а не на механизме.
Фундамент любого решения: явная версия события
Дальше речь про то, как получить порядок на открытой версии. Всё начинается с одного решения в модели данных, и без него не работает ни один из вариантов.
У каждого события должна быть позиция в потоке событий агрегата.
Поле | Что означает |
|---|---|
| к какому агрегату относится факт |
| какой это по счёту факт для данного агрегата |
| идентификатор конкретного факта, для идемпотентности потребителя |
Три требования к версии, каждое из которых легко нарушить.
Версия назначается в той же транзакции, что и изменение агрегата. Иначе она описывает порядок присвоения номеров, а не порядок фактов.
Версия нумерует события, а не изменения агрегата. Это тонкое место. Если версия живёт на строке агрегата и увеличивается на каждом изменении, а событие публикуется не на каждое изменение, то в потоке событий появятся дыры. Проверка «предыдущая опубликованная версия плюс один» на такой дыре встанет навсегда. Либо версия считается по потоку событий, либо проверка формулируется как «строго больше», а не «ровно на один больше», и тогда обнаружение потерянного события надо строить отдельно.
Версия защищена уникальным индексом:
ALTER TABLE domain_event ADD CONSTRAINT domain_event_position_unique UNIQUE (aggregate_id, aggregate_version);
Ограничение говорит базе, что у агрегата не может быть двух разных событий с одной позицией. Без него ошибка в назначении версии останется незамеченной до момента, когда её последствия проявятся у потребителя.
Сама транзакция выглядит так:
BEGIN; -- Захватываем строку агрегата и увеличиваем версию. -- UPDATE сам берёт блокировку строки и держит её до конца транзакции, -- поэтому отдельный SELECT FOR UPDATE здесь не нужен. UPDATE orders SET status = 'paid', version = version + 1 WHERE id = 123 RETURNING version; -- Записываем факт с той позицией, которую только что получили. INSERT INTO domain_event ( event_id, aggregate_id, aggregate_version, event_type, payload ) VALUES ( gen_random_uuid(), 123, $new_version, 'OrderPaid', $payload ); -- Ставим задачу на публикацию в том же коммите. INSERT INTO river_job (...) VALUES (...); COMMIT;
Почему это даёт правильный порядок. Блокировка строки, взятая UPDATE, удерживается до конца транзакции. Вторая транзакция, пришедшая по тому же заказу, на своём UPDATE встанет в ожидание и продолжит только после коммита первой. Значит, для одного агрегата транзакции выстраиваются в цепочку, и порядок присвоения версий совпадает с порядком коммитов. Для разных агрегатов блокировки разные, поэтому параллельность сохраняется.
Побочный эффект, о котором стоит знать: это точка сериализации на записи. Если по одному агрегату идёт высокая частота изменений, ожидание блокировки станет видно в задержках API. Обычно это приемлемо, но померить стоит.
Два способа обойтись без отдельной версии, оба не работают
Идентификатор задачи River как бизнес‑порядок не годится по причинам из первого раздела про потерю порядка.
Сортировка по created_at не годится тоже. Точности timestamp может не хватить на события, идущие подряд, часы разных экземпляров приложения расходятся, время создания записи не равно времени коммита, а при повторе появляется ещё и время обработки. Время это хороший атрибут для наблюдаемости и плохой механизм для строгого порядка.
Вариант A: проверка предыдущей версии и отложенный повтор
Самый простой по замыслу подход. Воркер, получив событие версии N, проверяет, опубликована ли версия N минус один. Если нет, откладывает задачу.

Откладывание в River делается возвратом river.JobSnooze(d) из воркера. River считает откладывание намеренным решением воркера, а не сбоем. Счётчик попыток растёт в момент взятия задачи в работу, поэтому при откладывании River уменьшает его обратно на единицу, и в сумме отложенная попытка лимит не расходует. Задача может откладываться сколько угодно раз, а в её метаданных растёт счётчик snoozes, и это готовая метрика.
Подход работает, но у него три цены, которые надо знать до того, как он попадёт в продакшен.
Цена первая, задержка. Отложенная задача становится доступной не сама по себе, а на очередном проходе планировщика River. Планировщик работает с постоянным интервалом в 5 секунд и не настраивается. Значит, каждое событие, приехавшее не в свою очередь, платит секунды, а не миллисекунды. Для потока, где перестановки редки, это незаметно. Для потока, где они фоновые, это означает, что заметная доля событий доезжает до Kafka с задержкой в единицы секунд.
Цена вторая, холостая работа воркеров. Отложенная задача воркера не удерживает: она просыпается, проверяет условие, возвращает snooze и освобождает его. Но каждая такая проверка это занятый воркер и запрос к базе. Если по агрегату скопилось десять событий, они будут просыпаться каждые пять секунд и уходить обратно, и на популярном агрегате эти холостые круги начнут отбирать пропускную способность у всех остальных.
Третье ограничение касается повторного выполнения задачи. Здесь легко ошибиться в обе стороны, поэтому по шагам
Если проверка написана строго, «моя версия равна последней опубликованной плюс один», и порядок действий такой: сначала публикация с ожиданием подтверждения, потом отметка в базе, то между разными версиями порядок держится и без всякой блокировки. Событие 7 физически не может начать публикацию, пока в базе нет отметки о шестом, а она появляется только после подтверждения от брокера. Проверка «строго больше» вместо «ровно на один больше» этого свойства не даёт: она разрешает публиковать седьмое, не дожидаясь шестого.
Ломается это на повторном выполнении одной и той же задачи. River даёт гарантию «не менее одного раза», а rescuer может запустить вторую копию задачи параллельно с первой, которая ещё выполняется. Две копии задачи шестой версии читают одно и то же состояние и обе идут публиковать. Пока одна висит на сети, вторая успевает отработать, состояние обновляется, публикуется седьмое событие, и только потом до брокера доходит запись от зависшей копии. В партиции окажется последовательность 6, 7, 6, а для потребителя, применяющего события подряд, это откат состояния назад.
Отсюда вывод: проверки версии достаточно для упорядочивания разных версий и недостаточно для защиты от повторного выполнения одной. Нужен механизм, который не даёт двум воркерам работать с одним агрегатом одновременно.
Вариант B: сериализация через строку состояния
Заводим отдельную таблицу, в которой на каждый агрегат приходится одна строка с последней опубликованной версией.
aggregate_id | last_published_version |
|---|---|
123 | 57 |
456 | 12 |
789 | 103 |
Воркер начинает с того, что берёт блокировку этой строки:
SELECT last_published_version FROM aggregate_publication WHERE aggregate_id = $1 FOR UPDATE;
FOR UPDATE даёт транзакции монопольное право на эту строку до конца транзакции. Пока первая транзакция её держит, любая другая транзакция, попытавшаяся сделать то же самое по тому же агрегату, будет ждать. По другим агрегатам блокировки независимы, поэтому параллельность сохраняется. Именно эта блокировка, а не сама проверка версии, превращает «проверить и опубликовать» в неделимую для конкурентов операцию.
Дальше воркер сравнивает last_published_version + 1 с версией своего события. Если совпало, публикует, обновляет строку и коммитит. Если нет, откладывает задачу.
Схема требует трёх вещей, без которых она ломается на нагрузке.
Строка состояния должна существовать. Если её нет, SELECT вернёт пустой результат, и воркер уйдёт по ветке, которую автор не предусмотрел. Создавать строку надо в той же бизнес‑транзакции, которая создаёт первое событие агрегата:
INSERT INTO aggregate_publication (aggregate_id, last_published_version) VALUES ($1, 0) ON CONFLICT (aggregate_id) DO NOTHING;
Ожидание на блокировке занимает соединение. Воркер, ждущий освобождения строки, держит соединение из пула и ничего не делает. Десятки таких воркеров на популярном агрегате исчерпают пул и уронят вместе с собой обычные запросы приложения. Правильное поведение здесь: не ждать. FOR UPDATE NOWAIT вернёт ошибку, FOR UPDATE SKIP LOCKED вернёт пустой результат, и в обоих случаях воркер должен отложить задачу и освободить соединение.
Блокировка держится только до конца транзакции. Значит, транзакция обязана быть открыта на момент публикации в Kafka. Транзакция, внутри которой ходят по сети, это отдельный риск: она удерживает соединение и мешает очистке старых версий строк. Ограничивайте время сетевой операции таймаутом продюсера и не делайте внутри такой транзакции ничего лишнего.
Альтернатива строке состояния это рекомендательная блокировка PostgreSQL (advisory lock) от хеша идентификатора агрегата. Транзакционный вариант, pg_advisory_xact_lock, снимается сам в конце транзакции, но требует держать транзакцию открытой на время публикации в Kafka. Сессионный вариант, pg_advisory_lock, живёт в соединении и позволяет разбить работу на короткие транзакции, но снимать его надо вручную, и незакрытая блокировка переживёт ошибку и уедет в пул вместе с соединением. Парные функции pg_try_advisory_xact_lock и pg_try_advisory_lock не ждут освобождения, а сразу сообщают, что ключ занят.
Механизм | Плюс | Минус |
|---|---|---|
Блокировка строки состояния | состояние публикации видно обычным SELECT | нужна своя таблица и её обслуживание |
| не нужна таблица, снимается сама | состояние блокировки не видно в бизнес‑данных, нужен отдельный запрос к системным представлениям |
Сессионная рекомендательная блокировка | не требует держать транзакцию | нужно гарантированно снимать вручную, утечка блокировки переживает ошибку |
Хеш идентификатора агрегата для рекомендательной блокировки означает коллизии: два разных агрегата могут получить один ключ блокировки и начать мешать друг другу. На корректность это не влияет, на параллельность влияет.
Вариант C: одна задача публикует всё, что накопилось
Этот вариант в исходной постановке задачи обычно пропускают, а он, на мой взгляд, лучший для открытой версии River.
Идея в смене единицы работы. Задача перестаёт быть «опубликовать событие версии 7» и становится «опубликовать всё неопубликованное по агрегату 123». Порядок при этом не нужно охранять между задачами, потому что внутри одной задачи он задан обычным ORDER BY.

Ключевой шаг здесь тот, что слева. Если блокировку агрегата держит другой воркер, текущая задача завершается успешно и ничего не делает. Ждать не нужно: тот, кто держит блокировку, всё равно опубликует и наши события тоже, потому что он публикует всё неопубликованное, а не своё конкретное. Это убирает и ожидание на блокировках, и цикл «проснуться, проверить, заснуть».
Функция за один заход публикует пачку событий одного агрегата и возвращает признак того, осталось ли что‑то ещё:
func (w *PublishWorker) Work(ctx context.Context, job *river.Job[PublishArgs]) error { for { done, err := w.publishBatch(ctx, job.Args.AggregateID) if err != nil { return err } if done { return nil } } } func (w *PublishWorker) publishBatch(ctx context.Context, aggregateID string) (bool, error) { tx, err := w.pool.Begin(ctx) if err != nil { return false, err } defer tx.Rollback(ctx) // Пробуем занять агрегат. SKIP LOCKED означает: если строку уже держит // другой воркер, вернётся пустой результат, и мы просто уходим. // Строка гарантированно существует: её создаёт бизнес-транзакция. var lastPublished int64 err = tx.QueryRow(ctx, ` SELECT last_published_version FROM aggregate_publication WHERE aggregate_id = $1 FOR UPDATE SKIP LOCKED`, aggregateID).Scan(&lastPublished) if errors.Is(err, pgx.ErrNoRows) { return true, nil // агрегат занят, публикация будет сделана без нас } if err != nil { return false, err } events, err := selectUnpublished(ctx, tx, aggregateID, lastPublished, batchSize) if err != nil { return false, err } if len(events) == 0 { return true, nil } // Публикуем строго по возрастанию версии и дожидаемся подтверждения // на каждой записи: без ожидания порядок теряется внутри продюсера. for _, e := range events { if err := w.producer.ProduceSync(ctx, e.ToRecord()); err != nil { return false, err } } if err := markPublished(ctx, tx, aggregateID, events); err != nil { return false, err } if err := tx.Commit(ctx); err != nil { return false, err } return len(events) < batchSize, nil }
Три момента реализации.
Публикация идёт по одному событию с ожиданием подтверждения. Пакетная отправка допустима, если клиент сохраняет порядок записей внутри одного запроса, но подтверждение всё равно надо дождаться до отметки об успехе.
Размер пачки ограничен. Без ограничения задача по агрегату с длинной историей будет выполняться часами, а rescuer через час сочтёт её зависшей и запустит вторую копию параллельно.
Гранулярность фиксации это осознанный выбор:
Как коммитим | Дублей при падении процесса | Транзакций на N событий |
|---|---|---|
Один коммит на всю пачку | до N | 1 |
Один коммит на событие | не более 1 | N |
Отдельно про уникальные задачи. В этом варианте в аргументах задачи лежит только aggregate_id, версии там нет: какие события публиковать, воркер выясняет запросом в момент выполнения. Уникальность в River считается по виду задачи и её аргументам, поэтому здесь возможен ровно один ключ уникальности, по aggregate_id. Схлопнуть пятьдесят задач в одну выглядит естественным улучшением, но набор состояний, по которым River считает задачи уникальными, по умолчанию включает running, и убрать это состояние из набора нельзя, оно обязательное. Значит, пока задача по агрегату выполняется, вставка новой задачи по тому же агрегату будет пропущена как дубль. Если выполняющаяся задача к этому моменту уже сделала свою выборку, новое событие не получит ни своей задачи, ни обработки в чужой, и останется неопубликованным до следующего события. Безопаснее не делать задачи уникальными: лишние задачи дёшевы, они завершаются немедленно, обнаружив занятый агрегат.
В вариантах A и B, где задача публикует конкретное событие и версия лежит в аргументах, уникальность по паре aggregate_id и version безопасна: совпадение ключа там означает настоящий дубль, а не потерянный сигнал.
Отравленное событие
Любой строгий порядок упирается в один вопрос: что делать, если конкретное событие не удаётся опубликовать никогда. Причины бывают разные: битый payload, схема, которую отвергает реестр схем (schema registry), если он у вас используется, слишком большое для топика сообщение.

Так ведёт себя любой строгий порядок. Если потребитель полагается на непрерывность, то опубликовать v4 без v3 означает солгать ему. Выбор придётся сделать явно.
Политика | Что получаем | Чем платим |
|---|---|---|
Блокировать агрегат до вмешательства человека | контракт не нарушается | поток по агрегату стоит, нужен дежурный и алерт |
Отправить событие в очередь недоставленных сообщений (dead letter queue), поток продолжить | поток живёт | контракт нарушен, потребитель должен это заметить |
Разрешить пропуск версии автоматически | поток не стоит никогда | это отказ от строгой семантики, и честнее признать его в контракте, чем в коде |
Пересоздать событие из текущего состояния агрегата | поток восстанавливается | нужен способ построить событие задним числом, не всегда возможен |
Отдельно про сроки. Число попыток задаётся параметром MaxAttempts, по умолчанию 25, и последняя попытка приходится примерно на третью неделю. То есть без явной политики агрегат может простоять заблокированным три недели, а потом задача будет отброшена, и поток встанет уже насовсем. Значение MaxAttempts для публикующих задач стоит выставить осознанно, а не оставлять умолчание.
Политика для отравленного события это бизнес‑решение, а не техническое. Ответ на вопрос «можно ли доставить v4 без v3» знает владелец потребителя, а не автор воркера.
Дубли неизбежны
Строгий порядок и доставка ровно один раз это разные требования, и второе достигается не там, где первое.
Публикация в Kafka и отметка в PostgreSQL это две разные системы, и одной транзакцией их накрыть нельзя. Любой порядок этих двух действий оставляет окно.

Обратный порядок даёт обратную беду: отметили в базе, упали, в Kafka ничего не уехало, событие потеряно навсегда. Из двух вариантов выбирают дубль: он лечится на потребителе, потеря не лечится ничем
River умеет немного сузить это окно. Есть механизм транзакционного завершения задачи: задачу можно пометить завершённой в той же транзакции, в которой воркер пишет свои изменения. Тогда «отметить событие опубликованным» и «завершить задачу» происходят атомарно, и остаётся ровно одно неатомарное место, отправка в Kafka. Убрать его нельзя, но одно окно вместо двух это заметно лучше.
Реалистичная формулировка цели: упорядоченная и надёжная доставка не менее одного раза плюс идемпотентные потребители. Для этого в сообщении нужен event_id, а потребитель хранит идентификаторы применённых событий и повтор игнорирует.
И важное разграничение: идемпотентность защищает от дубля, но не от нарушения порядка. Последовательность v7, v8, v6 не становится корректной от того, что у событий есть идентификаторы.
Вторая линия обороны у потребителя
Даже если продюсер всё делает правильно, проверка на стороне потребителя стоит дёшево и ловит то, что продюсер пропустил.
Потребитель хранит последнюю применённую версию по каждому агрегату и сравнивает её с версией пришедшего события.
Соотношение | Что это значит | Разумная реакция |
|---|---|---|
| всё по порядку | применить |
| дубль или устаревшее | отбросить |
| пропущено событие | зависит от типа потребителя |
Последняя строка это и есть точка выбора, о которой шла речь в начале. Потребителю, восстанавливающему состояние, дыра не мешает, он применяет и живёт дальше. Потребителю, реагирующему на переходы, дыра критична, и он должен либо ждать пропущенное событие в буфере, либо остановиться и позвать человека.
Буферизация на потребителе выглядит привлекательно, пока не начнёшь её проектировать. Сколько ждать пропущенное событие. Где держать буфер, чтобы он пережил перезапуск. Сколько памяти он займёт на плохом дне. Как отличить «событие задержалось» от «события не будет никогда». Что делать, если разрыв в сотню версий. Блокировать один агрегат или всю партицию, ведь остановка обработки партиции остановит и все остальные агрегаты в ней.
Это ровно та сложность, которую мы обсуждали на стороне продюсера, только перенесённая вниз по течению и размноженная по числу потребителей. Поэтому если порядок это системное свойство потока событий, обеспечивать его надо ближе к продюсеру, а на потребителе оставить проверку и алерт.
Что мерить
Своя реализация порядка без наблюдаемости неотличима от её отсутствия: нарушение проявится у потребителя через неделю, и связать его с причиной будет нечем.
№ | Метрика | На какой вопрос отвечает | Порог |
|---|---|---|---|
1 | Разрыв между текущей и опубликованной версией, максимум и p99* | насколько публикация отстаёт от изменений | максимум держится выше N дольше T** |
2 | Число агрегатов с неопубликованными событиями старше T | сколько потоков стоит прямо сейчас | больше нуля дольше T |
3 | Задержка от коммита до подтверждения Kafka, p99 | укладываемся ли в бюджет доставки | выше бюджета цепочки |
4 | Число откладываний задач | как часто срабатывает защита порядка | рост скорости относительно базовой линии |
5 | Число нарушений непрерывности версий | работает ли механизм вообще | больше нуля это инцидент |
6 | Число повторных публикаций одного | достаточно ли идемпотентны потребители | для калибровки, не для алерта |
*p99 это значение, которое не превышается в 99 процентах измерений,
**N и T это ваши пороги, их надо подставить по своему профилю нагрузки. Бюджет цепочки это суммарное время, которое вы готовы отдать на путь от коммита до применения события у потребителя
Четвёртая метрика в River достаётся бесплатно: при откладывании задачи растёт счётчик в её метаданных.
SELECT count(*) AS jobs_snoozed, coalesce(sum((metadata->>'snoozes')::int), 0) AS snoozes_total FROM river_job WHERE kind = 'publish_domain_event' AND metadata ? 'snoozes';
Пятая считается поиском дыр в потоке опубликованных версий:
WITH ordered AS ( SELECT aggregate_id, aggregate_version, lag(aggregate_version) OVER ( PARTITION BY aggregate_id ORDER BY aggregate_version ) AS prev_version FROM domain_event WHERE published_at IS NOT NULL ) SELECT aggregate_id, prev_version, aggregate_version FROM ordered WHERE prev_version IS NOT NULL AND aggregate_version <> prev_version + 1;
Запрос выбирает опубликованные события, для каждого берёт предыдущую опубликованную версию того же агрегата и оставляет только те пары, где разница не равна единице. Каждая строка результата это агрегат, у которого в опубликованном потоке дыра.
Чего я делать не стал бы
Ставить MaxWorkers = 1. Это не решает задачу и ломает пропускную способность. Не решает потому, что ограничение действует на экземпляр клиента: три пода дадут три воркера и ту же гонку. Ломает потому, что при публикации в 20 миллисекунд один воркер даёт 50 событий в секунду, и любое замедление брокера останавливает весь поток целиком, по всем агрегатам сразу.
Управлять порядком через приоритеты. Приоритет влияет на то, какие задачи будут выбраны раньше, и не создаёт между ними зависимости. Выборка это не выполнение.
Считать, что ключа Kafka достаточно. Ключ обеспечивает необходимое условие, попадание в одну партицию, и не обеспечивает достаточное, правильный порядок отправки.
Выводить бизнес‑порядок из идентификатора задачи или из времени создания записи. Причины разобраны выше, обе относятся к тому, что технические атрибуты описывают работу механизма, а не последовательность фактов.
Считать, что покупка River Pro закрывает вопрос целиком. Она закрывает сериализацию выполнения, оставляя вам назначение позиции, политику отравленного события и дисциплину ручных операций.
Что решить до выбора инструмента
Вопросы, которые я предлагаю обсудить командой. Разговор про River, Pro и подписки имеет смысл только после ответов на них.
Какой из четырёх порядков мы обязаны гарантировать. Мой ответ по умолчанию: порядок коммитов транзакций.
Кто конкретно из потребителей ломается при нарушении порядка. Если такого потребителя нет, дальше можно не читать: хватит версии и правила «побеждает последняя запись по версии».
Что считается доставкой. Подтверждение от брокера, доступность записи потребителю и успешная обработка потребителем это три разные гарантии с разной ценой.
Можно ли доставить событие N плюс один, если событие N не удаётся доставить никогда, и кто принимает это решение в момент инцидента.
Где живёт гарантия: у продюсера, у потребителя или на обоих уровнях.
Какие числа у нас на самом деле. Сколько пар событий попадает в опасное окно и сколько нарушений уже произошло за последнюю неделю.
Шестой пункт стоит выполнить до остальных: он превращает разговор о вероятностях в разговор о фактах.
Итог
Вариант | Порядок | Стоимость | Когда выбирать |
|---|---|---|---|
Без строгого порядка: версия в сообщении и правило «побеждает последняя запись» | не требуется, состояние сходится по версии | минимальная | ни один потребитель не реагирует на переходы, всем нужно текущее состояние |
River Pro Sequences | сериализует выполнение по ключу | подписка, задержка на первой задаче простаивающей последовательности, дисциплина ручных операций | нужен строгий порядок и не хочется писать и поддерживать свою сериализацию |
Вариант A: проверка предыдущей версии и отложенный повтор | даёт при строгой проверке «плюс один», монотонном обновлении состояния и проверке версии у потребителей | пять секунд на каждом переставленном событии, холостые круги воркеров, обязательный алерт на молчаливую остановку | минимум своего кода, события по агрегату идут редко, потребителей можно обязать проверять версию |
Вариант B: сериализация через строку состояния | даёт, конкурентов отсекает блокировка | транзакция открыта на время публикации, свои граничные случаи: создание строки состояния, NOWAIT или SKIP LOCKED, расход пула соединений | нужен строгий порядок, Pro недоступен, по агрегату идёт плотный поток |
Вариант C: одна задача публикует всё накопившееся по агрегату | даёт, порядок задан сортировкой по версии внутри задачи | требует лимита на размер пачки и явного выбора гранулярности коммита | мой выбор по умолчанию для OSS |
Что из этого следует. Каждая из трёх технологий в нашей цепочке работает корректно и делает ровно то, что обещает. PostgreSQL даёт атомарность, River даёт атомарную постановку задачи и повторы, Kafka даёт порядок внутри партиции. Проблема возникла не внутри компонента, а на стыке: мы сложили три частные гарантии и приняли сумму за общую, которой никто не давал.
Это, кажется, общее правило для распределённых систем. Гарантия существует только там, где её кто‑то явно предоставил, и путь целиком настолько же упорядочен, насколько упорядочено его самое слабое звено. Найти это звено дешевле на схеме, чем в инциденте.
Источники
Документация River:
Транзакционная постановка задач: https://riverqueue.com/docs/transactional‑enqueueing
Повторы задач и политика задержек: https://riverqueue.com/docs/job‑retries
Откладывание задач: https://riverqueue.com/docs/snoozing‑jobs
Служебные процессы, планировщик и спасение зависших задач: https://riverqueue.com/docs/maintenance‑services
Уникальные задачи: https://riverqueue.com/docs/unique‑jobs
Транзакционное завершение задачи: https://riverqueue.com/docs/transactional‑job‑completion
Надёжные воркеры и гарантия «не менее одного раза»: https://riverqueue.com/docs/reliable‑workers
Sequences в River Pro: https://riverqueue.com/docs/pro/sequences
Ограничения параллелизма в River Pro, там же про порядок выборки задач: https://riverqueue.com/docs/pro/concurrency‑limits
Служебные процессы: планировщик и rescuer, спасающий зависшие задачи: https://riverqueue.com/docs/maintenance‑services#rescuer
Kafka и клиенты:
Настройки продюсера Apache Kafka: https://kafka.apache.org/41/configuration/producer‑configs/
Публикация в franz‑go: https://github.com/twmb/franz‑go/blob/master/docs/producing‑and‑consuming.md
Конфигурация Sarama: https://github.com/IBM/sarama/blob/main/config.go
PostgreSQL:
Функции работы с последовательностями: https://www.postgresql.org/docs/current/functions‑sequence.html
Явные блокировки, включая рекомендательные: https://www.postgresql.org/docs/current/explicit‑locking.html
Отдельно про уровень уверенности. Всё, что касается настроек, умолчаний и заявленных гарантий, взято из документации по ссылкам выше и проверялось в августе 2026 года. Утверждение, что порядок выполнения внутри River Pro Sequences может разойтись с порядком коммитов из‑за видимости незакоммиченных задач, является моим выводом из документированного правила сортировки по идентификатору и правил видимости PostgreSQL. Я его не проверял экспериментально, и перед тем, как опираться на него в решении, стоит воспроизвести сценарий на своей версии.

