后台任务不应该都塞在 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)
}
}
优雅关闭:停止队列并等待任务完成
服务关闭时,不能简单退出——这会中断正在执行的任务。正确的关闭顺序是:
- 停止接收新任务(关闭 channel)
- 用 WaitGroup 等待所有 worker 处理完已取出的任务
- 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 或云队列,会更容易判断取舍。
关键设计原则回顾:
- 队列大小要匹配 worker 能力和业务并发量
- 队列满策略根据业务重要性选择丢弃、阻塞或落地持久化
- 优雅关闭保证正在执行的任务不被中断
- 区分错误类型:可重试的做退避重试,不可重试的立刻失败
- 幂等性避免重试导致的副作用
- 日志和指标让异步任务的可观测性不低于同步接口
从内存队列迁移到外部消息队列时,最自然的演进路径是保持 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 验证并发安全性。
生产环境注意事项
生产环境的代码比本地开发要求更高。以下是一些通用原则:
- 日志要克制:不要记录敏感信息,不要在热路径上打印大量日志。
- 超时和取消:所有外部调用都要有超时。使用
context.WithTimeout或context.WithDeadline。 - 资源限制:限制请求体大小、并发连接数、内存使用。
- 优雅关闭:http.Server 要设置 Shutdown 超时,goroutine 要有退出机制。
- 可观测性:至少记录关键指标(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,而是理解背后的设计原则和适用边界。先让代码工作,再让它正确,最后才考虑让它更快。清晰的代码比聪明的代码更有价值。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。