Как перестать писать свой планировщик: два уровня задач, защита от дублей и почему периодическая задача не должна делать работу сама.
Продолжение статьи «River в Go: фоновые задачи в одной транзакции с PostgreSQL». Пример тот же: API создаёт товар, воркер публикует его через медленный сервис склада.
Исходный код примера доступен в репозитории riverqueue-guide. Всё, что разобрано ниже, собрано в одном коммите: periodic jobs, part 2 c4d99a8. В конце статьи есть карта файлов и команды запуска.
Проблема
В поддержку приходит тикет: товар создан три часа назад и до сих пор не опубликован. Таких десятки: ночью сервис склада лежал два часа, задачи публикации перебрали попытки и ушли в discarded. По умолчанию River даёт задаче 25 попыток с растущим таймаутом, а дальше всё: river_job в discarded, товар навсегда с published_at = null. Восстанавливать такие записи руками это плохая работа.
Нужен механизм, который сам находит застрявшее и запускает задачи. Джоба вызывается не действием пользователя, а временем. Такая задача регулярно проходит по данным, ищет то, что осталось недоделанным, и возвращает это в обработку. Другие примеры таких джоб: периодическая синхронизация или прогрев кеша. Общая транзакция из первой статьи здесь не нужна. Там задача ставилась вместе с созданием товара, и обе записи были обязаны попасть в базу либо вместе, либо никак. Периодическая задача не создаёт бизнес-записей, она только читает и ставит другие задачи, поэтому связывать её вставку не с чем.
Далее рассмотрим варианты, начнем с самодельного тикера и четыре дыры, которые он оставляет, потом периодические задачи River, и наконец паттерн, который нам подошел.
Первая попытка: тикер в main
Первое, что приходит в голову, это горутина с тикером. В проекте из первой статьи два процесса: API и воркер. Тикеру место в воркере, рядом с River-клиентом. В internal/app/worker.go, сразу после client.Start(ctx):
// internal/app/worker.go (ДО)
if err := client.Start(ctx); err != nil {
log.Fatalf("could not start river worker client: %v", err)
}
// Sweeper: самодельное расписание.
go app.StartSweeper(ctx)
<-ctx.Done()
Сама горутина:
// internal/app/sweeper.go
func (a *App) StartSweeper(ctx context.Context) {
ticker := time.NewTicker(a.config.StuckPublicationsScanInterval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
// найди застрявшие и опубликовать, прямо здесь, без River
if err := a.productUseCase.PublishStuckProducts(ctx); err != nil {
a.logger.Error("sweeper failed", zap.Error(err))
}
}
}
}
ctx это тот же signal.NotifyContext, что и у River-клиента: SIGTERM отменит и его, и тикер. Кажется, всё на месте: расписание, контекст, лог ошибок, юзкейс с бизнес-логикой. Все просто и работает.
Но потом происходит инцидент на проде…
Сколько раз выполнится при двух репликах воркера? Дважды. Каждый тикер в каждом процессе независим: обе реплики в 12:00:00 найдут одну и ту же сотню застрявших товаров и начтут делать запросы к сервису склада. Можно делать свои костыли, типа блокировки в Redis или на уровне БД SELECT ... FOR UPDATE SKIP LOCKED.
Что случится при деплое? Процессу приходит SIGTERM: в Kubernetes его посылает kubelet, но в любой среде с корректной остановкой механика та же. Дальше ctx отменяется, и горутина выходит из select. Новый процесс стартует в случайный момент, и первый тик случится через полный интервал от старта горутины(например, при минутном интервале это от 0 до 60 секунд). В среднем полминуты простоя на каждый деплой. Тики, попавшие на паузу между процессами, не навёрстываются.
И это не единственное место, где оркестратор ломает самодельное расписание. При rolling update старый под живёт, пока новый проходит проверки готовности, поэтому какое-то время тикеров два даже при одной реплике. Переезд пода на другую ноду или перезапуск после OOM дают тот же эффект, что и деплой: фаза расписания сдвигается, и заметить это неоткуда.
Где история и повторы? Нигде. Если PublishStuckProducts упал на сороковом товаре из сотни, узнает об этом только лог. Что произошло с первыми тридцатью девятью и что с оставшимися шестьюдесятью, узнать негде: прогресс существовал только в локальных переменных и исчез вместе с упавшим вызовом. В базе не осталось ни отметки о попытке, ни точки, с которой можно продолжить.
Кто это видит? К сожеланию, никто, надо смотреть логи. У такого способа нет интерфейса, в котором можно наблюдать за джобами: ни списка запусков, ни их результатов. Только строки в логе, которые надо знать, как искать.
Все эти вопросы не про бизнес-логику. PublishStuckProducts написан правильно; проблемы начинаются выше, где решается «когда и сколько раз его вызвать».
Понимаем, что способ нам не подходит и идём в документацию River. Смотрим что такое periodic jobs.
Что такое periodic job
Термин из документации River, и первое, что нужно понять: периодическая задача не «выполняется по расписанию», она вставляет другую задачу по расписанию. В конфиге клиента задаётся не воркер с расписанием, а генератор:
river.NewPeriodicJob(
river.PeriodicInterval(interval), // когда
func() (river.JobArgs, *river.InsertOpts) {
return args, opts // что вставить
},
&river.PeriodicJobOpts{RunOnStart: true},
)
Одно срабатывание такого расписания дальше по тексту называется постановкой по расписанию. Устроена она предельно просто: River вызывает конструктор и делает обычный INSERT в таблицу river_job. На этом постановка закончилась, никакой работы она не выполняет. Дальше вставленная задача живёт по общим правилам: очередь, воркеры, попытки, уникальность.
Из такого устройства следуют три факта, от которых зависит всё остальное:
- Расписание обслуживает только лидер. Клиенты River в кластере выбирают лидера через таблицы River в PostgreSQL, и вставки делает только он. Десять реплик воркера дают одно расписание, а не десять. Но это относится только к клиентам с воркерами: insert-only клиент API из первой статьи расписание не обслуживает.
- Шедулер держит состояние в памяти. Он помнит только, когда была последняя постановка и когда должна быть следующая, и помнит это в памяти процесса. Процесс умер или выбрался новый лидер, и отсчёт начинается заново, без оглядки на предыдущий. Пропущенные постановки не навёрстываются: не было постановки, значит не было и вставки, а догонять River нечем, потому что записи о пропуске нигде нет.
RunOnStartэто страховка, которую рекомендует документация: вставить задачу сразу при старте шедулера, не дожидаясь первой постановки. - Конструктор не блокирует. Он вызывается синхронно внутри постановки, поэтому его работа вернуть args, а не ходить в базу: пока он думает, стоят постановки всех остальных периодических задач процесса. Вернул
nil, значит постановка прошла вхолостую и вставки не было.
Этих трёх фактов достаточно, чтобы переписать тикер. Первая версия получится не сказать чтобы совершенной, и именно её ошибка объясняет, почему рабочее решение будет выглядеть иначе.
Вторая попытка: одна задача на всю работу
Переписываем: не своя горутина, а PeriodicJobs в конфиге клиента, и вся работа прямо в Work:
// НЕ РАБОТАЕТ: одна задача на всю пачку
PeriodicJobs: []*river.PeriodicJob{
river.NewPeriodicJob(
river.PeriodicInterval(cfg.StuckPublicationsScanInterval),
func() (river.JobArgs, *river.InsertOpts) {
return jobs.SweepStuckPublicationsArgs{}, nil
},
&river.PeriodicJobOpts{RunOnStart: true},
),
},
func (w *SweepStuckPublicationsWorker) Work(
ctx context.Context,
job *river.Job[jobs.SweepStuckPublicationsArgs],
) error {
productIDs, err := w.productRepo.GetStuckPublicationIDs(ctx, w.batchLimit)
if err != nil {
return fmt.Errorf("select stuck: %w", err)
}
for _, id := range productIDs { // вся пачка в одной задаче
if err := w.publisher.Publish(ctx, id); err != nil {
return fmt.Errorf("publish %d: %w", id, err)
}
}
return nil
}
Деплоим, и сразу красиво: задача в River UI, история попыток, упавшая задача видна. Две реплики больше не дублируют работу: расписание обслуживает только лидер. Все четыре вопроса про тикер закрыты.
Открываем UI после ночного падения склада и видим: задача шла 40 минут, упала на 71-м товаре из 100, перезапустилась и начала пачку заново.
Почему сломалось
Одна постановка, одна задача, один счётчик попыток на всю пачку. Из этого следует всё:
- Падение теряет прогресс. 70 товаров опубликованы, но попытка не засчитана: ретрай начнёт пачку с начала. Идемпотентность воркера спасает от двойной публикации, но не от 40 минут повторной работы.
- Один медленный сервас склада отражается на задаче. Сотня товаров по 5–10 секунд даёт задачу длиной в десятки минут. Дальше подключается таймаут самого воркера и снимает задачу на середине.
- Ретраи задачи пересекаются с постановками по расписанию. Задача ещё ретраится, а следующая уже вставлена. Плюс
MaxAttemptsпо умолчанию равен 25: получаем два независимых цикла повтора одной и той же пачки.
Цепочка целиком:
периодическая задача «сделай всю работу»
└── Work: селект пачки + цикл публикации
├── один job = один счётчик попыток на 100 товаров
├── падение на 71-м: 70 опубликованы, попытка не засчитана
├── ретрай: пачка с начала
└── расписание продолжает вставлять новые задачи, пока старая ретраится
Рабочее решение: диспетчер и исполнители
Задача одна на сотню товаров, вот в чём ошибка второй попытки. Делим. Правильная форма это два уровня:
расписание (шедулер River, лидер)
│ раз в N
▼
dispatch_stuck_publications ← периодическая задача: только найти и поставить
│ селект пачки ID + InsertMany, одна транзакция
▼
river_job × batch_limit ← обычные задачи: по одной на товар
│
▼
publish_product воркер ← ретраи, изоляция ошибок и таймауты по одному
Периодическая задача теперь только находит работу и ставит обычные задачи, по одной на сущность. Ретраи, таймауты и изоляция ошибок из первой статьи достаются каждой такой задаче по отдельности, со своим счётчиком попыток.
Аргументы и уникальность
Диспетчер:
type DispatchStuckPublicationsArgs struct {
BatchLimit int `json:"batch_limit"`
}
func (DispatchStuckPublicationsArgs) Kind() string {
return "dispatch_stuck_publications"
}
func (DispatchStuckPublicationsArgs) InsertOpts() river.InsertOpts {
return river.InsertOpts{MaxAttempts: 1}
}
Здесь два решения, которые выглядят наоборот относительно советов первой статьи, и оба осознанные.
MaxAttempts: 1. Одна попытка это обычно плохо. Задача упадёт из-за случайной ошибки сети, уйдёт в discarded, и повторить её будет некому, хотя хватило бы одного повтора через секунду.
С диспетчером не так: повтор у него уже есть, и это следующая постановка по расписанию. Селект упал, задача ушла в discarded и в ErrorHandler, а через минуту расписание поставило новую. Терять при этом нечего, потому что диспетчер не меняет данные: он только читает базу и ставит задачи, а прочитать то же самое второй раз не страшно. Если добавить ему ещё и внутренние ретраи, повторов станет два, свои у задачи и по расписанию, и вернётся путаница из второй попытки.
Так можно не с любой периодической задачей. Если она отправляет письмо или списывает деньги, второй запуск это уже не то же самое, что первый, и одну попытку ей ставить нельзя.
BatchLimit в args, а не в замыкании конструктора. Аргументы задачи видны в River UI и в river_job: можно понять, с какими параметрами работал конкретный запуск. Значение из замыкания видно только в конфиге процесса.
Исполнитель это тот же PublishProductArgs из первой статьи, но теперь с уникальностью:
func (PublishProductArgs) InsertOpts() river.InsertOpts {
return river.InsertOpts{
UniqueOpts: river.UniqueOpts{ByArgs: true},
}
}
В первой статье уникальность была опцией: постановка одна на создание товара, дубли маловероятны. Теперь дубли штатная ситуация: застрявший товар попадает в выборку диспетчера, пока оригинальная задача ещё ретраится. ByArgs дедуплицирует по аргументам: задача с тем же product_id уже стоит в очереди или выполняется, поэтому повторная вставка пропускается. Unique jobs, как и в первый раз, не отменяют идемпотентность воркера, они лишь не плодят очередь.
Конфигурация клиента
// internal/app/worker.go (РАБОЧИЙ ВАРИАНТ)
client, err := river.NewClient(
riverdatabasesql.New(container.DB),
&river.Config{
Workers: workers,
Queues: map[string]river.QueueConfig{
river.QueueDefault: {MaxWorkers: maxWorkers},
},
PeriodicJobs: periodicJobs(cfg),
SoftStopTimeout: softStopTimeout,
},
)
Само расписание вынесено в отдельную функцию того же файла:
func periodicJobs(cfg *config.Config) []*river.PeriodicJob {
return []*river.PeriodicJob{
river.NewPeriodicJob(
river.PeriodicInterval(cfg.StuckPublicationsScanInterval),
func() (river.JobArgs, *river.InsertOpts) {
return jobs.DispatchStuckPublicationsArgs{
BatchLimit: cfg.StuckPublicationsBatchLimit,
}, &river.InsertOpts{
UniqueOpts: river.UniqueOpts{
ByPeriod: cfg.StuckPublicationsScanInterval,
},
}
},
&river.PeriodicJobOpts{RunOnStart: true},
),
}
}
Три решения в этом блоке.
Интервал и размер пачки живут в конфиге приложения, а не в коде. Частота проверок это эксплуатационный параметр: его подбирают по метрикам (сколько застрявших копится, сколько выдерживает сервис склада). В примере это STUCK_PUBLICATIONS_SCAN_INTERVAL и STUCK_PUBLICATIONS_BATCH_LIMIT из окружения.
UniqueOpts{ByPeriod: interval} на задачу диспетчера. Если предыдущий диспетчер ещё не выполнился (сервис склада лежит, а пачка ушла в работу), новая постановка не создаст дубль: вставка будет пропущена как дубликат периода.
RunOnStart: true. Страховка на случай смены лидера: длинный интервал может не добежать до срабатывания раньше, чем шедулер перезапустится. В связке с ByPeriod безопасно: даже если старый лидер успел вставить задачу в этом периоде, повторная вставка будет пропущена.
Селект: grace-период и очередь по возрасту
Селект это место, где живёт главный практический навык диспетчера. Здесь появляется grace-период: пауза между появлением товара и моментом, когда его можно считать застрявшим. Пока пауза не истекла, товар не забота диспетчера, им ещё занимается оригинальная задача публикации.
func (r *productRepo) GetStuckPublicationIDs(ctx context.Context, batchLimit int) ([]int64, error) {
var productIDs []int64
err := r.dbFromCtx(ctx).WithContext(ctx).
Model(&models.Product{}).
Where("published_at IS NULL").
Where("created_at <= NOW() - interval '10 minutes'").
Order("created_at ASC").
Limit(batchLimit).
Pluck("id", &productIDs).Error
return productIDs, err
}
Grace-период в условии. created_at <= NOW() - interval '10 minutes', а не просто published_at IS NULL: десять минут это время, за которое оригинальная задача публикации успевает отработать ретраи. Без него диспетчер лезет в работу, которая ещё идёт. Unique jobs прикроют часть таких случаев, но условие в запросе надёжнее: диспетчер забирает только то, что действительно застряло.
Order("created_at ASC") + Limit. Выборка застрявших является очередью с приоритетом по дате создания: самые старые идут первыми, а пачка ограничивает всплеск после инцидента. Если за час простоя склада застряло 10 000 товаров, первый запуск диспетчера возьмёт 100, а не все десять тысяч. Очередь из десяти тысяч задач долбила бы едва поднявшийся склад запросами без остановки, пока не разберётся. В лучшем случае на стороне склада сработает rate limiter: часть запросов отобьётся, и задачи уйдут в ретраи. В худшем склад ляжет снова. Пачка вместо всего бэклога даёт ему передышку между постановками.
Диспетчер: селект и вставка в одной транзакции
Теперь глянем на юзкейс:
func (u *ProductUseCase) DispatchStuckPublications(
ctx context.Context,
batchLimit int,
) (enqueued, skipped int, err error) {
productIDs, err := u.productRepo.GetStuckPublicationIDs(ctx, batchLimit)
if err != nil {
return 0, 0, fmt.Errorf("select stuck publications: %w", err)
}
if len(productIDs) == 0 {
return 0, 0, nil
}
err = u.tm.InTransaction(ctx, func(ctx context.Context) error {
enq, skp, err := u.jobEnqueuer.PublishProducts(ctx, productIDs)
if err != nil {
return fmt.Errorf("enqueue publish jobs: %w", err)
}
enqueued, skipped = enq, skp
return nil
})
if err != nil {
return 0, 0, err
}
return enqueued, skipped, nil
}
Транзакция здесь нужна не для того, для чего в первой статье: бизнес-записи, к которой её пришлось бы привязывать, нет. На атомарность она тоже не влияет. InsertManyTx отправляет всю пачку одной инструкцией INSERT, а одиночная инструкция в PostgreSQL и так выполняется целиком или не выполняется вовсе. Половины пачки в river_job не окажется и без всякого BEGIN.
Осталась она по другой причине. Постановщик из первой статьи умеет ровно одно: писать в транзакцию, которую нашёл в контексте. Без транзакции ему понадобился бы второй путь вставки, и в адаптере появилась бы развилка на ровном месте. Открыть транзакцию вокруг одной инструкции дешевле: это лишние BEGIN и COMMIT раз в минуту против нового кода, который пришлось бы поддерживать.
Постановщик является адаптером, но для пачки, InsertManyTx:
func (e *PublishProductEnqueuer) PublishProducts(
ctx context.Context,
productIDs []int64,
) (enqueued, skipped int, err error) {
sqlTx, err := sqlTxFromCtx(ctx)
if err != nil {
return 0, 0, err
}
params := make([]river.InsertManyParams, 0, len(productIDs))
for _, productID := range productIDs {
params = append(params, river.InsertManyParams{
Args: jobs.PublishProductArgs{ProductID: productID},
})
}
results, err := e.client.InsertManyTx(ctx, sqlTx, params)
if err != nil {
return 0, 0, err
}
for _, result := range results {
if result.UniqueSkippedAsDuplicate {
skipped++
continue
}
enqueued++
}
return enqueued, skipped, nil
}
InsertManyTx вместо цикла InsertTx даёт один round-trip к базе на пачку. Счётчики enqueued/skipped могут использоваться для разных метрик.
Вся логика в юзкейсе:
func (w *DispatchStuckPublicationsWorker) Work(
ctx context.Context,
job *river.Job[jobs.DispatchStuckPublicationsArgs],
) error {
enqueued, skipped, err := w.productDispatcher.DispatchStuckPublications(ctx, job.Args.BatchLimit)
if err != nil {
log.Printf("[river] job=%d kind=%s attempt=%d: FAILED (%v)",
job.ID, job.Kind, job.Attempt, err)
return fmt.Errorf("dispatch stuck publications: %w", err)
}
log.Printf("[river] job=%d kind=%s attempt=%d: dispatched enqueued=%d skipped_duplicates=%d",
job.ID, job.Kind, job.Attempt, enqueued, skipped)
return nil
}
Деплоим. В River UI аккуратные постановки диспетчера раз в минуту и по сто publish_product с индивидуальными попытками. Следующую ночь склад лежит два часа: утром enqueued показывает разбирание бэклога пачками, skipped_duplicates показывает нули, застрявших discarded после публикации нет. Тикеты из поддержки про товар, который создан три часа назад и до сих пор не опубликован, больше не приходят.
Как это наблюдать и проверять
River UI. Задачи диспетчера это история постановок; задачи исполнителей это индивидуальные попытки. RunOnStart показывает результат сразу после старта, не дожидаясь интервала.
Счётчики диспетчера. enqueued/skipped в логе: если skipped стабильно растёт, оригинальные задачи не успевают отрабатывать, grace-период стоит увеличить.
Ошибки диспетчера. У задачи диспетчера MaxAttempts: 1, поэтому ошибка сразу уходит в discarded, минуя ретраи. В примере воркер её логирует; в рабочем проекте на ErrorHandler клиента подключают Sentry или другой сборщик. Читается такая ошибка всегда одинаково: «этот запуск диспетчера не удался», а не «товар не будет опубликован никогда».
Тесты. Пачку вставок проверяет rivertest.RequireManyInsertedTx, а селект табличный тест: свежий товар в выборку не попадает, созданный 11 минут назад попадает, уже опубликованный не попадает, а когда застрявших больше, чем помещается в пачку, первыми идут самые старые.
rivertest.RequireManyInsertedTx[*riverdatabasesql.Driver](
ctx, t, env.sqlTx(t), []rivertest.ExpectedJob{
{Args: jobs.PublishProductArgs{ProductID: oldestID}},
{Args: jobs.PublishProductArgs{ProductID: newestID}},
},
)
Подводные камни
Планировщик не навёрстывает пропуски, и это не баг River. Документация NewPeriodicJob говорит прямо: шедулер стартует лидером, держит состояние в памяти и после рестарта начинается с нуля. Деплой в 03:00 при часовом интервале означает, что постановки в 03:30 не будет. Для нашего случая это нормально: пропущенная постановка просто сдвинула проверку. Критерий: если пропуск одного запуска это инцидент, OSS periodic jobs не подходят: смотрите durable-расписания River Pro или внешний планировщик, который делает Insert в River.
Только интервалы, и фаза плавает. PeriodicInterval это интервал от старта шедулера, не cron-выражение и не «круглое» время: рестарт сдвигает все последующие запуски. Нужны «по будням в 9:00» с таймзонами: интерфейс PeriodicSchedule публичный, подключайте сторонний cron-пакет, но переходы на летнее время останутся на вас.
Интервалы от минуты. Рекомендация из документации: не меньше минуты, и никогда меньше секунды. Желание «каждые 5 секунд» обычно означает, что случай не про расписание, а про реакцию на событие, и тогда это первая статья, InsertTx в транзакции пользователя.
MaxAttempts: 1 годится только для идемпотентного диспетчера. У которого за спиной расписание. Для периодической задачи с несимметричным эффектом (письмо, квота) одна попытка это потерянный сбой.
Unique ≠ идемпотентность. ByPeriod и ByArgs защищают от повторной постановки, не от повторного выполнения. Условный UPDATE ... WHERE published_at IS NULL из первой статьи остаётся на месте.
Уникальность по умолчанию учитывает и completed. Набор состояний ByState, который River подставляет сам, включает завершённые задачи, поэтому повторно поставить publish_product для того же товара нельзя, пока чистильщик не уберёт старую строку (по умолчанию через сутки). Диспетчеру это не мешает: опубликованный товар в выборку не попадает. Зато удивит, если сбросить published_at руками и ждать новую задачу.
Не ходите в базу из конструктора. Он вызывается синхронно, и запрос из него подвешивает шедулер вместе со всеми периодическими задачами процесса. Пусть конструктор всегда возвращает args, а «работы нет» решает воркер, как диспетчер в примере.
Выводы
- Если выбрать самодельный тикер надо учитывать четыре проблемы: дубли на репликах, потерянные тики при деплое, отсутствие истории, невидимость. Periodic jobs River закрывают всё сразу за счёт одного факта: расписание обслуживает только лидер.
- «Жирная» периодическая задача возвращает проблемы тикера в новом виде: один счётчик попыток на пачку, падение теряет прогресс, ретраи пересекаются с постановками.
- «Диспетчер + исполнители»: периодическая задача только находит работу и ставит обычные задачи по одной на сущность. Ретраи, таймауты и изоляция ошибок остаются у исполнителей, каждая со своим счётчиком.
- Уникальность из опции становится обязательной:
ByPeriodна диспетчера,ByArgsна исполнителей. А grace-период в запросе, который ищет застрявших, не даёт диспетчеру забирать товары, над которыми оригинальная задача публикации ещё работает.
Вывод: расписание отвечает только за «когда»: шедулер River раз в интервал вставляет задачу диспетчера → диспетчер находит пачку застрявших и в одной транзакции ставит по одной задаче на сущность → исполнители делают работу с индивидуальными ретраями. И чем меньше периодическая задача знает о том, что происходит дальше, тем надёжнее вся конструкция.
Код примера
Всё разобранное выше лежит в репозитории riverqueue-guide, в коммите periodic jobs, part 2 c4d99a8 от 5 сентября 2026. Он идёт сразу за коммитом первой статьи, create jobs via river, part 1 534dfb0, поэтому git show c4d99a8 покажет ровно ту разницу, которую разбирает этот текст.
Как читать пример
Путь одной постановки по файлам репозитория:
шедулер River (обслуживает только лидер)
│ раз в STUCK_PUBLICATIONS_SCAN_INTERVAL
▼
internal/app/worker.go
periodicJobs(cfg) интервал, конструктор args, RunOnStart
│ INSERT dispatch_stuck_publications (ByPeriod, MaxAttempts: 1)
▼
river/product/worker/dispatch_stuck_publications.go
DispatchStuckPublicationsWorker.Work тонкий адаптер, счётчики в лог
│
▼
usecases/product/product_usecase.go
ProductUseCase.DispatchStuckPublications
│
├── productRepo.GetStuckPublicationIDs grace-период, старые первыми, batch limit
│ repositories/product/product_repo.go
│
└── tm.InTransaction
└── jobEnqueuer.PublishProducts InsertManyTx, по одной задаче на товар
river/product/enqueuer/publish_product.go
│
▼
river_job × N: publish_product, уникальные ByArgs
│
▼
PublishProductWorker.Work из первой статьи
Если открываете репозиторий впервые, читайте файлы в таком порядке:
internal/app/worker.go: регистрация воркера диспетчера и функцияperiodicJobsс интервалом, конструктором args иRunOnStart.internal/infrastructure/river/product/jobs/dispatch_stuck_publications.go: аргументы диспетчера иMaxAttempts: 1.internal/domain/product/usecases/product/product_usecase.go:DispatchStuckPublications, селект и вставка в одной транзакции.internal/domain/product/repositories/product/product_repo.go: селектGetStuckPublicationIDsс grace-периодом.internal/infrastructure/river/product/enqueuer/publish_product.go:PublishProductsнаInsertManyTxи счётчикиenqueued/skipped.internal/infrastructure/river/product/jobs/publish_product.go: уникальностьByArgsу исполнителя, единственное изменение в коде из первой статьи.internal/infrastructure/river/product/worker/dispatch_stuck_publications.go: воркер диспетчера.
В таком порядке сначала видно расписание и решение «кого поставить», а уже потом детали адаптеров. Тесты лежат рядом с кодом: product_repo_test.go проверяет селект таблицей, product_usecase_test.go проверяет пачку вставок и повторную постановку с skipped_duplicates.
Как запустить
Всё живёт в docker compose. Сначала поднимаем стек, и только потом создаём застрявшие товары: make stuck ходит в уже работающий PostgreSQL.
docker compose up -d --build
Лог воркера смотрим с фильтром. Это важная деталь: в dev-режиме воркер запущен через air, который при старте печатает десятки строк watching ..., и строки River в них теряются.
docker compose logs -f worker | grep river
Первая постановка диспетчера появляется сразу, не дожидаясь интервала, благодаря RunOnStart. Работы пока нет, поэтому оба счётчика нулевые:
2026/09/06 11:22:50 [river] job=49 kind=dispatch_stuck_publications attempt=1: dispatched enqueued=0 skipped_duplicates=0
Теперь создаём то, что диспетчеру придётся забрать. make stuck вставляет товар старше grace-периода и без задачи публикации за спиной, то есть ровно то состояние, которое остаётся после discarded:
make stuck
В пределах одного интервала сканирования диспетчер забирает пачку и ставит по задаче на товар. Дальше видно то самое разделение на два уровня: одна строка расписания и три независимые задачи, у каждой свой job и свой attempt.
2026/09/06 11:23:10 [river] job=50 kind=dispatch_stuck_publications attempt=1: dispatched enqueued=3 skipped_duplicates=0
2026/09/06 11:23:11 [river] job=53 kind=publish_product attempt=1 product_id=11: started
2026/09/06 11:23:11 [river] job=51 kind=publish_product attempt=1 product_id=9: started
2026/09/06 11:23:11 [river] job=52 kind=publish_product attempt=1 product_id=10: started
Подсвеченные строки здесь и дальше это постановки по расписанию, всё остальное работа исполнителей.
Дедупликацию удобно увидеть, если сделать так, чтобы публикация не успевала закончиться между постановками: поднять WAREHOUSE_MIN_DELAY и WAREHOUSE_MAX_DELAY и опустить STUCK_PUBLICATIONS_SCAN_INTERVAL. Тогда следующий запуск диспетчера находит те же товары всё ещё неопубликованными и пытается поставить их снова:
2026/09/06 11:23:30 [river] job=54 kind=dispatch_stuck_publications attempt=1: dispatched enqueued=0 skipped_duplicates=3
2026/09/06 11:23:50 [river] job=58 kind=dispatch_stuck_publications attempt=1: dispatched enqueued=0 skipped_duplicates=3
Вот ради этих двух строк и стоит UniqueOpts{ByArgs: true}: enqueued=0 при skipped_duplicates=3 означает, что расписание продолжает работать, а очередь не растёт. Без уникальности каждая постановка добавляла бы ещё три задачи на те же товары.
Та же картина со стороны очереди доступна через ./scripts/river_jobs.sh list и River UI на localhost:8083. Тесты и линтер запускаются оттуда же:
docker compose run --rm tests
docker compose run --rm lint