Go Worker 队列入门:用 channel 和 context 处理后台任务

本文详解 Go 中使用 channel、goroutine、WaitGroup 和 context 实现 Worker 队列的方法,附带重试、观测和队列边界说明。

后台任务不应该都塞在 HTTP 请求里

很多 Web 服务会遇到一些不适合在请求里同步完成的事情:发送邮件、生成报表、处理图片、同步第三方数据、写审计日志。简单项目里你可能先在 handler 里直接调用这些逻辑,但请求会变慢,也更容易受外部系统影响。

假设发送一封邮件平均耗时 200ms,如果放在注册接口里同步发送,用户注册完等待 200ms 才能看到成功页面。如果邮件服务偶尔超时,整个注册流程可能卡在 5 秒超时然后失败。这就是典型的"不该同步做的事被做成了同步"。

Go 的 goroutine 和 channel 很适合实现入门级后台 worker。你可以把任务放进 channel,由固定数量 worker 消费。再配合 context,就能在服务关闭时停止接收新任务,并等待已有任务结束。

这篇文章实现一个内存里的 Worker 队列。它不替代真正的消息队列,比如 Redis、RabbitMQ、Kafka;但适合理解并发任务处理的基本结构,也适合内部工具、原型开发和小型服务的后台任务。

定义任务和队列结构

先定义任务结构体和队列:

package worker

import (
	"context"
	"fmt"
	"sync"
)

// Job 定义一个后台任务
type Job struct {
	ID       int64
	Type     string
	Email    string
	Body     string
	Attempts int // 已尝试次数
}

// Queue 是一个有缓冲 channel 的内存队列
type Queue struct {
	jobs chan Job
	wg   sync.WaitGroup
}

// NewQueue 创建指定缓冲大小的队列
func NewQueue(size int) *Queue {
	return &Queue{
		jobs: make(chan Job, size),
	}
}

缓冲区大小的选择取决于业务场景:

  • 太小:任务容易堆积,Submit 频繁返回队列满错误
  • 太大:占用更多内存,且服务关闭时需要等待更多积压任务
  • 一般建议:worker 数量 × 每个 worker 最大并发任务数

任务提交与队列满策略

提交任务时,有几种队列满的处理策略:

package worker

import (
	"context"
	"fmt"
)

// Submit 将任务放入队列,队列满时返回错误
func (q *Queue) Submit(ctx context.Context, job Job) error {
	select {
	case q.jobs <- job:
		return nil
	case <-ctx.Done():
		return ctx.Err()
	default:
		return fmt.Errorf("queue is full: job %d dropped", job.ID)
	}
}

这里用了 default,表示队列满时立刻返回错误,而不是阻塞请求。是否阻塞要看业务需求。对 HTTP 请求来说,队列满了返回 503(Service Unavailable)可能比一直卡住更好。客户端收到 503 后可以等待重试,或者把任务落地到数据库稍后处理。

另一种策略是带超时等待:

// SubmitWithTimeout 队列满时等待一段时间
func (q *Queue) SubmitWithTimeout(ctx context.Context, job Job, timeout time.Duration) error {
	ctx, cancel := context.WithTimeout(ctx, timeout)
	defer cancel()

	select {
	case q.jobs <- job:
		return nil
	case <-ctx.Done():
		return fmt.Errorf("submit job %d timed out or cancelled: %w", job.ID, ctx.Err())
	}
}

对非关键任务可以用"丢弃策略",对关键任务应该使用"阻塞+超时"或持久化到可靠存储。

启动 worker 处理任务

func (q *Queue) Start(ctx context.Context, workerCount int, handle func(context.Context, Job) error) {
	for i := 0; i < workerCount; i++ {
		q.wg.Add(1)
		go func(workerID int) {
			defer q.wg.Done()
			for {
				select {
				case job, ok := <-q.jobs:
					if !ok {
						log.Printf("worker=%d channel closed, exiting", workerID)
						return
					}
					if err := handle(ctx, job); err != nil {
						log.Printf("worker=%d job=%d type=%s error=%v", workerID, job.ID, job.Type, err)
					}
				case <-ctx.Done():
					log.Printf("worker=%d context cancelled, exiting", workerID)
					return
				}
			}
		}(i + 1)
	}
}

注意这里 job, ok := <-q.jobs 同时处理了 channel 关闭和 context 取消两种情况。ok 为 false 表示 channel 已关闭,此时 worker 应该退出。

处理函数示例:

func SendEmail(ctx context.Context, job Job) error {
	select {
	case <-time.After(200 * time.Millisecond):
		log.Printf("send email to %s: %s", job.Email, job.Body)
		return nil
	case <-ctx.Done():
		return ctx.Err()
	}
}

启动队列:

ctx, cancel := context.WithCancel(context.Background())
defer cancel()

queue := NewQueue(100)
queue.Start(ctx, 4, SendEmail)

// 提交任务
for i := 1; i <= 10; i++ {
	job := Job{ID: int64(i), Type: "email", Email: fmt.Sprintf("user%d@example.com", i), Body: "Welcome!"}
	if err := queue.Submit(ctx, job); err != nil {
		log.Printf("submit failed: %v", err)
	}
}

优雅关闭:停止队列并等待任务完成

服务关闭时,不能简单退出——这会中断正在执行的任务。正确的关闭顺序是:

  1. 停止接收新任务(关闭 channel)
  2. 用 WaitGroup 等待所有 worker 处理完已取出的任务
  3. worker 在 channel 关闭后退出循环
func (q *Queue) Stop() {
	close(q.jobs)
	q.wg.Wait()
	log.Println("all workers finished, queue stopped")
}

完整 worker 循环中,case job, ok := <-q.jobs 同时在处理 context 取消和 channel 关闭两种情况。服务关闭时,你可以先停止接收 HTTP 请求,再调用 queue.Stop() 等待 worker 结束。

完整的优雅关闭示例:

package main

import (
	"context"
	"fmt"
	"log"
	"os"
	"os/signal"
	"syscall"
	"time"
	"worker"
)

func main() {
	ctx, cancel := context.WithCancel(context.Background())
	defer cancel()

	queue := worker.NewQueue(100)
	queue.Start(ctx, 4, worker.SendEmail)

	// 模拟提交任务
	go func() {
		for i := 1; i <= 100; i++ {
			job := worker.Job{ID: int64(i), Type: "email", Email: fmt.Sprintf("user%d@example.com", i), Body: "Hello"}
			if err := queue.Submit(ctx, job); err != nil {
				log.Printf("submit failed: %v", err)
			}
			time.Sleep(10 * time.Millisecond)
		}
	}()

	// 等待中断信号
	sigCh := make(chan os.Signal, 1)
	signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
	<-sigCh

	log.Println("shutting down gracefully...")
	cancel()       // 通知 context 取消
	queue.Stop()   // 关闭 channel 并等待 worker
	log.Println("shutdown complete")
}

这个模式是 Go 服务优雅关闭的标准做法。注意调用顺序:cancel() 通知所有 goroutine 应该停止,然后 queue.Stop() 关闭 channel 并等待。如果反过来先关闭 channel 再 cancel,可能导致正在取任务的 worker 永远等待 ctx.Done()。

增加指数退避重试

入门版可以在 worker 内部做有限重试。简单的线性重试:

func handleWithRetry(ctx context.Context, job Job, handle func(context.Context, Job) error) error {
	var lastErr error
	for attempt := 1; attempt <= 3; attempt++ {
		if err := handle(ctx, job); err != nil {
			lastErr = err
			select {
			case <-time.After(time.Duration(attempt) * 200 * time.Millisecond):
			case <-ctx.Done():
				return ctx.Err()
			}
			continue
		}
		return nil
	}
	return fmt.Errorf("job %d failed after %d retries: %w", job.ID, 3, lastErr)
}

但注意:这不适合所有错误。比如参数格式错误、用户不存在、权限不足,重试多少次都不会成功。真实系统应该区分可重试错误和不可重试错误:

var ErrNotRetryable = errors.New("not retryable")

type RetryableError struct {
	Cause error
}

func (e *RetryableError) Error() string {
	return fmt.Sprintf("retryable: %v", e.Cause)
}

func handleWithSmartRetry(ctx context.Context, job Job, handle func(context.Context, Job) error) error {
	var lastErr error
	for attempt := 1; attempt <= 3; attempt++ {
		err := handle(ctx, job)
		if err == nil {
			return nil
		}
		lastErr = err

		// 不可重试错误,立刻返回
		if errors.Is(err, ErrNotRetryable) {
			return err
		}

		// 可重试错误,等待后重试
		backoff := time.Duration(attempt*attempt) * 100 * time.Millisecond // 指数退避
		select {
		case <-time.After(backoff):
		case <-ctx.Done():
			return ctx.Err()
		}
	}
	return fmt.Errorf("job %d failed after retries: %w", job.ID, lastErr)
}

还要考虑幂等性。如果发送邮件接口不支持幂等,重试可能导致用户收到多封邮件。后台任务不是简单加 goroutine 就结束,任务语义同样重要。实现幂等性的常见方法是为每个任务生成唯一 ID 并在服务端去重:

func SendEmailWithID(ctx context.Context, job Job) error {
	// 假设服务端支持 X-Idempotency-Key	headers := map[string]string{
		"X-Idempotency-Key": fmt.Sprintf("email-%d", job.ID),
	}
	_ = headers
	return nil
}

给队列加观测信息和指标

后台任务如果没有日志和指标,出问题很难查。至少记录任务开始、失败和最终成功:

func loggedHandle(ctx context.Context, job Job, handle func(context.Context, Job) error) error {
	start := time.Now()
	log.Printf("job started id=%d type=%s email=%s", job.ID, job.Type, job.Email)

	err := handle(ctx, job)
	if err != nil {
		log.Printf("job failed id=%d type=%s duration=%s err=%v", job.ID, job.Type, time.Since(start), err)
		return err
	}

	log.Printf("job finished id=%d type=%s duration=%s", job.ID, job.Type, time.Since(start))
	return nil
}

还可以记录队列深度、处理成功数、失败数。即使不用完整监控系统,简单计数也能回答几个关键问题:任务是不是堆积了?失败是不是突然变多?处理耗时是不是变长?

用原子操作做无锁计数:

type QueueMetrics struct {
	Submitted  atomic.Int64
	Processed  atomic.Int64
	Failed     atomic.Int64
	QueueDepth atomic.Int64 // 近似值
}

func (m *QueueMetrics) RecordSubmit() {
	m.Submitted.Add(1)
	m.QueueDepth.Add(1)
}

func (m *QueueMetrics) RecordDone(success bool) {
	m.Processed.Add(1)
	m.QueueDepth.Add(-1)
	if !success {
		m.Failed.Add(1)
	}
}

在 HTTP 接口中暴露简单的状态页:

func (m *QueueMetrics) Handler() http.HandlerFunc {
	return func(w http.ResponseWriter, r *http.Request) {
		fmt.Fprintf(w, "submitted: %d\n", m.Submitted.Load())
		fmt.Fprintf(w, "processed: %d\n", m.Processed.Load())
		fmt.Fprintf(w, "failed: %d\n", m.Failed.Load())
		fmt.Fprintf(w, "queue_depth: %d\n", m.QueueDepth.Load())
	}
}

如果任务重要,日志里要有任务 ID 或业务 ID,方便从用户反馈一路追到后台执行记录。异步系统最怕请求已经返回成功,但后台到底有没有执行没人知道。可以通过 webhook 或回调机制通知调用方最终结果。

入门队列的边界:何时该用消息队列

内存队列有明显限制:

  • 进程崩溃任务会丢失:没有持久化,重启后所有未处理任务消失
  • 无法跨多实例共享:没有分布式协调,适合单实例部署
  • 没有重试持久化:超出重试次数的任务直接失败,不会进入死信队列
  • 没有延迟任务:所有任务立即消费,不适合定时执行场景
  • 没有可视化管理:无法查看任务状态、重新触发、手动删除

它适合轻量异步处理和学习,不适合关键业务任务。如果任务不能丢,比如支付后发货、订单状态同步,应该使用可靠消息队列或数据库任务表,并设计重试、幂等和监控。

常用替代方案对比:

方案优点缺点适用场景
内存 channel简单、零依赖、低延迟无持久化、单机原型、非关键任务
Redis + List/Stream持久化、分布式、可观测需要 Redis 运维中小规模后台任务
RabbitMQ成熟、路由灵活、死信运维复杂企业级消息系统
Kafka高吞吐、可回溯延迟高、复杂大数据流水线
PostgreSQL 任务表事务一致性、已有数据库性能受限已有 PG 的业务系统

完整可运行的 Worker 示例

把上述概念组合起来的完整实现:

package main

import (
	"context"
	"fmt"
	"log"
	"os"
	"os/signal"
	"sync"
	"syscall"
	"time"
)

type Job struct {
	ID    int64
	Email string
	Body  string
}

type Queue struct {
	jobs chan Job
	wg   sync.WaitGroup
}

func NewQueue(size int) *Queue {
	return &Queue{jobs: make(chan Job, size)}
}

func (q *Queue) Submit(job Job) error {
	select {
	case q.jobs <- job:
		return nil
	default:
		return fmt.Errorf("queue full")
	}
}

func (q *Queue) Start(ctx context.Context, workers int, handler func(context.Context, Job) error) {
	for i := 0; i < workers; i++ {
		q.wg.Add(1)
		go func(id int) {
			defer q.wg.Done()
			for {
				select {
				case job, ok := <-q.jobs:
					if !ok {
						return
					}
					if err := handler(ctx, job); err != nil {
						log.Printf("worker=%d job=%d error=%v", id, job.ID, err)
					}
				case <-ctx.Done():
					return
				}
			}
		}(i + 1)
	}
}

func (q *Queue) Stop() {
	close(q.jobs)
	q.wg.Wait()
}

func main() {
	ctx, cancel := context.WithCancel(context.Background())
	defer cancel()

	queue := NewQueue(100)
	queue.Start(ctx, 4, func(ctx context.Context, job Job) error {
		log.Printf("processing job=%d email=%s", job.ID, job.Email)
		time.Sleep(50 * time.Millisecond)
		return nil
	})

	// 提交任务
	for i := 1; i <= 50; i++ {
		_ = queue.Submit(Job{ID: int64(i), Email: fmt.Sprintf("user%d@example.com", i)})
	}

	// 优雅关闭
	sig := make(chan os.Signal, 1)
	signal.Notify(sig, syscall.SIGINT, syscall.SIGTERM)
	<-sig

	log.Println("shutting down...")
	cancel()
	queue.Stop()
	log.Println("done")
}

小结

Go 可以用 channel、goroutine、WaitGroup 和 context 实现简单 Worker 队列。核心结构是:有缓冲 channel 存任务,固定数量 worker 消费,提交时处理队列满,关闭时停止接收并等待 worker 退出。

这种队列适合入门和轻量任务,但不要把它当成可靠消息系统。理解它的结构和边界后,再学习 Redis 队列、Kafka 或云队列,会更容易判断取舍。

关键设计原则回顾:

  1. 队列大小要匹配 worker 能力和业务并发量
  2. 队列满策略根据业务重要性选择丢弃、阻塞或落地持久化
  3. 优雅关闭保证正在执行的任务不被中断
  4. 区分错误类型:可重试的做退避重试,不可重试的立刻失败
  5. 幂等性避免重试导致的副作用
  6. 日志和指标让异步任务的可观测性不低于同步接口

从内存队列迁移到外部消息队列时,最自然的演进路径是保持 Submit 和 Start 接口不变,只替换底层存储和分发机制。这样早期验证过的业务逻辑不需要重写,也能享受到生产级消息队列的持久化和分布式能力。

性能对比与基准测试

理解 Go Worker 队列入门 的最佳方式是通过基准测试观察实际行为。下面是一个基本的测试框架:

func BenchmarkMain(b *testing.B) {
    for i := 0; i < b.N; i++ {
        _ = i
    }
}

运行 go test -bench=. -benchmem 可以得到每个操作的耗时和内存分配数据。对比不同实现时,建议固定输入规模,跑多次取平均值。机器负载、CPU 频率和缓存状态都会影响结果,所以重要的优化应该在稳定环境中反复验证。

常见错误与最佳实践

错误一:性能优化过早

很多初学者在代码刚写好就开始担心性能,结果引入了不必要的复杂度。正确的做法是先用清晰的写法实现功能,在性能问题真实出现时再通过 profile 定位热点,再针对性优化。

错误二:忽略边界条件

空输入、超大输入、并发场景、系统资源耗尽等边界条件往往是 bug 的来源。写代码时养成习惯:每个函数都问自己,空值怎么办?错误怎么处理?资源泄漏有没有可能?

错误三:错误处理不完整

Go 的错误处理要求显式检查。常见问题是只在最外层处理错误,中间层把 error 吞掉或转换后丢失了上下文。使用 fmt.Errorf 配合 %w 保留原始错误链,上层可以用 errors.Is 判断。

错误四:并发代码缺少同步

Go 的并发模型很简洁,但共享内存访问必须同步。不要凭感觉认为"这里应该不会并发访问"就省略锁或原子操作。用 go test -race 验证并发安全性。

生产环境注意事项

生产环境的代码比本地开发要求更高。以下是一些通用原则:

  1. 日志要克制:不要记录敏感信息,不要在热路径上打印大量日志。
  2. 超时和取消:所有外部调用都要有超时。使用 context.WithTimeout 或 context.WithDeadline。
  3. 资源限制:限制请求体大小、并发连接数、内存使用。
  4. 优雅关闭:http.Server 要设置 Shutdown 超时,goroutine 要有退出机制。
  5. 可观测性:至少记录关键指标(QPS、延迟、错误率)。

测试策略

好的测试应该覆盖正常路径、错误路径和边界条件。表驱动测试是 Go 社区推荐的方式:

func TestExample(t *testing.T) {
    tests := []struct {
        name string
        input string
        want  string
    }{
        {"valid", "hello", "HELLO"},
        {"empty", "", ""},
    }
    for _, tt := range tests {
        t.Run(tt.name, func(t *testing.T) {
            got := strings.ToUpper(tt.input)
            if got != tt.want {
                t.Fatalf("ToUpper(%q) = %q, want %q", tt.input, got, tt.want)
            }
        })
    }
}

实战 FAQ

Q: 这个功能在旧版 Go 中能用吗?
A: 需要看具体功能引入的版本。建议使用最新的稳定版 Go。

Q: 第三方库更好还是标准库更好?
A: 能标准库解决先用标准库。第三方库引入依赖成本和许可证风险。

Q: 写测试时发现代码难测怎么办?
A: 这通常意味着代码耦合度太高。考虑把大函数拆成小函数,把外部依赖抽象成接口。

Q: 怎么判断代码算不算过度设计?
A: 问自己:这个抽象让调用方更简单了吗?减少了多少重复?维护成本是增加还是减少了?

小结

Go Worker 队列入门 是 Go 开发中非常实用的技能。关键不是记住所有 API,而是理解背后的设计原则和适用边界。先让代码工作,再让它正确,最后才考虑让它更快。清晰的代码比聪明的代码更有价值。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「golang」更多文章

  1. 熔断、降级与限流:Go 微服务韧性设计完全指南
  2. 事件溯源与 CQRS 在 Go 中的实践:复杂业务系统的架构升级
  3. TinyGo 嵌入式开发与物联网实战:微控制器编程完全指南