Если фоновая задача появляется вслед за записью в PostgreSQL, возникает неприятный вопрос: что случится, если данные уже сохранены, а задача ещё не попала в очередь? River позволяет убрать это окно отказа: бизнес-запись и задача коммитятся одной транзакцией.

Пример написан на Go, River, GORM и Gin. API создаёт товар, а отдельный воркер резервирует остаток и завершает публикацию. За остатки отвечает сервис склада: чужая система учёта по HTTP, которая живёт за пределами нашей базы, отвечает 5–10 секунд и ничего не знает про наши транзакции. Дальше по тексту «сервис склада» означает именно его. По ходу примера проверим атомарную постановку задачи, длину транзакции и поведение при повторном выполнении.

Исходный код примера доступен в репозитории riverqueue-guide.

Пример рассчитан на разработчика, который уже знаком с context.Context, database/sql и обычной транзакцией в PostgreSQL. Для локального запуска нужны Docker и Docker Compose.

Проблема, которую решает River

Две независимые записи

Представим обычный HTTP-сценарий:

  1. API добавляет товар в таблицу products.
  2. После коммита отправляет задачу публикации в отдельный брокер.
  3. Воркер получает задачу и резервирует остаток в сервисе склада.

Между первым и вторым действием процесс может упасть. Товар уже существует, но задача потеряна, поэтому публикация никогда не закончится. Обратный порядок тоже небезопасен: брокер может принять задачу, а транзакция товара откатится.

Это классическая проблема двух независимых записей, или dual write.

Почему enqueue внутри транзакции не всегда безопасен

Можно попробовать поставить задачу до COMMIT и написать такой код:

err := db.Transaction(func(tx *gorm.DB) error {
	product, err := createProduct(tx, input)
	if err != nil {
		return err
	}

	return queue.Enqueue(PublishProductArgs{ProductID: product.ID})
})

if err != nil {
	return err
}

Выглядит так, будто бизнес-логика и задача находятся в одной транзакции. Но если за queue стоит отдельный брокер, он ничего не знает про *gorm.DB и PostgreSQL-транзакцию. Брокер может принять и отдать задачу воркеру сразу после Enqueue, пока приложение ещё не выполнило COMMIT.

Тогда события пойдут в таком порядке:

  1. API вставляет товар, но ещё не коммитит транзакцию.
  2. Брокер принимает задачу и немедленно запускает воркер.
  3. Воркер делает SELECT товара из другого соединения.
  4. PostgreSQL не показывает ему незакоммиченную строку, поэтому воркер получает not found.
  5. После этого API коммитит транзакцию, но задача уже могла завершиться ошибкой или быть отменена как невыполнимая.

PostgreSQL не допускает dirty read: при стандартном Read Committed запрос видит только данные, зафиксированные до начала запроса. Даже запрошенный уровень Read Uncommitted в PostgreSQL фактически работает как Read Committed. Поэтому воркер не «подсмотрит» строку из незавершённой транзакции, и строить на этом расчёт нельзя.

Автоматический retry иногда замаскирует гонку: следующая попытка запустится после коммита и найдёт товар. Но корректность не должна зависеть от того, успел ли воркер стартовать и как он разобрал ошибку not found. Кроме того, если бизнес-транзакция откатится, задача во внешнем брокере всё равно останется.

River хранит задачи в PostgreSQL, поэтому запись в products и запись в river_job можно сделать через один *sql.Tx. Если транзакция откатилась, не останется ни товара, ни задачи. Если коммит прошёл, задача уже находится в базе и переживёт остановку API. Это и есть transactional enqueueing.

До коммита строка river_job не видна другим транзакциям по тем же правилам PostgreSQL. Поэтому River не сможет забрать задачу раньше, чем станут видны связанные с ней бизнес-данные.

POST /api/v1/products                    worker
        │                                  │
        │  одна короткая транзакция        │
        ├─ INSERT INTO products            │
        └─ INSERT INTO river_job ─────────►│
                                           │
                                  Reserve(product_id)
                                  5–10 секунд, без транзакции
                                           │
                                  UPDATE products

Запуск и проверка примера

В демонстрационном проекте Docker Compose поднимает PostgreSQL, миграции, API, воркер и River UI:

docker compose up -d --build
curl http://localhost:8100/healthz

Ожидаемый ответ health check:

{"status":"ok"}

Создадим товар:

curl -s -X POST http://localhost:8100/api/v1/products \
  -H "Content-Type: application/json" \
  -d '{"name":"Mechanical keyboard","description":"Tactile switches","price_cents":7990}'

API вернёт 201 Created. Сразу после создания поля in_stock и published_at равны null: медленная часть работы уже стоит в очереди, но HTTP-запрос её не ждёт. Сокращённо ответ выглядит так:

{
  "id": 1,
  "name": "Mechanical keyboard",
  "price_cents": 7990,
  "in_stock": null,
  "published_at": null
}

Через 5–10 секунд повторный запрос покажет зарезервированный остаток и время публикации:

curl -s http://localhost:8100/api/v1/products/1
docker compose logs worker | grep river

Очередь также доступна в River UI. Там видно аргументы, состояние и историю попыток задачи.

После успешного выполнения задача появится в разделе Completed. В карточке задачи можно проверить итоговое состояние, аргументы и количество выполненных попыток.

Архитектура приложения

Почему не выполнить всё в HTTP-обработчике

Фоновая задача полезна не потому, что Go не умеет параллельно обрабатывать запросы. Gin работает поверх net/http, и один медленный запрос не блокирует весь сервер.

Проблема возникает, когда внешний вызов делают внутри бизнес-транзакции. Пока сервис склада отвечает 5–10 секунд, транзакция удерживает соединение из пула. В примере разрешено десять открытых соединений: десять одновременных публикаций могут занять весь пул, после чего даже быстрые чтения начнут ждать.

В примере обязанности разделены так:

  • HTTP-обработчик быстро сохраняет товар и задачу;
  • воркер ограничивает число одновременных обращений к сервису склада;
  • медленный внешний вызов не удерживает транзакцию;
  • River повторяет задачу после временной ошибки.

Сам по себе перенос кода в воркер не делает его надёжным. Воркер всё равно должен укладываться в таймауты, возвращать ошибки наверх и не ломаться, если задача выполнится повторно. River подробно разбирает эти требования в руководстве Writing reliable workers.

Как читать код примера

В проекте много каталогов, но основной путь короткий:

HTTP handler
    │
    ▼
ProductUseCase.Create
    │
    ├── productRepo.Create
    └── PublishProductEnqueuer.PublishProduct

River worker
    │
    ▼
ProductUseCase.Publish
    │
    ├── warehouse.Reserve
    └── productRepo.MarkPublished

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

  1. internal/domain/product/usecases/product/product_usecase.go: вся бизнес-последовательность.
  2. internal/infrastructure/river/product/jobs/publish_product.go: контракт аргументов задачи.
  3. internal/infrastructure/river/product/enqueuer/publish_product.go: соединение GORM-транзакции с River.
  4. internal/infrastructure/river/product/worker/publish_product.go: классификация ошибок.
  5. pkg/transaction/transaction.go: передача текущей GORM-транзакции через контекст.
  6. internal/app/container.go и internal/app/worker.go: сборка зависимостей и запуск процессов.

В таком порядке сначала видна транзакция и путь задачи, а уже потом HTTP DTO, конфигурация и Docker-инфраструктура.

Аргументы задачи и Kind

У каждого типа задачи River есть структура аргументов и метод Kind():

type PublishProductArgs struct {
	ProductID int64 `json:"product_id"`
}

func (PublishProductArgs) Kind() string { return "publish_product" }

В очередь попадает только идентификатор товара, а не вся модель. Аргументы сериализуются в JSON и могут пролежать в river_job долго, а состояние товара за это время изменится. Поэтому воркер читает актуальные данные из PostgreSQL сам, а маленький payload заодно проще хранить и версионировать.

Имя из Kind() тоже часть контракта, который лежит в базе. Переименовать Go-тип не страшно, а вот менять саму строку "publish_product" опасно: задачи со старым именем уже стоят в очереди, и для них придётся отдельно продумать переход.

Общая транзакция PostgreSQL и River

Use case открывает транзакцию через GORM, создаёт товар и с тем же контекстом вызывает постановщика задачи:

func (u *ProductUseCase) Create(
	ctx context.Context,
	input product.ProductCreateInput,
) (models.Product, error) {
	var created models.Product

	err := u.tm.InTransaction(ctx, func(ctx context.Context) error {
		item, err := u.productRepo.Create(ctx, input)
		if err != nil {
			return fmt.Errorf("create product: %w", err)
		}

		if err := u.jobEnqueuer.PublishProduct(ctx, item.ID); err != nil {
			return fmt.Errorf("enqueue publish job: %w", err)
		}

		created = item
		return nil
	})
	if err != nil {
		return models.Product{}, err
	}

	return created, nil
}

Use case не импортирует River. Он знает только локальный интерфейс jobEnqueuer, поэтому решение о конкретной очереди остаётся в инфраструктурном слое.

Менеджер транзакций кладёт *gorm.DB в производный контекст. Репозиторий достаёт его через ExtractDB, а адаптер River через ExtractTx. Затем адаптер извлекает обёрнутый *sql.Tx и передаёт его в InsertTx:

func (e *PublishProductEnqueuer) PublishProduct(
	ctx context.Context,
	productID int64,
) error {
	gormTx, err := transaction.ExtractTx(ctx)
	if err != nil {
		return err
	}

	sqlTx, ok := gormTx.Statement.ConnPool.(*sql.Tx)
	if !ok {
		return errors.New("enqueuer: failed to extract sql transaction from gorm")
	}

	_, err = e.client.InsertTx(
		ctx,
		sqlTx,
		jobs.PublishProductArgs{ProductID: productID},
		nil,
	)
	return err
}

Транзакцию здесь обеспечивает PostgreSQL, а GORM остаётся обычной ORM и просто отдаёт приложению *sql.Tx. В этом примере River и GORM работают с одной транзакцией database/sql. Подробности такой интеграции есть в руководстве Using River with GORM.

После этого возможны только два итоговых состояния:

COMMIT   → товар существует, задача существует
ROLLBACK → товара нет, задачи нет

Отдельные клиенты для API и воркера

API создаёт insert-only клиент без зарегистрированных воркеров:

riverEnqueueClient, err := river.NewClient(
	riverdatabasesql.New(sqlDB),
	&river.Config{},
)

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

Отдельный процесс регистрирует воркер и задаёт параллелизм:

workers := river.NewWorkers()
river.AddWorker(
	workers,
	worker.NewPublishProductWorker(container.ProductUsecase),
)

client, err := river.NewClient(
	riverdatabasesql.New(container.DB),
	&river.Config{
		Workers: workers,
		Queues: map[string]river.QueueConfig{
			river.QueueDefault: {MaxWorkers: 10},
		},
		SoftStopTimeout: 30 * time.Second,
	},
)

MaxWorkers ограничивает нагрузку на сервис склада и PostgreSQL. SoftStopTimeout даёт выполняющимся задачам время завершиться после SIGTERM; затем их контексты отменяются. Поэтому любой внешний вызов внутри воркера должен прерываться по ctx.Done(), а не продолжать работу в фоне.

За два процесса приходится платить вторым composition root: конфиг, пул соединений, миграции и graceful shutdown собираются дважды. Взамен число обработчиков очереди перестаёт зависеть от числа реплик API. Если бы клиент с воркерами жил внутри API, пять инстансов под HTTP-нагрузкой превратились бы в пятьдесят одновременных обращений к сервису склада, хотя его пропускная способность к объёму HTTP-трафика отношения не имеет.

Надёжность воркера

Идемпотентность при повторном выполнении

River использует семантику at least once: одна и та же задача может быть выполнена повторно. Например, воркер успел изменить внешнюю систему, но упал раньше, чем River записал успех.

В примере Publish сначала проверяет PublishedAt, затем обращается к сервису склада и выполняет условный UPDATE:

func (u *ProductUseCase) Publish(ctx context.Context, productID int64) error {
	item, err := u.productRepo.Get(ctx, productID)
	if err != nil {
		return err
	}

	if item.PublishedAt != nil {
		return nil
	}

	inStock, err := u.warehouse.Reserve(ctx, productID)
	if err != nil {
		return err
	}

	_, err = u.productRepo.MarkPublished(ctx, productID, inStock)
	return err
}

Репозиторий обновляет только ещё не опубликованную строку:

result := db.Model(&models.Product{}).
	Where("id = ? AND published_at IS NULL", id).
	Updates(map[string]any{
		"in_stock":     inStock,
		"published_at": time.Now(),
	})

Условие в SQL защищает состояние от гонки, но оно не отменяет повторный вызов сервиса склада: если процесс упал после Reserve, следующий запуск обратится к нему снова.

Поэтому реальный сервис склада должен принимать идемпотентный ключ, например productID или отдельный reservationID, и возвращать прежний результат для повтора. Unique jobs защищают от повторной постановки, но не меняют семантику повторного выполнения.

Какие ошибки повторять, а какие отменять

River определяет дальнейшее состояние задачи по ошибке, возвращённой из Work:

  • nil: задача выполнена;
  • обычная ошибка: временный сбой, River запланирует повтор;
  • river.JobCancel(err): постоянная ошибка, повторы бессмысленны;
  • river.JobSnooze(duration): задачу нужно отложить, и это не считается ошибкой.

Воркер остаётся тонким адаптером:

func (w *PublishProductWorker) Work(
	ctx context.Context,
	job *river.Job[jobs.PublishProductArgs],
) error {
	err := w.productPublisher.Publish(ctx, job.Args.ProductID)
	if err == nil {
		return nil
	}

	if errors.Is(err, product.ErrProductNotFound) {
		return river.JobCancel(fmt.Errorf(
			"publish product %d: %w",
			job.Args.ProductID,
			err,
		))
	}

	return fmt.Errorf("publish product %d: %w", job.Args.ProductID, err)
}

Не стоит логировать ошибку и возвращать nil: River решит, что работа выполнена, и повтор потеряется. Логи полезны для диагностики, но именно возвращаемое значение управляет жизненным циклом задачи.

Миграции и диагностика

River добавляет в PostgreSQL собственные таблицы, включая river_job, river_queue и river_migration. В примере их миграции лежат рядом с миграцией products и применяются одним шагом через golang-migrate.

Как получить миграции River

River поставляет SQL-миграции внутри Go-модуля драйвера riverdatabasesql. В проекте их извлекает скрипт scripts/fetch_river_migrations.sh:

docker compose run --rm --entrypoint "" migrate \
  ./scripts/fetch_river_migrations.sh

Локальный Go toolchain тоже подходит:

go mod download
./scripts/fetch_river_migrations.sh

После выполнения в db/migrations появятся пары *.up.sql и *.down.sql. Перед коммитом их нужно прочитать и проверить обычным diff, как любые другие изменения схемы.

Скрипт выполняет четыре действия:

  1. Через go list -m находит каталог riverdatabasesql в Go module cache.
  2. Берёт миграции именно той версии River, которая зафиксирована в go.mod.
  3. Копирует up- и down-файлы в общий каталог db/migrations.
  4. Присваивает каждому файлу уникальную числовую версию в формате golang-migrate.

Ручное копирование легко рассинхронизировать с библиотекой: код уже обновлён, а SQL остался от прежней версии. Скрипт фиксирует, откуда взят SQL и какой он версии, держит River-схему рядом с бизнес-схемой и позволяет применять всё одной командой:

migrate -path db/migrations -database "$DATABASE_URL" up

Альтернатива: river migrate-up. Она подходит, но тогда схема River и схема приложения управляются двумя разными командами. В примере выбран единый migration pipeline, чтобы порядок изменений был виден в репозитории и одинаково воспроизводился локально, в CI и при деплое.

Скрипт рассчитан прежде всего на первоначальный импорт. При обновлении River нельзя удалять или переименовывать миграции, уже применённые в production. Сохраняйте их версии и добавляйте только новые River-миграции; иначе golang-migrate воспримет старый SQL под новым номером как ещё не выполненный.

Как наблюдать за очередью

Для повседневной диагностики достаточно River UI. SQL-запросы к river_job полезны в CI и при разборе инцидента:

SELECT state, count(*) AS jobs, max(attempt) AS max_attempt
FROM river_job
GROUP BY state
ORDER BY state;

Ручное изменение строк очереди через UPDATE лучше оставить учебным экспериментом. Для retry, cancel и управления очередями безопаснее использовать River UI или публичный API библиотеки: они соблюдают правила переходов между состояниями, а разовый SQL-скрипт про эти правила ничего не знает.

Проверка и ограничения

Что проверяют тесты

Тест с rivertest.RequireInsertedTx ищет задачу через тот же *sql.Tx, внутри которого был создан товар.

job := rivertest.RequireInsertedTx[*riverdatabasesql.Driver](
	ctx,
	t,
	sqlTx,
	jobs.PublishProductArgs{},
	nil,
)

assert.Equal(t, created.ID, job.Args.ProductID)

Кроме этого, тесты проходят через реальные Gin-обработчики, GORM-репозиторий и PostgreSQL. Вместо сервиса склада подставлена локальная заглушка с нулевой задержкой, чтобы прогон занимал миллисекунды.

docker compose run --rm tests

Отдельный тест должен смоделировать ошибку постановки задачи после успешного INSERT products и убедиться, что GORM откатил обе записи. Отмена контекста до начала транзакции проверяет ранний выход, но не этот путь.

Что остаётся добавить для production

Код намеренно небольшой, поэтому перед production понадобятся дополнительные решения:

  • явная обработка ситуации, когда условный UPDATE не изменил ни одной строки;
  • политика для задач, исчерпавших все попытки;
  • подбор числа воркеров, таймаутов и размера пула по измерениям;
  • корректное закрытие соединений с базой при остановке процессов;
  • проверка стратегии миграций и выбранного драйвера River под нагрузкой.

Когда выбирать River

River подходит, когда задача прямо следует из изменения данных в БД: отправить уведомление после создания заказа, перестроить индекс после обновления документа, подготовить файл после сохранения отчёта. Общее у этих случаев то, что задачу нельзя потерять, а её появление привязано к успешной записи.

Отдельная очередь, в том числе на Redis, остаётся разумным выбором, если она уже есть в платформе или архитектуре нужен самостоятельный брокер.

Атомарная связь с SQL-транзакцией не обязательна, когда:

  • постановка задачи не сопровождается изменением бизнес-данных;
  • задача запускается по расписанию: очистка, прогрев кеша или периодический отчёт;
  • все необходимые данные уже были закоммичены до постановки задачи;
  • аргументы самодостаточны и воркеру не нужно читать только что созданную строку;
  • между записью в базу и постановкой задачи нет инварианта «обе операции или ни одной».

Если общая транзакция не нужна, это ещё не значит, что задачу можно терять. Надёжность хранения, повторы и идемпотентность всё равно нужны. Просто их не требуется атомарно связывать с конкретной записью в PostgreSQL.

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