用 channel + worker pool 实现轻量级异步任务队列:基于纯 Go 内存,适用于单机、低延迟、无持久化场景;需设置带缓冲的 taskCh 和固定数量 worker 以控并发、防 goroutine 泛滥。

用 channel + worker pool 实现轻量级异步任务队列
纯 Go 内存级任务队列适合单机、低延迟、无持久化要求的场景,比如内部通知触发、日志归档、配置热加载。它不依赖 Redis 或其他中间件,启动快、调试直观,但任务会随进程退出丢失。
关键点是控制并发和防止 goroutine 泛滥:
-
taskCh必须带缓冲,例如make(chan Task, 100),否则生产者可能被阻塞 - worker 数量要固定,用
for i := 0; i 启动,避免动态创建 - 每个
worker必须用for task := range taskCh循环消费,不能只取一次 - 任务函数内部建议包一层
defer func() { recover() }(),防止单个 panic 杀掉整个 worker
示例中常见错误:把 taskCh 声明为无缓冲 chan Task,导致高并发提交时主流程卡死;或忘记在 main 退出前 close(taskCh),导致 worker 永不退出。
对接 Redis 实现分布式可靠队列(TaskQ / Asynq)
多实例部署、任务不可丢、需重试/延时/去重时,必须换用持久化后端。Redis 是最常用选择,TaskQ 和 Asynq 都基于它,但抽象层级不同。
立即学习“go语言免费学习笔记(深入)”;
TaskQ 的 Queue 接口统一了 Push/Pop 行为,但实际使用时要注意:
-
RedisQueue默认用 List(LPUSH/BRPOP),简单但不支持消息确认与重入;若需 Exactly-Once 语义,得切到 Streams 模式(需手动配置) - TaskQ 不内置调度器,
Delay字段仅靠客户端等待,不是服务端定时投递;真要延时任务,得自己套一层time.AfterFunc或改用 Asynq - 并发消费者数由
Worker实例数量决定,不是靠 channel 缓冲大小——每个 Worker 独立连接 Redis,不存在“队列满”阻塞生产者的问题
容易踩的坑:直接拿 TaskQ 的内存队列 memqueue 上生产环境,重启后所有 pending 任务清零;或没设 Redis 连接池大小,高并发下连接耗尽报 connection refused。
任务结构体设计影响扩展性与可观测性
一个裸 func() 类型的任务无法记录状态、超时、重试次数或返回结果。真实项目里,Task 结构体至少要包含这些字段:
-
ID string:用于幂等去重和日志追踪,建议用uuid.NewString() -
Fn func(context.Context) error:接受context.Context,便于超时控制与取消传播 -
Retry int和MaxRetry int:配合中间件做指数退避,TaskQ 默认只重试 3 次且无退避 ResultCh chan(可选):若需同步获取结果,避免轮询 DB 或 Redis
别把参数直接塞进闭包,例如 func() { doX(a, b) } —— 这会让序列化、重试、日志打点全部失效。所有数据必须显式存在 Task 字段里,才能被中间件正确处理。
优雅关闭与上下文传递常被忽略
很多实现能跑,但一关服务就 panic 或任务中断。根本原因是没处理好生命周期:
- worker 启动时应接收
context.Context,并在select中监听ctx.Done(),而不是硬等range taskCh - 关闭流程必须是:先
close(taskCh)或cancel()ctx,再wg.Wait()等所有 worker 退出,最后才释放资源(如 Redis client) - HTTP handler 提交任务时,别直接用
req.Context()传给后台执行——请求结束,ctx 就 cancel 了;得用context.WithTimeout(context.Background(), 30*time.Second)
最隐蔽的问题:worker 里调用了另一个带 context 的外部服务(如 HTTP 调用),但没把父 context 传进去,导致超时设置失效,整个任务 hang 死。


















