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.
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.
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.
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.
Performance Tuning
- Concurrency per process. Goroutines are cheap, but the resources jobs touch are not. Size
Concurrency(Asynq) orMaxWorkers(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 nearGOMAXPROCS. - Connection pools. River workers use the pgx pool for claiming, completing, and for your job code. Set
MaxConnsaboveMaxWorkersplus 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
- Getting Started with Asynq in Go — a Redis-backed worker from zero to production.
- River: a Postgres Job Queue for Go — transactional enqueue and Postgres-backed workers.
- Graceful Shutdown for Go Workers — signals, contexts, and grace periods.
- Choosing a Go Job Queue Library — a requirements-driven comparison.
- Database-Backed Job Queues — the design River is built on.