Go worker 取消入门:后台循环如何听懂 context

用消息消费 worker 示例讲 context 取消、select 循环、任务超时、资源关闭和测试后台 goroutine 退出。

后台 worker 是 Go 服务里很常见的结构:消费队列、定时清理、同步外部数据、处理导入任务。很多 worker 写起来很简单,但停不下来。服务收到关闭信号后,HTTP server 已经退出,worker 还卡在 sleep、网络请求或 channel receive 上,进程迟迟不结束。

本文讲 worker 如何使用 context 取消。核心原则是:循环要监听 ctx.Done(),任务处理要传递 context,阻塞等待要能被取消。

一个基础 worker

type Job struct {
	ID string
}

func RunWorker(ctx context.Context, jobs <-chan Job, handle func(context.Context, Job) error) {
	for {
		select {
		case <-ctx.Done():
			return
		case job, ok := <-jobs:
			if !ok {
				return
			}
			if err := handle(ctx, job); err != nil {
				log.Printf("handle job %s: %v", job.ID, err)
			}
		}
	}
}

这个 worker 会在 context 取消或 jobs channel 关闭时退出。处理函数也接收 context,便于内部数据库查询和 HTTP 调用停止。

给单个任务设置超时

整个 worker 的 ctx 表示服务生命周期;单个任务还可以有自己的超时:

func handleWithTimeout(parent context.Context, job Job) error {
	ctx, cancel := context.WithTimeout(parent, 30*time.Second)
	defer cancel()
	return processJob(ctx, job)
}

这样某个任务卡住不会无限占用 worker。注意从 parent 派生,这样服务关闭时任务也会立即取消。

不要用 time.Sleep 阻塞退出

坏写法:

for {
	process()
	time.Sleep(time.Minute)
}

取消时可能要等一分钟。用 timer 或 ticker 配合 select:

timer := time.NewTimer(time.Minute)
select {
case <-timer.C:
case <-ctx.Done():
	timer.Stop()
	return
}

定时循环用 time.Ticker 也可以,但记得 defer ticker.Stop()

处理中的任务怎么办

服务关闭时,worker 有两种策略:尽快取消当前任务,或者给当前任务一段时间收尾。选择取决于业务。发送通知可以取消后下次重试;写关键数据可能希望完成当前事务。

可以在 shutdown 时给外层 context 一个宽限期:

shutdownCtx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
defer cancel()
RunWorker(shutdownCtx, jobs, handle)

更常见的是 worker 在程序启动时运行,收到信号后取消 root context,然后等待 WaitGroup。

用 WaitGroup 等待退出

ctx, cancel := context.WithCancel(context.Background())
var wg sync.WaitGroup

for i := 0; i < 4; i++ {
	wg.Add(1)
	go func() {
		defer wg.Done()
		RunWorker(ctx, jobs, handle)
	}()
}

// 收到信号
cancel()
wg.Wait()

不要启动 goroutine 后就不管。长期运行服务应该知道自己启动了哪些后台任务,并在退出时等待它们收尾。

测试 worker 能退出

func TestWorkerStopsOnCancel(t *testing.T) {
	ctx, cancel := context.WithCancel(context.Background())
	jobs := make(chan Job)

	done := make(chan struct{})
	go func() {
		RunWorker(ctx, jobs, func(ctx context.Context, job Job) error {
			return nil
		})
		close(done)
	}()

	cancel()
	select {
	case <-done:
	case <-time.After(time.Second):
		t.Fatal("worker did not stop")
	}
}

这类测试很有价值。后台 goroutine 泄漏通常不会马上报错,但会让测试进程、服务退出和资源使用变得不稳定。

channel 关闭由谁负责

并发代码里一个常见约定是:谁发送,谁关闭。worker 只是从 jobs 里读任务,不应该关闭 jobs。如果多个地方都可能关闭同一个 channel,程序很容易 panic。

func produce(ctx context.Context, jobs chan<- Job) {
	defer close(jobs)

	for i := 0; i < 100; i++ {
		select {
		case <-ctx.Done():
			return
		case jobs <- Job{ID: i}:
		}
	}
}

worker 看到 jobs 被关闭后退出循环。这样职责很清楚:生产者决定什么时候没有新任务,worker 只负责消费。若任务来自消息队列或数据库,通常不关闭 channel,而是让 context 控制退出。

失败任务如何处理

示例里 worker 遇到错误只是记录日志,真实项目要决定失败任务的去向。一般有三种方式:直接丢弃、重试、写入失败表。入门项目可以先实现有限次数重试,避免无限循环。

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

重试间隔不要写成固定的零等待,否则下游服务抖动时 worker 会快速打满数据库或第三方 API。即使只是 200ms、400ms、600ms 这样的简单退避,也比立刻重试友好很多。

限制并发

worker 数量不是越多越好。CPU 密集任务可以接近 CPU 核数,IO 密集任务可以多一些,但最终要看数据库连接池、外部服务限流和机器资源。一个简单的经验是:先用小并发上线,再通过指标调大。

func runPool(ctx context.Context, n int, jobs <-chan Job) {
	var wg sync.WaitGroup
	wg.Add(n)

	for i := 0; i < n; i++ {
		id := i + 1
		go func() {
			defer wg.Done()
			worker(ctx, id, jobs)
		}()
	}

	wg.Wait()
}

如果 worker 内部还会访问数据库,记得让 SetMaxOpenConns 和 worker 数量匹配。二十个 worker 配五个数据库连接,很多任务会卡在等待连接;五十个 worker 配五十个连接,又可能把数据库压垮。并发控制要从整条链路看。

退出时的可观测性

服务关闭时,只打印“退出了”信息不够。最好记录收到信号、停止取新任务、等待 worker、剩余任务数量、最终耗时。排查发布卡住时,这些日志比猜测有用。

started := time.Now()
log.Println("worker pool stopping")
cancel()
wg.Wait()
log.Printf("worker pool stopped in %s", time.Since(started))

如果使用队列系统,还可以在退出前停止拉取新消息,等正在处理的消息 ack 完成。不要在收到信号后继续抢新任务,否则发布窗口会被拉长。

小结

Go worker 要能体面退出,循环里必须监听 ctx.Done(),任务处理要传递 context,sleep 和 ticker 也要可取消。单个任务可以派生超时,多个 worker 用 WaitGroup 等待退出。

后台任务的正确性不只在“能处理任务”,也在“该停的时候能停”。把取消路径写进设计和测试,worker 才适合长期运行。

常见问题与解答

context 取消后任务还在运行怎么办?

context 取消只是发送信号,不会直接中断正在运行的代码。handler 需要在合适的检查点看 ctx.Done()。如果任务里有阻塞调用(如数据库查询、网络请求),要确保这些调用也接受 context。

WaitGroup 能不能和 context 取消一起用?

可以。先 cancel() 通知 goroutine 停止,再 wg.Wait() 等待它们真正退出。顺序不要反:如果先 wg.Wait()cancel(),而 goroutine 正在阻塞,你就永远等不到它完成。

workers 数量应该设多少?

没有固定答案。CPU 密集型任务可以接近 CPU 核数,IO 密集型可以更多,但最终取决于下游容量。推荐先小后大:从少量 worker 开始,根据监控逐步调整。

worker pool 完整示例

type JobPool struct {
    workers int
    wg      sync.WaitGroup
    jobs    chan Job
    results chan Result
}

func NewJobPool(workers int, queueSize int) *JobPool {
    return &JobPool{
        workers: workers,
        jobs:    make(chan Job, queueSize),
        results: make(chan Result, workers),
    }
}

func (p *JobPool) Start(ctx context.Context, handler func(context.Context, Job) Result) {
    for i := 0; i < p.workers; i++ {
        p.wg.Add(1)
        go func(id int) {
            defer p.wg.Done()
            for job := range p.jobs {
                select {
                case <-ctx.Done():
                    return
                default:
                    p.results <- handler(ctx, job)
                }
            }
        }(i)
    }

    go func() {
        p.wg.Wait()
        close(p.results)
    }()
}

回填任务队列的注意事项

如果 worker 从消息队列(如 RabbitMQ、Kafka、SQS)消费消息,context 取消时的处理更复杂:

  1. 先停止拉取新消息:在 shutdown 时先停止 Consume 调用。
  2. 等已取消息处理完成:给当前正在处理的消息一段时间。
  3. 再确认/拒绝/重入队:根据处理状态决定怎么回馈队列。

如果 worker 从数据库轮询任务,要注意:

  • 最后一次轮询可能取到很多任务
  • 不要处理到一半就强退,可能导致任务状态不一致
  • 考虑给任务加 “处理中” 状态,重启后可以继续处理

实践练习

完成以下练习以巩固所学知识:

  1. 阅读 Go 官方文档相关章节
  2. 编写一个完整的示例程序
  3. 为示例程序编写单元测试
  4. 使用 go testgo benchmark 验证实现
  5. 尝试优化内存分配和运行时间

推荐阅读

真实项目用例

在实际团队协作中,下面是几个推荐的工作流:

代码审查清单

  • 函数是否处理了所有 error 返回值
  • 并发代码是否有明确的退出路径和 WaitGroup
  • 用户输入是否经过校验和清洗
  • 敏感配置是否通过环境变量或加密存储注入
  • 测试是否覆盖了正常路径和至少一个错误路径
  • 日志是否包含足够的上下文信息但不泄露敏感数据
  • 接口设计是否符合最小接口原则

CI/CD 集成建议

  • 每次提交前运行 go fmt ./...
  • CI 中运行 go vet ./...golangci-lint run
  • 单元测试使用 go test -race ./... 检测数据竞争
  • 关键路径的 benchmark 加入回归测试
  • 使用 go mod verify 确保依赖完整性

性能调优检查点

  • 使用 pprof 分析 CPU 和内存使用
  • 关注 benchmark 的 allocs/op,减少高频路径的堆分配
  • 检查数据库查询是否使用索引
  • 确认外部 HTTP 调用有合理的超时设置
  • 缓存热点数据,但注意缓存一致性和过期策略

面试高频考点

如果你正在准备 Go 相关面试,以下概念是高频考点:

  1. goroutine 和线程的区别
  2. channel 的缓冲和非缓冲用法
  3. defer 的执行顺序和与返回值的关系
  4. map 的并发不安全性和解决方案
  5. interface 的隐式实现和类型断言
  6. slice 的底层数组和 append 机制
  7. GC 的基本原理和调优参数
  8. context 的使用场景和超时控制
  9. error 的包装和 errors.Is/errors.As
  10. sync.Mutex vs sync.RWMutex vs atomic

掌握这些概念意味着你具备了独立开发 Go 服务的基础能力。继续在实际项目中磨练,你会越来越熟悉 Go 的工程风格和最佳实践。

常见问题(FAQ)

Q: 这个特性在实际项目中真的有用吗?
A: 是的。本文介绍的技术来源于真实后端开发场景。无论是标准库工具还是工程实践,在日常服务开发中都会反复用到。

Q: Go 版本会影响示例代码吗?
A: 本文代码主要针对 Go 1.20+ 编写。较新版本(如 1.22、1.23)的语法可能有微调,但核心概念保持不变。如有版本差异,文中会特别说明。

Q: 学习 Go 应该先学标准库还是直接上框架?
A: 强烈建议先学标准库。框架是对标准库的封装和扩展。只有理解了标准库的能力边界,才能正确选择和使用框架,也才能在框架出问题时快速定位。

Q: 代码里的错误处理为什么都是显式的 if err != nil
A: 这是 Go 的设计哲学。显式错误处理让失败路径清晰可见,不会隐藏在任何 try-catch 之后。习惯了之后,你会发现这种写法实际上降低了排查错误的难度。

Q: 并发相关代码怎么测试?
A: 使用 Go 内置的 -race 标志检测数据竞争:go test -race ./...。结合 sync.WaitGroupcontext.WithTimeout 编写有退出路径的并发测试,避免 goroutine 泄漏。

常见坑与避坑指南

  1. 不要信任用户输入:无论表单、JSON、Cookie 还是 HTTP Header,都当作不可信数据处理,做校验和转义。
  2. 资源要释放:文件、数据库连接、HTTP 响应体都要及时关闭。defer 是一个好习惯。
  3. 不要忽略错误:即使 defer file.Close() 可能返回错误,至少记录日志。完全忽略错误是 bug 的温床。
  4. 不要滥用 goroutine:每个 goroutine 都要有明确的退出路径。使用 sync.WaitGroupcontext 管理生命周期。
  5. 不要硬编码配置:端口、路径、超时时间、密钥都应该从配置读取,让程序适应不同环境。
  6. 不要过早优化:先让代码正确和可读,再用 benchmark 和 profile 找到真正的热点。

延伸阅读与实践建议

读完本文后,建议完成以下实践:

  1. 把文中所有示例代码在自己的机器上跑一遍
  2. 给示例代码补充错误分支的测试用例
  3. 尝试基于本文内容构建一个小型完整项目
  4. 在 review 他人的 Go 代码时,检查本文提到的边界是否被覆盖
  5. 订阅 Go 官方博客,关注语言演进和最佳实践更新

参考资源

  • Go 官方网站:https://go.dev/
  • Go 标准库文档:https://pkg.go.dev/std
  • Go by Example:https://gobyexample.com/
  • Effective Go:https://go.dev/doc/effective_go
  • Go 常见问题:https://go.dev/doc/faq
  • Go 项目实战社区案例和开源项目源码

本文力求在讲解技术细节的同时兼顾工程实用性。Go 语言的设计简洁但不简单,掌握它需要持续的实践和反思。希望这篇文章能成为你学习道路上的一个可靠参考。

小结补充

Worker 的取消路径和正常执行路径同等重要。没有取消机制的 goroutine 会在服务关闭时成为孤儿进程,占着资源不释放。养成写 ctx.Done() 分支的习惯,你的并发代码会可靠得多。测试时也要专门验证取消路径,确保 worker 能在规定时间内干净退出。

继续阅读

探索更多技术文章

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

全部文章 返回首页

「golang」更多文章

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