Most job queue tutorials start with Redis, RabbitMQ, or some managed queue service. That's fine for teams with existing infrastructure — but if you're building an internal tool, a side project, or a microservice that needs to process tasks asynchronously, pulling in a message broker is overkill. Go's standard library gives you everything you need: goroutines, channels, sync primitives, and context for cancellation.
What we're building
A concurrent job queue with:
- A fixed pool of workers pulling from a buffered channel
- Non-blocking enqueue — full queue returns an error, no caller hangs
- Per-job retry with a configurable limit
- Graceful shutdown on SIGINT/SIGTERM — no jobs dropped mid-flight
- An error channel so worker failures don't crash the process
Around 150 lines of Go. No external libraries.
The core structure
package queue
import (
"context"
"fmt"
"sync"
)
type Job struct {
ID string
Payload any
Retries int
MaxRetry int
}
type Handler func(ctx context.Context, job Job) error
type Queue struct {
jobs chan Job
wg sync.WaitGroup
handler Handler
errors chan error
}
func New(bufferSize, workers int, handler Handler) *Queue {
q := &Queue{
jobs: make(chan Job, bufferSize),
handler: handler,
errors: make(chan error, bufferSize),
}
for i := 0; i < workers; i++ {
q.wg.Add(1)
go q.worker(context.Background())
}
return q
}
func (q *Queue) Enqueue(job Job) error {
select {
case q.jobs <- job:
return nil
default:
return fmt.Errorf("queue full: dropped job %s", job.ID)
}
}
func (q *Queue) worker(ctx context.Context) {
defer q.wg.Done()
for job := range q.jobs {
err := q.handler(ctx, job)
if err != nil && job.Retries < job.MaxRetry {
job.Retries++
q.jobs <- job // re-queue
continue
}
if err != nil {
q.errors <- fmt.Errorf("job %s failed after %d retries: %w", job.ID, job.Retries, err)
}
}
}
func (q *Queue) Stop() {
close(q.jobs)
q.wg.Wait()
close(q.errors)
}
func (q *Queue) Errors() <-chan error {
return q.errors
}
A few design decisions worth noting:
-
Buffered channel:
Enqueueuses a non-blockingselect. The queue either accepts the job immediately or returns an error — the caller decides how to handle backpressure. - No global state: the queue is a plain struct. You can run multiple independent queues in the same process without interference.
-
Retry loop: re-queuing happens inside the worker rather than a separate scheduler. A job with
MaxRetry: 3occupies a worker slot for all attempts before giving up — acceptable for low-volume queues, worth revisiting if workers become a bottleneck.
Graceful shutdown with OS signals
package main
import (
"context"
"fmt"
"log"
"os"
"os/signal"
"syscall"
"time"
"yourmodule/queue"
)
func main() {
q := queue.New(100, 4, func(ctx context.Context, job queue.Job) error {
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(200 * time.Millisecond):
log.Printf("processed job %s: %v", job.ID, job.Payload)
return nil
}
})
go func() {
for err := range q.Errors() {
log.Println("worker error:", err)
}
}()
for i := 0; i < 20; i++ {
_ = q.Enqueue(queue.Job{
ID: fmt.Sprintf("job-%d", i),
Payload: i * 10,
MaxRetry: 2,
})
}
sig := make(chan os.Signal, 1)
signal.Notify(sig, syscall.SIGINT, syscall.SIGTERM)
<-sig
log.Println("shutting down...")
q.Stop()
log.Println("all jobs completed")
}
When SIGINT arrives, new enqueues stop and workers drain the remaining jobs before exiting. The sync.WaitGroup in Stop() blocks until every goroutine returns — no job is silently dropped.
Adding priorities with two channels
Standard channels have no priority concept. The practical solution is two channels — one for high-priority, one for normal — with a worker select that checks high-priority first.
type PriorityQueue struct {
high chan Job
normal chan Job
wg sync.WaitGroup
}
func (pq *PriorityQueue) worker(ctx context.Context, handler Handler, errors chan<- error) {
defer pq.wg.Done()
for {
// Always drain high-priority before picking normal
select {
case job, ok := <-pq.high:
if !ok {
return
}
if err := handler(ctx, job); err != nil {
errors <- err
}
continue
default:
}
select {
case job, ok := <-pq.high:
if !ok {
return
}
if err := handler(ctx, job); err != nil {
errors <- err
}
case job, ok := <-pq.normal:
if !ok {
return
}
if err := handler(ctx, job); err != nil {
errors <- err
}
}
}
}
The double select — first checking high alone with a default, then both — is the canonical Go pattern for priority queues. Without the first pass, a select with two ready channels picks randomly, so high-priority jobs can be starved under sustained normal load.
What you trade away vs. a real broker
Be clear-eyed about the limitations:
- No persistence: if the process crashes, queued jobs are gone. For anything financial or critical, write pending jobs to Postgres or SQLite before enqueuing.
- No distributed consumers: this queue lives in one process. Horizontal scaling requires distributing work at the application level before it reaches the queue.
- No built-in dead-letter queue: exhausted retries emit an error to the error channel. Logging or storing those failed jobs is your responsibility.
- Coarse backpressure: the buffer either accepts or rejects. Real brokers give you consumer acknowledgement, flow control, and visibility into queue depth across replicas.
For background jobs within a single service — sending emails, generating reports, processing uploads — this design handles thousands of jobs per second with zero operational overhead. When you need durability or cross-process coordination, move to a proper broker. If your pipeline handles sensitive data, review access controls before shipping — our free security checklists cover API and service hardening in detail.
The takeaway
Go's standard library is sufficient for the majority of job queue use cases. A buffered channel with a worker pool, a WaitGroup for shutdown, and an error channel for failures gives you a solid foundation without Redis, RabbitMQ, or any external dependency to operate.
Build this first. Add a broker when you can demonstrate — with real production data — that you actually need one.
I run AYI NEDJIMI Consultants, a cybersecurity consulting firm. We publish free security hardening checklists — PDF and Excel.
Top comments (0)