DEV Community

Ayi NEDJIMI
Ayi NEDJIMI

Posted on

How to build a job queue in Go without external dependencies

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
}
Enter fullscreen mode Exit fullscreen mode

A few design decisions worth noting:

  • Buffered channel: Enqueue uses a non-blocking select. 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: 3 occupies 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")
}
Enter fullscreen mode Exit fullscreen mode

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
            }
        }
    }
}
Enter fullscreen mode Exit fullscreen mode

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)