Graceful Shutdown for Go Workers
Every deploy of a worker fleet is a shutdown test, and this guide makes Go workers pass it, as part of Task Queues in Go in Backend Frameworks & Worker Scaling. It covers the shutdown sequence common to Asynq, River, and hand-rolled consumers: catch the signal, stop taking new work, give in-flight jobs a bounded time to finish, cancel what remains, and exit before the orchestrator's hard kill.
Problem Statement
A Go worker consuming from SQS processes document conversions that take 5–90 seconds. On every Kubernetes rollout, some conversions run twice and a few customers receive two notification emails. Logs show the pattern: the pod receives SIGTERM, keeps polling SQS for new messages, starts conversions it cannot finish, and is SIGKILLed 30 seconds later. The messages it held become visible again after the visibility timeout and are processed by the new pods — while the old pod had already sent the email. You want shutdown to stop fetching immediately, let short jobs finish, hand long jobs back cleanly, and never exceed the grace period.
Prerequisites
- A Go worker with a fetch loop and a pool of handler goroutines (library-based or hand-written).
- Handlers that accept a
context.Contextand pass it to every blocking call. - Idempotent side effects, since a job cancelled mid-way will be redelivered.
- Access to the Kubernetes Deployment spec for
terminationGracePeriodSecondsandpreStop.
Step 1 — Understand the Kubernetes Shutdown Timeline
When a pod is terminated, Kubernetes runs the preStop hook (if any), then sends SIGTERM to PID 1 in each container, and waits up to terminationGracePeriodSeconds (default 30, counted from the start of termination and including the preStop time) before sending SIGKILL. Your process must finish everything inside that window.
spec:
terminationGracePeriodSeconds: 120 # longest job you will let finish + margin
containers:
- name: worker
image: registry.internal/converter:${GIT_SHA}
lifecycle:
preStop:
exec:
command: ["/bin/sh", "-c", "sleep 2"] # let endpoint/queue-scaler updates settle
Two consequences shape the Go code. First, if your binary is not PID 1 (for example, started by a shell script), it may never receive SIGTERM; use exec in entrypoint scripts or a minimal init like tini. Second, the application's own drain timeout must be comfortably less than the grace period, or SIGKILL arrives mid-cleanup.
Step 2 — Turn the Signal into a Context
signal.NotifyContext gives you a context cancelled on SIGTERM or SIGINT. Use it as the fetch context only — not as the parent of job contexts, or every in-flight job would be cancelled the instant the signal arrives.
func main() {
// Cancelled on SIGTERM: stops the fetch loop, nothing else
fetchCtx, stopFetch := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGINT)
defer stopFetch()
// Parent of all job contexts: cancelled only when the drain window expires
jobsCtx, cancelJobs := context.WithCancel(context.Background())
defer cancelJobs()
w := NewWorker(sqsClient, queueURL, 16)
done := make(chan struct{})
go func() { w.Run(fetchCtx, jobsCtx); close(done) }()
<-fetchCtx.Done() // SIGTERM received
slog.Info("shutdown: stopped fetching, draining", "in_flight", w.InFlight())
select {
case <-w.Drained(): // all in-flight jobs finished
case <-time.After(90 * time.Second): // drain window expired
slog.Warn("drain timeout, cancelling stragglers", "in_flight", w.InFlight())
cancelJobs() // jobs see ctx.Done()
select {
case <-w.Drained():
case <-time.After(10 * time.Second):
slog.Error("stragglers ignored cancellation", "in_flight", w.InFlight())
}
}
<-done
}
Separating the two contexts is the core of the pattern: one signal stops intake, a later deadline stops work.
Step 3 — Stop Fetching Before Anything Else
The fetch loop must check the fetch context before every receive and must not start a job after shutdown begins. With long polling, pass the fetch context into the receive call so a 20-second long poll returns immediately on SIGTERM.
func (w *Worker) Run(fetchCtx, jobsCtx context.Context) {
sem := make(chan struct{}, w.concurrency)
for {
select {
case <-fetchCtx.Done():
return // never start new work after SIGTERM
case sem <- struct{}{}: // wait for a free slot
}
out, err := w.sqs.ReceiveMessage(fetchCtx, &sqs.ReceiveMessageInput{
QueueUrl: &w.queueURL,
MaxNumberOfMessages: 1,
WaitTimeSeconds: 20, // cancelled early via fetchCtx
VisibilityTimeout: 120,
})
if err != nil || len(out.Messages) == 0 {
<-sem
continue
}
msg := out.Messages[0]
w.wg.Add(1)
go func() {
defer func() { <-sem; w.wg.Done() }()
w.handle(jobsCtx, msg)
}()
}
}
Receiving one message per free slot (rather than ten at a time) matters at shutdown: a worker that prefetched ten messages but can only run four holds six it will never process, and they stay invisible for the whole visibility timeout. The same effect with other brokers is discussed in tuning prefetch and consumer concurrency.
Step 4 — Hand Back Cancelled Jobs Immediately
When a job is cancelled by the drain deadline, it should return its message to the queue now rather than let it sit invisible until the visibility timeout expires. On SQS, that is ChangeMessageVisibility to zero; on RabbitMQ, a nack with requeue; in Asynq and River, the library does it for you when the handler returns a context error.
func (w *Worker) handle(ctx context.Context, msg types.Message) {
err := w.convert(ctx, msg) // passes ctx to every I/O call
switch {
case err == nil:
w.sqs.DeleteMessage(context.Background(), &sqs.DeleteMessageInput{
QueueUrl: &w.queueURL, ReceiptHandle: msg.ReceiptHandle})
case errors.Is(err, context.Canceled):
// Shutdown cut us off: make the message visible to the new pods right away
w.sqs.ChangeMessageVisibility(context.Background(), &sqs.ChangeMessageVisibilityInput{
QueueUrl: &w.queueURL, ReceiptHandle: msg.ReceiptHandle, VisibilityTimeout: 0})
default:
// ordinary failure: leave it; it reappears after the visibility timeout (backoff)
}
}
Note the context.Background() for the delete and visibility calls: using the cancelled job context there would make the cleanup call fail instantly, which is a common and confusing bug.
Step 5 — Make Long Jobs Resumable or Short
A drain window cannot be longer than the grace period, and very long grace periods slow every rollout. For jobs that can run for many minutes, either split them into shorter steps, or checkpoint progress so a cancelled run resumes where it stopped.
func (w *Worker) convert(ctx context.Context, msg types.Message) error {
job := parse(msg)
cp, _ := w.checkpoints.Load(ctx, job.ID) // last completed page, if any
for page := cp.NextPage; page < job.Pages; page++ {
if err := ctx.Err(); err != nil {
return err // cancelled between pages
}
if err := w.renderPage(ctx, job, page); err != nil {
return err
}
w.checkpoints.Save(context.Background(), job.ID, page+1) // durable progress
}
return w.finalize(ctx, job)
}
With checkpoints, the redelivered job picks up at the next page, so cancellation costs at most one page of repeated work.
The same idea is behind the heartbeat-and-resume approach in configuring visibility timeouts for long-running workers.
Step 6 — Use the Library's Shutdown When You Have One
Asynq and River implement this sequence internally; configure their timeouts to fit inside the grace period.
// Asynq: stops fetching on SIGTERM, waits ShutdownTimeout, then re-queues unfinished tasks
srv := asynq.NewServer(redisOpt, asynq.Config{Concurrency: 16, ShutdownTimeout: 90 * time.Second})
srv.Run(mux) // handles signals itself
// River: Stop waits for running jobs; StopAndCancel cancels their contexts
<-fetchCtx.Done()
softCtx, cancel := context.WithTimeout(context.Background(), 90*time.Second)
defer cancel()
if err := riverClient.Stop(softCtx); err != nil {
hard, cancel2 := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel2()
riverClient.StopAndCancel(hard)
}
In both cases handler code still has to honour ctx.Done(); the library can cancel the context but cannot interrupt a goroutine that ignores it.
Verification
Run a deploy under load and look for three signals: no duplicated side effects, no messages stuck invisible, and no SIGKILLs.
# Pods killed rather than exiting cleanly show exit code 137
kubectl get pods -l app=converter -o jsonpath='{range .items[*]}{.metadata.name}{" "}{.status.containerStatuses[0].lastState.terminated.exitCode}{"\n"}{end}'
# During a rollout, in-flight (not visible) messages should drop quickly, not plateau
aws cloudwatch get-metric-statistics --namespace AWS/SQS \
--metric-name ApproximateNumberOfMessagesNotVisible \
--dimensions Name=QueueName,Value=conversions --period 60 --statistics Maximum \
--start-time "$(date -u -d '-15 min' +%FT%TZ)" --end-time "$(date -u +%FT%TZ)"
Pair it with a local test: start the worker, enqueue a 60-second job, send SIGTERM after 5 seconds with a 10-second drain, and assert the message becomes visible again within a second of cancellation.
Gotchas & Edge Cases
Background goroutines inside handlers. A handler that spawns its own goroutine and returns leaves work the drain logic cannot see. Use errgroup.WithContext and wait on it inside the handler.
Health checks during drain. If the liveness probe fails once the worker stops fetching, Kubernetes may restart the container mid-drain. Keep liveness passing until exit.
Autoscalers fighting the drain. A queue-length autoscaler may start replacement pods while old ones drain, which is fine — but a scale-down that picks a busy pod triggers the same shutdown path. Protect long jobs with the techniques in zero-downtime worker deploys on Kubernetes.
Log flushing. Buffered log handlers lose the final lines if the process exits without flushing. Flush loggers and tracers as the last step before returning from main.
FAQ
How long should the grace period be? Long enough for your p99 job duration plus the cancel budget and a margin — but not so long that rollouts crawl. Jobs longer than a couple of minutes should be split or checkpointed instead of stretching the grace period.
Should SIGTERM cancel running jobs immediately? No. Cancel intake immediately and running jobs only after the drain window. Immediate cancellation turns every deploy into a mass redelivery.
What about SIGINT and local development?
Handle both. signal.NotifyContext(ctx, syscall.SIGTERM, syscall.SIGINT) gives the same drain behaviour on Ctrl-C, which makes the shutdown path easy to exercise locally.
Related
- Task Queues in Go — the Go worker landscape.
- Graceful Shutdown & Worker Deployments — the cross-framework principles.
- Handling SIGTERM in Celery Workers — the same sequence in Python.
- Getting Started with Asynq in Go — ShutdownTimeout in context.