River: a Postgres Job Queue for Go
River is a Go job queue that stores jobs in Postgres and inserts them inside your application's transactions, and this guide sets it up end to end as part of Task Queues in Go in Backend Frameworks & Worker Scaling. It is the Go embodiment of the design in Database-Backed Job Queues: SKIP LOCKED claiming, LISTEN/NOTIFY wake-ups, and no extra infrastructure beyond the database you already run.
Problem Statement
A Go billing service creates invoices in Postgres and enqueues "charge invoice" and "email invoice" jobs to Redis. Twice this quarter an invoice was created but never charged, because the process crashed between the database commit and the Redis enqueue. The service already uses pgx and has spare database capacity; it runs about 40 jobs per second at peak. You want the job to exist if and only if the invoice commits, retries with backoff for transient payment-provider errors, at most one pending charge job per invoice, and a way to delay a job without counting it as a failure when the provider asks you to come back later.
Prerequisites
- Go 1.22+, Postgres 13+ (River supports current and recent Postgres versions), and
github.com/jackc/pgx/v5. go get github.com/riverqueue/river github.com/riverqueue/river/riverdriver/riverpgxv5.- The River CLI for migrations:
go install github.com/riverqueue/river/cmd/river@latest, or River'srivermigratepackage to run them from your own migration tooling. - A connection budget that can absorb the worker pool (Step 4).
Step 1 — Run River's Migrations
River owns a small set of tables (river_job, river_leader, river_queue, and others). Apply them like any schema change, ideally through your existing migration pipeline so they are versioned with the application.
# One-off or in CI: apply all River migrations to the target database
river migrate-up --database-url "$DATABASE_URL"
# Preview what would change after upgrading the library
river migrate-get --up --version 6 --database-url "$DATABASE_URL"
// Or from Go, inside your own migration step
migrator, err := rivermigrate.New(riverpgxv5.New(pool), nil)
if err != nil {
return err
}
if _, err := migrator.Migrate(ctx, rivermigrate.DirectionUp, nil); err != nil {
return err
}
Run migrations before deploying a River version that needs them; the library checks the schema version and refuses to start against an outdated schema, which is the behaviour you want during rolling deploys.
Step 2 — Define Job Args and Workers
Each job kind is a struct implementing Kind(), and a worker type that embeds river.WorkerDefaults and implements Work. Optional methods on the args set per-kind insert defaults.
// internal/jobs/charge.go
type ChargeInvoiceArgs struct {
InvoiceID int64 `json:"invoice_id"`
}
func (ChargeInvoiceArgs) Kind() string { return "charge_invoice" }
// Per-kind defaults applied on every insert of this kind
func (ChargeInvoiceArgs) InsertOpts() river.InsertOpts {
return river.InsertOpts{
Queue: "billing",
MaxAttempts: 12,
UniqueOpts: river.UniqueOpts{
ByArgs: true, // same invoice_id...
ByState: []rivertype.JobState{rivertype.JobStateAvailable,
rivertype.JobStateRunning, rivertype.JobStateRetryable,
rivertype.JobStateScheduled}, // ...while not finished
},
}
}
type ChargeInvoiceWorker struct {
river.WorkerDefaults[ChargeInvoiceArgs]
Payments *payments.Client
DB *pgxpool.Pool
}
func (w *ChargeInvoiceWorker) Timeout(*river.Job[ChargeInvoiceArgs]) time.Duration {
return 45 * time.Second
}
The UniqueOpts above say "at most one charge job per invoice in any non-final state". River enforces it with a unique index, so concurrent inserts of the same job cannot both succeed — the database constraint does the work that a Redis lock with a TTL does elsewhere.
Step 3 — Implement Work with Errors, Cancels, and Snoozes
Returning an error schedules a retry with River's default backoff (roughly attempt^4 seconds with jitter). Two special return values change that: river.JobCancel(err) stops retrying permanently, and river.JobSnooze(d) reschedules the job after d without consuming an attempt.
func (w *ChargeInvoiceWorker) Work(ctx context.Context, job *river.Job[ChargeInvoiceArgs]) error {
inv, err := loadInvoice(ctx, w.DB, job.Args.InvoiceID)
if errors.Is(err, pgx.ErrNoRows) {
return river.JobCancel(fmt.Errorf("invoice %d not found", job.Args.InvoiceID))
}
if err != nil {
return err // transient DB error: retry
}
if inv.Status == "paid" {
return nil // idempotent: already charged
}
res, err := w.Payments.Charge(ctx, payments.ChargeRequest{
Amount: inv.AmountCents,
Customer: inv.CustomerRef,
IdempotencyKey: fmt.Sprintf("invoice-%d", inv.ID), // provider-side dedup
})
var rl *payments.RateLimitedError
switch {
case errors.As(err, &rl):
return river.JobSnooze(rl.RetryAfter) // come back later, no attempt used
case errors.Is(err, payments.ErrCardDeclined):
return river.JobCancel(err) // business outcome, not a failure
case err != nil:
return err // network/5xx: retry with backoff
}
return markPaid(ctx, w.DB, inv.ID, res.ChargeID)
}
Snoozing is the right answer to rate limiting: a provider saying "retry after 30 seconds" is not a failure, and spending attempts on it would exhaust MaxAttempts during a traffic spike. The provider idempotency key makes a retry after a crash safe even if the first charge went through, as covered in preventing duplicate job execution with idempotency.
The three outcomes are worth keeping distinct in dashboards as well as in code. Retries measure unhealthy dependencies; snoozes measure back-pressure from a healthy dependency that is simply busy; cancels measure business outcomes such as declined cards. Lumping them into one "failed jobs" graph hides which of the three is growing, and each one calls for a different response.
Step 4 — Configure the Client and Size Its Pool
The same river.Client type is used for inserting (in API processes) and working (in worker processes). A client with no Queues configured can insert but never works jobs, which is the right setup for API pods.
func newWorkerClient(pool *pgxpool.Pool, deps Deps) (*river.Client[pgx.Tx], error) {
workers := river.NewWorkers()
river.AddWorker(workers, &ChargeInvoiceWorker{Payments: deps.Payments, DB: pool})
river.AddWorker(workers, &EmailInvoiceWorker{Mailer: deps.Mailer})
return river.NewClient(riverpgxv5.New(pool), &river.Config{
Queues: map[string]river.QueueConfig{
"billing": {MaxWorkers: 10}, // limited by payment-provider concurrency
river.QueueDefault: {MaxWorkers: 40}, // emails and light work
},
Workers: workers,
FetchCooldown: 100 * time.Millisecond,
FetchPollInterval: 2 * time.Second, // fallback poll if a notification is missed
RescueStuckJobsAfter: 5 * time.Minute, // return jobs from crashed workers
JobTimeout: time.Minute, // default when a worker sets no Timeout
})
}
// pgxpool: workers + River's own connections (listener, leader election, completer)
cfg, _ := pgxpool.ParseConfig(os.Getenv("DATABASE_URL"))
cfg.MaxConns = 60 // >= sum(MaxWorkers) + ~10 headroom
If handlers also query the database, MaxConns must cover both River's overhead and the handlers' concurrent queries. Multiply by worker pod count and compare against max_connections before scaling out — the same arithmetic shown for Rails in running Rails Solid Queue in production.
Step 5 — Insert Jobs Transactionally
InsertTx takes a pgx.Tx, so the job row is part of the invoice transaction.
func CreateInvoice(ctx context.Context, pool *pgxpool.Pool, rc *river.Client[pgx.Tx], in NewInvoice) (int64, error) {
tx, err := pool.Begin(ctx)
if err != nil {
return 0, err
}
defer tx.Rollback(ctx) // no-op after commit
var id int64
if err := tx.QueryRow(ctx,
`INSERT INTO invoices (customer_ref, amount_cents, status) VALUES ($1, $2, 'open') RETURNING id`,
in.CustomerRef, in.AmountCents).Scan(&id); err != nil {
return 0, err
}
if _, err := rc.InsertManyTx(ctx, tx, []river.InsertManyParams{
{Args: ChargeInvoiceArgs{InvoiceID: id}},
{Args: EmailInvoiceArgs{InvoiceID: id}, InsertOpts: &river.InsertOpts{
ScheduledAt: time.Now().Add(2 * time.Minute)}}, // after the charge normally lands
}); err != nil {
return 0, err
}
return id, tx.Commit(ctx)
}
If the commit fails, neither the invoice nor the jobs exist. If it succeeds, a NOTIFY fires on commit and an idle worker picks up the charge within milliseconds. This is the property the Redis-based setup could not provide without a transactional outbox.
Step 6 — Start, Stop, and Observe
Start the client in worker processes, and stop it gracefully on SIGTERM: Stop lets running jobs finish, and StopAndCancel cancels their contexts if they overrun.
if err := client.Start(ctx); err != nil {
log.Fatal(err)
}
sig := make(chan os.Signal, 1)
signal.Notify(sig, syscall.SIGINT, syscall.SIGTERM)
<-sig
softCtx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
defer cancel()
if err := client.Stop(softCtx); err != nil { // wait for running jobs
hardCtx, cancel2 := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel2()
_ = client.StopAndCancel(hardCtx) // cancel stragglers' contexts
}
For visibility, River UI (a separate binary or embeddable handler) lists jobs by state, shows errors per attempt, and lets operators retry or cancel. For metrics, query river_job directly or use River's OpenTelemetry middleware.
-- Backlog and oldest available job per queue
SELECT queue, count(*) AS available,
extract(epoch FROM now() - min(scheduled_at)) AS oldest_s
FROM river_job WHERE state = 'available' GROUP BY queue;
Verification
func TestInvoiceAndJobsCommitTogether(t *testing.T) {
ctx := context.Background()
// Rollback path: force an error after inserting jobs
_, err := createInvoiceThenFail(ctx, pool, riverClient)
require.Error(t, err)
rivertest.RequireNotInserted(ctx, t, riverpgxv5.New(pool), &ChargeInvoiceArgs{}, nil)
// Commit path
id, err := CreateInvoice(ctx, pool, riverClient, NewInvoice{CustomerRef: "c1", AmountCents: 1000})
require.NoError(t, err)
job := rivertest.RequireInserted(ctx, t, riverpgxv5.New(pool), &ChargeInvoiceArgs{}, nil)
require.Equal(t, id, job.Args.InvoiceID)
}
In staging, kill a worker pod mid-job and confirm the job reappears as available after RescueStuckJobsAfter, then completes once.
Gotchas & Edge Cases
Timeouts longer than the rescue window. If a job's timeout exceeds RescueStuckJobsAfter, the rescuer can return a still-running job to available and a second worker starts it. Keep RescueStuckJobsAfter well above your longest job timeout.
PgBouncer in transaction mode. River's listener needs a session-level connection for LISTEN. Point the River pool at a session-mode pool or directly at Postgres, or disable the notifier and rely on polling — details in Postgres LISTEN/NOTIFY for job wake-ups.
Completed job retention. River keeps finished jobs for a configurable period (24 hours for completed by default) and its maintenance process deletes them. Shorten retention for high-volume kinds to limit table size and vacuum load.
Large args. Args are stored as JSON in each row. Keep them to ids; fetch data inside the worker.
FAQ
How does River compare with Asynq? River gives transactional enqueue and database-grade durability with no extra infrastructure; Asynq gives higher throughput and lower latency on Redis. See choosing a Go job queue library.
Can I run River against a read replica?
No. Claiming, completing, and inserting jobs are writes, and LISTEN/NOTIFY works only on the primary.
What throughput can River sustain?
Thousands of jobs per second on a well-provisioned Postgres, bounded by write capacity and vacuum. Batch inserts with InsertMany and short completed-job retention help at the high end.
Related
- Task Queues in Go — River in the context of other Go options.
- Getting Started with Asynq in Go — the Redis-backed alternative.
- Graceful Shutdown for Go Workers — signal handling and grace periods.
- Building a Postgres Job Queue with SKIP LOCKED — the mechanics River automates.