Task Queues in Go

Go services have a different relationship with background work than Python or Ruby services, and this guide covers the Go job-queue landscape as part of Backend Frameworks & Worker Scaling. Goroutines make in-process concurrency cheap, so the temptation is to go doWork() and move on — which works until the process restarts and every in-flight goroutine vanishes with it. A durable task queue is what turns "fire a goroutine" into "this work will happen, even across deploys and crashes".

The two libraries most Go teams reach for are Asynq, a Redis-backed queue modelled on Sidekiq, and River, a Postgres-backed queue built on SKIP LOCKED with transactional enqueue. Both give you typed task payloads, retries with backoff, scheduling, and a web UI; they differ in where the queue lives and what that implies for consistency and operations.

The Scenario: Goroutines That Disappear

A Go API handles signups. After inserting the user, the handler starts a goroutine to send a welcome email and another to provision a workspace. It is fast and simple. Then the team adopts rolling deploys on Kubernetes. Every deploy sends SIGTERM to old pods, the http.Server shuts down gracefully — and the process exits, taking any running email and provisioning goroutines with it. About 0.3% of signups never get a workspace. Nobody notices for weeks, because there is no record that the work was ever supposed to happen.

The fix is not better goroutine management; it is recording the work durably before acknowledging the request, then executing it from that record with retries. That is exactly what a task queue provides.

Fire-and-forget vs durable enqueue In the top path, the HTTP handler starts a goroutine to provision a workspace; when the pod receives SIGTERM and exits, the goroutine is killed and the work is lost with no record. In the bottom path, the handler enqueues a durable task in Redis or Postgres; if the worker processing it dies, the task is retried by another worker. What happens to the work on SIGTERM? handler go provision(user) pod exits: work lost handler enqueue task (durable) any worker runs or retries it A goroutine is concurrency, not durability. The queue record is what survives the process.

Architectural Overview: Client, Server, and Handlers

Both Asynq and River split the system the same way. A client enqueues typed tasks from any process — usually the API. A server (Asynq) or client in worker mode (River) runs in worker processes, pulls tasks, and dispatches them to handlers registered by task type. Concurrency is a pool of goroutines per process, and each handler receives a context.Context that is cancelled on timeout or shutdown.

// The shape shared by Go job libraries: typed args, a handler per kind
type WelcomeEmailArgs struct {
    UserID int64  `json:"user_id"`
    Locale string `json:"locale"`
}

// River: Kind() names the job; a worker type handles it
func (WelcomeEmailArgs) Kind() string { return "welcome_email" }

type WelcomeEmailWorker struct {
    river.WorkerDefaults[WelcomeEmailArgs]
    Mailer *mail.Client
}

func (w *WelcomeEmailWorker) Work(ctx context.Context, job *river.Job[WelcomeEmailArgs]) error {
    // ctx is cancelled if the job exceeds its timeout or the client is stopping
    return w.Mailer.SendWelcome(ctx, job.Args.UserID, job.Args.Locale)
}

The type-safety is a real advantage over dynamic-language queues: a payload mismatch between producer and handler is a compile error when both share the args struct, instead of a runtime KeyError in production. It also means payload evolution needs the same care as any JSON contract — add fields as optional, never rename — as covered in versioning job payload schemas.

Client, store, worker pool API processes use the library client to enqueue typed tasks into the backing store, Redis for Asynq or Postgres for River. Worker processes each run a fixed pool of goroutines that fetch tasks and dispatch them to the handler registered for the task kind, passing a context that is cancelled on timeout or shutdown. Go job system layout API pod: client cron pod: client Redis (Asynq) or Postgres (River) worker pod goroutine pool (N) handler by kind ctx cancelled on stop Producers and workers share the args structs, so payload mismatches fail at compile time.

Implementation 1: Asynq on Redis

Asynq stores tasks in Redis lists and sorted sets, uses Lua scripts for atomic state transitions, and supports weighted queue priorities, unique tasks, scheduled tasks, and aggregation of small tasks into groups. A complete worker is short:

package main

import (
    "context"
    "encoding/json"
    "log"
    "time"

    "github.com/hibiken/asynq"
)

const TypeWelcomeEmail = "email:welcome"

func NewWelcomeEmailTask(userID int64) (*asynq.Task, error) {
    payload, err := json.Marshal(map[string]int64{"user_id": userID})
    if err != nil {
        return nil, err
    }
    return asynq.NewTask(TypeWelcomeEmail, payload,
        asynq.MaxRetry(10),
        asynq.Timeout(30*time.Second),        // ctx deadline inside the handler
        asynq.Queue("default")), nil
}

func handleWelcomeEmail(ctx context.Context, t *asynq.Task) error {
    var p struct{ UserID int64 `json:"user_id"` }
    if err := json.Unmarshal(t.Payload(), &p); err != nil {
        return fmt.Errorf("bad payload: %v: %w", err, asynq.SkipRetry) // poison: don't retry
    }
    return mailer.SendWelcome(ctx, p.UserID)
}

func main() {
    srv := asynq.NewServer(
        asynq.RedisClientOpt{Addr: "redis:6379"},
        asynq.Config{
            Concurrency: 20,                                   // goroutines per process
            Queues:      map[string]int{"critical": 6, "default": 3, "low": 1},
            ShutdownTimeout: 25 * time.Second,                 // < k8s terminationGracePeriod
        },
    )
    mux := asynq.NewServeMux()
    mux.HandleFunc(TypeWelcomeEmail, handleWelcomeEmail)
    if err := srv.Run(mux); err != nil {                      // blocks; handles SIGTERM
        log.Fatal(err)
    }
}

The Queues map implements weighted priority: critical is polled six times as often as low, so low-priority work still makes progress. Strict priority is available with StrictPriority: true, at the risk of starvation. A full walkthrough, including the Asynqmon UI and task inspection, is in getting started with Asynq in Go.

Implementation 2: River on Postgres

River keeps jobs in a Postgres table and can insert them inside the caller's transaction, which closes the dual-write gap between a business write and an enqueue. It uses LISTEN/NOTIFY for low-latency pickup, with polling as a fallback.

package main

import (
    "context"
    "log"

    "github.com/jackc/pgx/v5"
    "github.com/jackc/pgx/v5/pgxpool"
    "github.com/riverqueue/river"
    "github.com/riverqueue/river/riverdriver/riverpgxv5"
)

func main() {
    ctx := context.Background()
    pool, err := pgxpool.New(ctx, "postgres://app@pg:5432/app")
    if err != nil {
        log.Fatal(err)
    }

    workers := river.NewWorkers()
    river.AddWorker(workers, &WelcomeEmailWorker{Mailer: mailer})

    client, err := river.NewClient(riverpgxv5.New(pool), &river.Config{
        Queues: map[string]river.QueueConfig{
            river.QueueDefault: {MaxWorkers: 50},     // goroutines for this queue
            "reports":          {MaxWorkers: 4},      // heavy jobs isolated
        },
        Workers: workers,
    })
    if err != nil {
        log.Fatal(err)
    }
    if err := client.Start(ctx); err != nil {
        log.Fatal(err)
    }
    // ... wait for signal, then client.Stop(ctx) — see graceful shutdown below
}

// In the API: enqueue in the same transaction as the signup
func createUser(ctx context.Context, tx pgx.Tx, riverClient *river.Client[pgx.Tx], u User) error {
    if err := insertUser(ctx, tx, u); err != nil {
        return err
    }
    _, err := riverClient.InsertTx(ctx, tx, WelcomeEmailArgs{UserID: u.ID, Locale: u.Locale}, nil)
    return err                                          // commit both or neither
}

River's per-queue MaxWorkers gives each queue its own goroutine budget, which is simpler to reason about than weights. Unique jobs, periodic jobs, and job snoozing are built in. The Postgres-side trade-offs — write load, vacuum, connection counts — are those of any database-backed job queue; River, a Postgres job queue for Go goes deeper.

Scheduled and Periodic Jobs

Both libraries handle the two kinds of time-based work that come up in almost every service: a one-off job that should run later ("send a reminder in 24 hours") and a periodic job that should run on a schedule ("reconcile payments every night at 02:00").

One-off delays are an enqueue option. Asynq stores the task in a Redis sorted set scored by its due time and a forwarder process moves it to the ready list when it is due; River inserts the row with a future scheduled_at and the claim query ignores it until then.

// Asynq: run in 24 hours, or at an absolute time
client.Enqueue(task, asynq.ProcessIn(24*time.Hour))
client.Enqueue(task, asynq.ProcessAt(time.Date(2026, 10, 1, 9, 0, 0, 0, time.UTC)))

// River: the same with InsertOpts
riverClient.Insert(ctx, ReminderArgs{UserID: id}, &river.InsertOpts{
    ScheduledAt: time.Now().Add(24 * time.Hour),
})

Periodic jobs need a single scheduler, or every worker pod enqueues the same nightly job. Asynq provides a separate Scheduler process (run exactly one replica, or use its PeriodicTaskManager with a config provider); River's periodic jobs are enqueued only by the elected leader among running clients, so running many worker pods is safe by default.

// River: periodic job registered on the client; only the leader enqueues it
periodic := []*river.PeriodicJob{
    river.NewPeriodicJob(
        river.PeriodicInterval(15*time.Minute),
        func() (river.JobArgs, *river.InsertOpts) { return RefreshRatesArgs{}, nil },
        &river.PeriodicJobOpts{RunOnStart: true},
    ),
}
client, _ := river.NewClient(riverpgxv5.New(pool), &river.Config{
    Queues: queues, Workers: workers, PeriodicJobs: periodic,
})

Leader election for schedulers is a general problem — the same one described in preventing duplicate scheduled jobs with leader election — and it is worth checking which of your periodic tasks run once per cluster versus once per pod.

Every pod works, one pod schedules Three River worker pods are running. Pod 2 holds leadership and is the only one that enqueues the periodic refresh job every fifteen minutes. All three pods claim and process jobs from the shared queue, including the periodic ones. If pod 2 exits, another pod is elected and takes over scheduling. Scheduling is leader-only; processing is not pod 1: worker pod 2: worker + leader pod 3: worker river_job table refresh_rates every 15 min leader exits: re-election in seconds enqueue

Observability for Go Workers

Go job libraries expose less out of the box than Sidekiq's dashboard or Flower, so plan the metrics you need. Three layers cover it. Queue metrics — depth and oldest-job age per queue — come from the library's inspector (Asynq) or a SQL query over the job table (River). Handler metrics — duration, outcome, and attempt number per job kind — are best emitted from a middleware so every handler gets them without code changes. Tracing propagates the producer's trace context through the task so a slow job can be tied to the request that enqueued it.

// Asynq middleware: duration histogram and outcome counter per task type
var (
    jobDuration = promauto.NewHistogramVec(prometheus.HistogramOpts{
        Name:    "job_duration_seconds",
        Buckets: []float64{.01, .05, .1, .5, 1, 5, 15, 60, 300},
    }, []string{"type"})
    jobResults = promauto.NewCounterVec(prometheus.CounterOpts{
        Name: "job_results_total",
    }, []string{"type", "outcome"})
)

func metricsMiddleware(next asynq.Handler) asynq.Handler {
    return asynq.HandlerFunc(func(ctx context.Context, t *asynq.Task) error {
        start := time.Now()
        err := next.ProcessTask(ctx, t)
        jobDuration.WithLabelValues(t.Type()).Observe(time.Since(start).Seconds())
        outcome := "success"
        if err != nil {
            outcome = "error"
        }
        jobResults.WithLabelValues(t.Type(), outcome).Inc()
        return err
    })
}

mux.Use(metricsMiddleware)

River supports the same pattern through its middleware and hook interfaces, and ships an OpenTelemetry integration. Keep label cardinality low — task type and outcome, never user ids — and pick histogram buckets that bracket your real job durations, as discussed in choosing histogram buckets for job duration.

Testing Handlers and Enqueue Paths

Go's type system catches payload mismatches, but it cannot catch the behaviours that make background jobs fail in production: a handler that is not idempotent, an enqueue that happens outside the business transaction, or a retry classification that retries a permanent error forever. Test those three things explicitly.

Handlers are ordinary functions, so the fastest tests call them directly with a constructed job and a context, and assert on side effects. Idempotency is tested by calling the handler twice with the same job and asserting the side effect happened once. Cancellation is tested by passing an already-cancelled context and asserting the handler returns promptly with context.Canceled.

func TestWelcomeEmailIsIdempotent(t *testing.T) {
    mailer := &fakeMailer{}
    w := &WelcomeEmailWorker{Mailer: mailer}
    job := &river.Job[WelcomeEmailArgs]{JobRow: &rivertype.JobRow{ID: 42},
        Args: WelcomeEmailArgs{UserID: 7, Locale: "en"}}

    require.NoError(t, w.Work(context.Background(), job))
    require.NoError(t, w.Work(context.Background(), job))   // redelivery
    require.Equal(t, 1, mailer.SentTo(7))                   // sent once
}

func TestWelcomeEmailHonoursCancellation(t *testing.T) {
    ctx, cancel := context.WithCancel(context.Background())
    cancel()
    err := (&WelcomeEmailWorker{Mailer: slowMailer()}).Work(ctx, job)
    require.ErrorIs(t, err, context.Canceled)
}

For the enqueue path, River's rivertest package can assert that a job was inserted within a transaction, and that it was not inserted when the transaction rolled back — the property that motivated choosing River in the first place. For Asynq, run a real Redis in a container and use the Inspector to list pending tasks after the code under test runs. The broader testing strategy, including crash-and-redelivery tests, is covered under testing background jobs.

Trade-off Analysis

Concern Asynq (Redis) River (Postgres) Plain channels / goroutines
Durability Depends on Redis AOF/RDB config Full database durability None
Transactional enqueue with business data No (needs an outbox) Yes (InsertTx) No
Throughput ceiling Very high (Redis) Moderate (database writes) Highest, but not durable
Pickup latency ~1 ms ~5–20 ms with LISTEN/NOTIFY Immediate
Priority model Weighted or strict queues Per-queue worker limits, job priority Whatever you build
Uniqueness asynq.Unique(ttl) with Redis lock Unique opts backed by DB constraint Manual
Web UI Asynqmon River UI None
Extra infrastructure Redis None beyond Postgres None

The deciding factor is usually where your system of record lives. If the work is a consequence of a Postgres write — most application jobs — River's transactional enqueue removes a class of bugs outright. If jobs are high-volume and loosely coupled to database state, or Redis is already a first-class component, Asynq's throughput and lower latency win. Choosing a Go job queue library scores the options, including Machinery and cloud queues, against concrete requirements.

Failure Modes & Recovery

Handlers that ignore context. A handler that calls an HTTP API without passing ctx keeps running after the job's timeout and after shutdown has begun. The library considers the job failed or abandoned and may retry it while the original is still in flight — a duplicate. Remediation: thread ctx into every I/O call (http.NewRequestWithContext, db.QueryContext), and lint for calls without it.

Panics in handlers. Both libraries recover panics in handlers and record them as failures, but a panic in a goroutine spawned by a handler crashes the whole worker process, taking every in-flight job with it. Remediation: never start unmanaged goroutines inside handlers; use errgroup.WithContext so errors and cancellation propagate.

Shutdown shorter than the longest job. Kubernetes sends SIGTERM, waits terminationGracePeriodSeconds (30 s by default), then SIGKILLs. A job still running is killed mid-way. Remediation: set the library's shutdown timeout below the grace period, make long jobs checkpoint and resume, and rely on redelivery — Asynq requeues unfinished tasks on shutdown, and River's rescuer returns stuck jobs. The mechanics are in graceful shutdown for Go workers.

Redis eviction. Asynq on a Redis instance with an eviction policy other than noeviction can silently lose tasks under memory pressure. Remediation: dedicated Redis with maxmemory-policy noeviction, as covered in Redis maxmemory policy for queues.

Context cancellation must reach every call The worker passes a context with a deadline to the handler. The handler passes the same context to the database query and the HTTP call. When the deadline passes or shutdown begins, the context is cancelled and both calls return immediately, so the job stops instead of continuing unobserved. A call made without the context keeps running. One ctx, threaded through every call worker ctx deadline 30s handler(ctx) db.QueryContext(ctx): stops http req with ctx: stops http.Get(url): keeps running A call without ctx outlives the job; the retry then runs alongside it.

Performance Tuning

  • Concurrency per process. Goroutines are cheap, but the resources jobs touch are not. Size Concurrency (Asynq) or MaxWorkers (River) to the scarcest downstream: database connections, an API's rate limit, or memory per job. For I/O-bound jobs, 20–100 goroutines per process is common; for CPU-bound jobs, stay near GOMAXPROCS.
  • Connection pools. River workers use the pgx pool for claiming, completing, and for your job code. Set MaxConns above MaxWorkers plus a few for the library's own use, and account for every pod in the database's connection budget. For Asynq, the Redis pool size defaults to ten per CPU; raise it if handlers also use Redis.
  • Payload size. Keep payloads to ids and small parameters. Large payloads cost Redis memory in Asynq and TOAST rewrites in River.
  • Batching tiny tasks. Asynq's task aggregation groups many small tasks (for example, individual notification events) into one handler call after a delay or size threshold, which cuts per-task overhead dramatically for chatty producers.
  • Separate queues by weight. Put slow, heavy jobs in their own queue with a small goroutine budget so they cannot occupy every slot — the same isolation principle as in right-sizing worker concurrency per CPU.
// Prometheus: expose Asynq queue metrics with the bundled collector
reg := prometheus.NewRegistry()
reg.MustRegister(metrics.NewQueueMetricsCollector(asynq.NewInspector(redisOpt)))
http.Handle("/metrics", promhttp.HandlerFor(reg, promhttp.HandlerOpts{}))
# Oldest pending task age per queue (Asynq collector)
max by (queue) (asynq_queue_latency_seconds) > 60

FAQ

Can I just use goroutines and a channel as a queue? For work that may be lost — best-effort cache warming, metrics flushing — yes. For anything a user or another system depends on, no: an in-memory channel disappears with the process. Durable enqueue before acknowledging the request is the only way to guarantee the work happens.

Asynq or River? River if jobs are consequences of Postgres writes and you want transactional enqueue with no extra infrastructure. Asynq if you need higher throughput, lower latency, or already operate Redis well. Both are production-grade.

How do I test Go job handlers? Test handlers as plain functions with a constructed task or job and a context. For integration tests, both libraries run against a real Redis or Postgres in a container; River also offers test helpers to assert that a job was inserted within a transaction.

Should workers run in the same binary as the API? Usually build one binary with two entry points (an api and a worker subcommand) and deploy them as separate processes. They share the args structs and handler code, which keeps producer and consumer in lockstep, while separate deployments let you scale workers on queue depth and restart them without touching request-serving pods. Running workers inside API pods couples their scaling and makes every API deploy interrupt background work.

Is Machinery still a good choice? Machinery (Celery-inspired, multiple brokers) is in maintenance mode. For new projects, Asynq, River, or a managed cloud queue are better-supported choices.

Related