后台 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 取消时的处理更复杂:
- 先停止拉取新消息:在 shutdown 时先停止
Consume调用。 - 等已取消息处理完成:给当前正在处理的消息一段时间。
- 再确认/拒绝/重入队:根据处理状态决定怎么回馈队列。
如果 worker 从数据库轮询任务,要注意:
- 最后一次轮询可能取到很多任务
- 不要处理到一半就强退,可能导致任务状态不一致
- 考虑给任务加 “处理中” 状态,重启后可以继续处理
实践练习
完成以下练习以巩固所学知识:
- 阅读 Go 官方文档相关章节
- 编写一个完整的示例程序
- 为示例程序编写单元测试
- 使用
go test和go benchmark验证实现 - 尝试优化内存分配和运行时间
推荐阅读
- Go 官方博客: https://go.dev/blog/
- Effective Go: https://go.dev/doc/effective_go
- Go by Example: https://gobyexample.com/
- Go 标准库文档: https://pkg.go.dev/std
真实项目用例
在实际团队协作中,下面是几个推荐的工作流:
代码审查清单
- 函数是否处理了所有 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 相关面试,以下概念是高频考点:
- goroutine 和线程的区别
- channel 的缓冲和非缓冲用法
- defer 的执行顺序和与返回值的关系
- map 的并发不安全性和解决方案
- interface 的隐式实现和类型断言
- slice 的底层数组和 append 机制
- GC 的基本原理和调优参数
- context 的使用场景和超时控制
- error 的包装和 errors.Is/errors.As
- 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.WaitGroup 和 context.WithTimeout 编写有退出路径的并发测试,避免 goroutine 泄漏。
常见坑与避坑指南
- 不要信任用户输入:无论表单、JSON、Cookie 还是 HTTP Header,都当作不可信数据处理,做校验和转义。
- 资源要释放:文件、数据库连接、HTTP 响应体都要及时关闭。
defer是一个好习惯。 - 不要忽略错误:即使
defer file.Close()可能返回错误,至少记录日志。完全忽略错误是 bug 的温床。 - 不要滥用 goroutine:每个 goroutine 都要有明确的退出路径。使用
sync.WaitGroup和context管理生命周期。 - 不要硬编码配置:端口、路径、超时时间、密钥都应该从配置读取,让程序适应不同环境。
- 不要过早优化:先让代码正确和可读,再用 benchmark 和 profile 找到真正的热点。
延伸阅读与实践建议
读完本文后,建议完成以下实践:
- 把文中所有示例代码在自己的机器上跑一遍
- 给示例代码补充错误分支的测试用例
- 尝试基于本文内容构建一个小型完整项目
- 在 review 他人的 Go 代码时,检查本文提到的边界是否被覆盖
- 订阅 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 能在规定时间内干净退出。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。