Getting Started with Asynq in Go
Asynq is the most widely used Redis-backed task queue for Go, and this guide takes a service from no background jobs to a deployed, observable Asynq worker, as part of Task Queues in Go in Backend Frameworks & Worker Scaling. Each step adds one production concern — typed tasks, retries, priorities, uniqueness, scheduling, monitoring — so you can stop at the level your service needs.
Problem Statement
An image-hosting API written in Go resizes uploads synchronously inside the HTTP handler. Large uploads push p99 request latency past 8 seconds and occasionally hit the load balancer's 30-second timeout, after which the client retries and the image is processed twice. You want the handler to return in milliseconds after recording the upload, resizing to happen in background workers with retries, urgent re-processing requests to jump ahead of bulk backfills, and duplicate requests for the same image to collapse into one task.
Prerequisites
- Go 1.22+ and a module for the service.
- Redis 6.2+ (Asynq uses Lua scripts and requires
noevictionfor durability), reachable from API and worker pods. go get github.com/hibiken/asynqand optionally theasynqmonweb UI.- Handlers that can be safely retried — resizing writes deterministic output keys, so a repeat overwrites the same files.
Step 1 — Define Tasks in a Shared Package
Put task type names, payload structs, and constructors in one package imported by both the API and the worker. Constructors are the single place where default options live.
// internal/tasks/tasks.go
package tasks
import (
"encoding/json"
"fmt"
"time"
"github.com/hibiken/asynq"
)
const TypeImageResize = "image:resize"
type ImageResizePayload struct {
ImageID string `json:"image_id"`
Sizes []int `json:"sizes"`
}
func NewImageResizeTask(imageID string, sizes []int) (*asynq.Task, error) {
b, err := json.Marshal(ImageResizePayload{ImageID: imageID, Sizes: sizes})
if err != nil {
return nil, fmt.Errorf("marshal: %w", err)
}
return asynq.NewTask(TypeImageResize, b,
asynq.MaxRetry(8), // then archived (Asynq's dead set)
asynq.Timeout(2*time.Minute), // handler ctx deadline
asynq.Retention(24*time.Hour), // keep completed task info for debugging
), nil
}
Retention keeps a completed task's record for the given duration so it shows in Asynqmon; without it, successful tasks disappear immediately. Keep it short — every retained task costs Redis memory.
Step 2 — Enqueue from the API
The client is safe for concurrent use; create one per process and reuse it. Enqueue after the upload is durably recorded, and choose the queue per call.
// cmd/api/main.go (excerpt)
var asynqClient = asynq.NewClient(asynq.RedisClientOpt{Addr: os.Getenv("REDIS_ADDR")})
func (h *Handler) Upload(w http.ResponseWriter, r *http.Request) {
img, err := h.store.SaveOriginal(r.Context(), r.Body) // S3 + DB row
if err != nil {
http.Error(w, "upload failed", http.StatusInternalServerError)
return
}
task, _ := tasks.NewImageResizeTask(img.ID, []int{128, 512, 2048})
info, err := asynqClient.EnqueueContext(r.Context(), task,
asynq.Queue("default"),
asynq.TaskID("resize:"+img.ID), // deterministic id: duplicates are rejected
)
if err != nil && !errors.Is(err, asynq.ErrTaskIDConflict) {
// The image is stored but not queued: record it for the reconciler (Step 6)
h.log.Error("enqueue failed", "image_id", img.ID, "err", err)
}
writeJSON(w, http.StatusAccepted, map[string]string{"image_id": img.ID, "task_id": info.ID})
}
A deterministic TaskID makes enqueue idempotent: a client retry that uploads the same image produces the same id, and Asynq returns ErrTaskIDConflict instead of creating a second task while the first still exists. This is the cheapest duplicate protection Asynq offers and fixes the double-processing from the problem statement.
Step 3 — Write the Handler with Retry Classification
Returning an error makes Asynq retry with exponential backoff. Wrapping asynq.SkipRetry sends the task straight to the archive — use it for failures that will never succeed.
// internal/worker/resize.go
func HandleImageResize(store *storage.Client) asynq.HandlerFunc {
return func(ctx context.Context, t *asynq.Task) error {
var p tasks.ImageResizePayload
if err := json.Unmarshal(t.Payload(), &p); err != nil {
return fmt.Errorf("decode payload: %v: %w", err, asynq.SkipRetry)
}
src, err := store.GetOriginal(ctx, p.ImageID)
if errors.Is(err, storage.ErrNotFound) {
return fmt.Errorf("image %s deleted: %w", p.ImageID, asynq.SkipRetry)
}
if err != nil {
return fmt.Errorf("fetch original: %w", err) // transient: retry
}
for _, size := range p.Sizes {
if err := ctx.Err(); err != nil {
return err // timeout or shutdown
}
out, err := resize(src, size)
if err != nil {
return fmt.Errorf("resize %d: %v: %w", size, err, asynq.SkipRetry) // corrupt input
}
if err := store.PutVariant(ctx, p.ImageID, size, out); err != nil {
return fmt.Errorf("store variant: %w", err) // transient: retry
}
}
return nil
}
}
The classification mirrors the approach in retrying only transient errors by exception type: decode failures, deleted inputs, and corrupt images skip retries; network and storage errors retry. Checking ctx.Err() between sizes lets a long job stop promptly on timeout or shutdown.
Getting this classification right has a direct cost impact. A corrupt upload that is retried eight times with exponential backoff occupies a worker slot eight times over roughly half an hour and produces eight alerts' worth of error logs, all for an input that could never succeed. Marking it SkipRetry sends it to the archive on the first attempt, where it waits for a human with its payload and error intact. Conversely, wrapping a storage timeout in SkipRetry turns a two-second network blip into a permanently unprocessed image. When in doubt, classify by asking whether the same input could succeed on a later attempt.
Step 4 — Run the Server with Weighted Queues
The worker process runs an asynq.Server with a mux mapping task types to handlers.
// cmd/worker/main.go
func main() {
redisOpt := asynq.RedisClientOpt{Addr: os.Getenv("REDIS_ADDR"), PoolSize: 40}
srv := asynq.NewServer(redisOpt, asynq.Config{
Concurrency: 16, // CPU-bound resizing: ~2x cores
Queues: map[string]int{
"critical": 6, // user-triggered reprocessing
"default": 3, // new uploads
"backfill": 1, // bulk re-renders
},
RetryDelayFunc: func(n int, err error, t *asynq.Task) time.Duration {
d := time.Duration(math.Pow(2, float64(n))) * time.Second
if d > 10*time.Minute {
d = 10 * time.Minute
}
return d/2 + time.Duration(rand.Int63n(int64(d/2))) // capped exponential + jitter
},
ShutdownTimeout: 25 * time.Second,
ErrorHandler: asynq.ErrorHandlerFunc(func(ctx context.Context, t *asynq.Task, err error) {
retried, _ := asynq.GetRetryCount(ctx)
maxRetry, _ := asynq.GetMaxRetry(ctx)
slog.Error("task failed", "type", t.Type(), "retry", retried, "max", maxRetry, "err", err)
}),
})
mux := asynq.NewServeMux()
mux.Handle(tasks.TypeImageResize, HandleImageResize(storage.New()))
if err := srv.Run(mux); err != nil {
log.Fatalf("asynq server: %v", err)
}
}
With weights 6:3:1, a worker with work in every queue picks critical about 60% of the time and backfill about 10%. A million-image backfill therefore slows new-upload processing by roughly 10% instead of blocking it. If backfills must never delay anything, set StrictPriority: true and accept that backfill runs only when the other queues are empty.
Step 5 — Add Uniqueness and Scheduling Where Needed
Beyond deterministic TaskIDs, Asynq offers Unique(ttl), which rejects duplicate tasks (same type, payload, and queue) for a time window, and delayed processing.
// Collapse repeated "regenerate thumbnails" clicks for 5 minutes
client.Enqueue(task, asynq.Queue("critical"), asynq.Unique(5*time.Minute))
// Re-check a pending moderation result in 10 minutes
client.Enqueue(checkTask, asynq.ProcessIn(10*time.Minute))
Unique holds a Redis lock for the TTL, even after the task completes — a second enqueue inside the window is rejected although nothing is running. Prefer TaskID when you want "at most one pending", and Unique when you want "at most one per time window". Uniqueness trade-offs in general are covered in deduplicating jobs with Redis SET NX keys.
Step 6 — Deploy Workers and the Monitoring UI
Run the worker as its own deployment. Keep ShutdownTimeout below the pod's grace period so in-flight tasks finish or are re-queued before SIGKILL.
apiVersion: apps/v1
kind: Deployment
metadata: { name: image-worker }
spec:
replicas: 4
template:
spec:
terminationGracePeriodSeconds: 30 # > ShutdownTimeout (25s)
containers:
- name: worker
image: registry.internal/images:${GIT_SHA}
args: ["worker"]
env:
- { name: REDIS_ADDR, value: "redis-jobs:6379" }
resources:
requests: { cpu: "2", memory: "1Gi" }
limits: { memory: "2Gi" }
---
# Asynqmon: web UI for queues, retries, archived tasks (put behind SSO)
apiVersion: apps/v1
kind: Deployment
metadata: { name: asynqmon }
spec:
replicas: 1
template:
spec:
containers:
- name: asynqmon
image: hibiken/asynqmon:latest
args: ["--redis-addr=redis-jobs:6379", "--enable-metrics-exporter"]
ports: [{ containerPort: 8080 }]
--enable-metrics-exporter exposes queue size, latency, and processed/failed counters at /metrics for Prometheus. Scale the worker deployment on queue latency rather than CPU, as described in scaling workers with KEDA on queue length.
Verification
# Enqueue a test task and watch it move through the states
asynq task ls --queue=default --state=pending
asynq stats # processed / failed counts per queue
# Prove duplicate suppression: second enqueue with the same TaskID must conflict
go run ./cmd/enqueue --image img-test-1 && go run ./cmd/enqueue --image img-test-1
# expect: "task ID conflicts with another task"
In Prometheus, asynq_queue_latency_seconds{queue="default"} should stay low under normal load, and asynq_tasks_processed_total should track upload rate.
Gotchas & Edge Cases
Redis eviction deletes tasks. Asynq's data lives in Redis keys; any maxmemory-policy other than noeviction lets Redis delete them under pressure. Use a dedicated instance with noeviction and alert on memory usage.
Archived tasks are the dead-letter set. Tasks that exhaust retries or return SkipRetry move to the archive, which has its own size and age caps. Check it regularly or alert on its size; it is not visible unless someone looks.
Handlers must be registered in every worker. A task whose type has no handler fails with "handler not found" and retries. Deploy handlers before producers start enqueueing a new type.
Payload changes during rolling deploys. Old workers may receive tasks from new producers. Add fields as optional and ignore unknown ones; never change the meaning of an existing field.
FAQ
How many goroutines should Concurrency be? For CPU-bound handlers like image resizing, about one to two per core. For I/O-bound handlers, tens per core is normal. Measure: raise it until throughput stops improving or a downstream limit is hit.
Does Asynq guarantee exactly-once processing? No — it is at-least-once. A worker crash after the side effect but before acknowledgement causes a retry. Make handlers idempotent.
Can Asynq use Redis Cluster or Sentinel?
Yes: asynq.RedisClusterClientOpt and asynq.RedisFailoverClientOpt support both. Sentinel failover can lose the most recent writes if replication was behind; see Redis Sentinel for queue high availability.
Related
- Task Queues in Go — where Asynq fits among Go options.
- River: a Postgres Job Queue for Go — the Postgres-backed alternative.
- Graceful Shutdown for Go Workers — getting deploys right with Asynq.
- Choosing a Go Job Queue Library — a requirements-driven comparison.