用channel+worker pool实现轻量级任务队列需:taskch带缓冲、worker固定数量且持续消费、加recover防panic、退出前close;goroutine限流用带容量的sem channel。

用 channel + worker pool 实现轻量级任务队列
单机、低延迟、无持久化需求的场景下,纯内存方案最直接。但容易卡死或泄漏,关键在初始化和消费逻辑。
-
taskCh必须带缓冲,比如make(chan Task, 100);声明为无缓冲chan Task会导致高并发提交时主流程阻塞 - worker 数量要固定启动,不能动态增减;每个 worker 必须用
for task := range taskCh持续消费,只取一次就退出会丢任务 - 每个 worker 内部建议加
defer func() { recover() }(),否则一个 panic 会让整个 goroutine 退出,后续任务无人处理 - 程序退出前记得
close(taskCh),否则 worker 会永远阻塞在range上
goroutine 泛滥时怎么限流
直接写 go fn() 看似简单,但请求峰值一来,几百个 goroutine 同时跑,内存和调度开销立刻飙升。这不是异步,是自爆。
- 用带容量的
semchannel 控制并发数,例如sem := make(chan struct{}, 3),每次执行前sem ,结束后 <code> - 配合
sync.WaitGroup等待所有任务完成,避免提前退出 - 别依赖
time.Sleep做等待——它不释放 CPU,且无法响应取消;改用context.WithTimeout+select配合
为什么 Task 结构体必须带 context.Context
没有 context.Context 的任务,等于没刹车。超时、取消、链路追踪全靠它传递。
- 函数签名应为
func(ctx context.Context) error,不是裸func() - worker 执行时需传入带 timeout 的 context,例如
ctx, cancel := context.WithTimeout(parentCtx, 5*time.Second) - 任务内部所有 I/O 操作(HTTP 请求、DB 查询、Redis 调用)都得接收并传递该 context,否则超时无效
- 别把参数硬编码进闭包,如
go func() { doX(a, b) }()——这导致无法序列化、无法重试、日志里也看不到原始参数
什么时候必须换 Redis-based 队列
只要出现以下任一情况,channel 方案就得切走:进程重启后任务不能丢、需要重试、要延时执行、多实例部署。
-
asynq是首选:支持Retry、Delay、Timeout、Web UI 查看队列状态,底层用 Redis Streams 保证至少一次投递 - 别用
TaskQ的默认List模式上生产——它不支持消息确认,失败任务会丢失;切到Streams模式需手动配置 - Redis 连接池大小必须显式设置,否则高并发下
connection refused错误频发;asynq默认连接池是 10,通常不够,建议设为 50+ -
memqueue只适合开发调试,上线即踩坑:重启后所有 pending 任务清零,连告警都来不及发
真实项目里,最常被忽略的是任务结构体设计——ID 不唯一、没上下文、没重试计数、没日志埋点,结果就是出问题时既查不到谁发起的,也看不出重试了几次,更没法做幂等判断。











