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's rivermigrate package 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.

River job states A job is inserted as available, or scheduled if it has a future run time. A worker moves it to running. Success moves it to completed. An error moves it to retryable with a backoff time, after which it becomes available again. Exhausting max attempts or a cancel error moves it to discarded or cancelled. Unique opts cover every non-final state available running completed retryable discarded scheduled Retryable and scheduled jobs return to available when their time arrives.

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.

Error, snooze, or cancel A plain error from Work schedules a retry with exponential backoff and counts against MaxAttempts. JobSnooze reschedules the job after the given duration without using an attempt, which suits rate-limit responses. JobCancel moves the job to cancelled immediately, which suits business outcomes like a declined card. What Work returns decides what happens next return err retry with backoff uses one attempt river.JobSnooze(d) run again after d attempt not counted river.JobCancel(err) stop permanently state: cancelled Graph the three separately: they signal different problems.

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.

Separate goroutine budgets per queue One River worker process runs two queues. The billing queue has a limit of ten concurrent jobs, matching the payment provider's concurrency allowance. The default queue has forty. A burst of email jobs can fill the default budget without taking any slots from billing, and a billing backlog cannot starve emails. One process, two isolated budgets worker pod billing MaxWorkers 10 default MaxWorkers 40: emails, PDFs, webhooks Set billing's limit from the provider's allowance, not from CPU.

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