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

Продолжение статьи «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 из первой статьи

Если открываете репозиторий впервые, читайте файлы в таком порядке:

  1. internal/app/worker.go: регистрация воркера диспетчера и функция periodicJobs с интервалом, конструктором args и RunOnStart.
  2. internal/infrastructure/river/product/jobs/dispatch_stuck_publications.go: аргументы диспетчера и MaxAttempts: 1.
  3. internal/domain/product/usecases/product/product_usecase.go: DispatchStuckPublications, селект и вставка в одной транзакции.
  4. internal/domain/product/repositories/product/product_repo.go: селект GetStuckPublicationIDs с grace-периодом.
  5. internal/infrastructure/river/product/enqueuer/publish_product.go: PublishProducts на InsertManyTx и счётчики enqueued/skipped.
  6. internal/infrastructure/river/product/jobs/publish_product.go: уникальность ByArgs у исполнителя, единственное изменение в коде из первой статьи.
  7. 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