很多业务请求并不需要在 HTTP 请求返回前完成全部工作。发送邮件、生成缩略图、同步第三方订单,这些操作可能很慢,也可能因为外部服务不稳定而失败。如果把它们直接放在请求处理函数里,用户要等待,服务也更容易在高峰期被拖垮。任务队列解决的正是这个问题:请求只负责登记一项工作,后台 Worker 稍后完成它。
本系列文章将会学习一个极简的 Go 异步任务队列 Asynq,本文先不急着钻进所有内部细节,而是用一个可以运行的例子认识 Asynq,再沿着这个例子看清 Client、Redis、Server 和 Handler 各自做了什么。
Asynq 是什么
Asynq 是一个用 Go 编写的、基于 Redis 的分布式任务队列。它把系统分成两端:
- Client 是生产者,把任务写入 Redis;
- Server 是消费者,从 Redis 取出任务并交给 Worker 执行;
- Handler 是业务代码,决定一项任务具体如何完成;
- Redis 保存任务及其状态,让不同进程、不同机器上的 Client 和 Server 能够协作。
最简的工作流程可以画成这样:
1 | 业务请求 |
这个模型带来几个直接的好处:请求线程不用等待耗时操作;任务可以暂存在 Redis 中等待处理;增加 Server 实例就能增加处理能力。Asynq 还提供延迟执行、重试、任务去重、队列优先级和运维查询等能力,后续文章会分别展开。
源码结构
第一次打开 Asynq 仓库,容易被根目录里许多 Go 文件吓到。其实它们不是一团混在一起的代码,而是按“公开 API、服务端组件、存储实现、工具和扩展”分层摆放。先掌握这张地图,后面的源码阅读会轻松很多。
1 | asynq/ |
根目录的 go.mod 是核心库。tools/ 和 x/ 各自有独立的 go.mod,这样 CLI、指标和限流扩展可以单独管理依赖,不会把工具库带进核心运行时。internal 是 Go 的内部包,仓库外的程序不能直接导入它;这让 Redis Key、Lua 脚本等实现细节可以在不承诺公共 API 兼容性的情况下演进。
阅读源码时,可以沿着下面这条路径前进:
- 先看
asynq.go,了解任务模型和 Redis 连接配置; - 再看
client.go与server.go,找到生产和消费的公开入口; - 接着进入
processor.go,观察 Server 如何取出任务并启动 Worker; - 需要追查数据落盘时,从
internal/base的Broker接口跳到internal/rdb; - 运维查询看
inspector.go,定时任务看scheduler.go,可选能力看x/; - 最后对照各组件旁边的
*_test.go,通过测试理解边界行为。
先把示例跑起来
准备 Redis
Asynq 要求 Redis 4.0 或更高版本。已经有 Redis 的读者可以直接使用;如果本机安装了 Docker,可以这样启动一个临时实例:
1 | docker run --name asynq-redis -p 6379:6379 -d redis:7 |
然后创建 Go 项目并安装 Asynq:
1 | mkdir asynq-demo && cd asynq-demo |
下面的代码分成 tasks、生产者和 Worker 三部分。实际项目也可以把它们放在不同的服务中,这正是任务队列常见的部署方式。
定义任务
一个 Asynq 任务由类型名和 payload 组成。类型名用于路由,payload 保存业务数据,通常使用 JSON 或 protobuf 编码。先创建 tasks/tasks.go:
1 | package tasks |
这里有一个值得注意的设计:NewEmailDeliveryTask 只创建内存中的 Task,还没有访问 Redis。任务何时入队由调用方决定,这使得同一种任务可以在不同业务流程中使用不同的队列或调度选项。
发送任务
生产者可以放在 Web 服务或命令行程序中。创建 producer/main.go:
1 | package main |
NewClient 接收的是 RedisConnOpt。最常用的 RedisClientOpt 连接单个 Redis;同一个接口还支持 Sentinel 和 Cluster。Client 可以被多个 goroutine 并发使用,程序退出时调用 Close 释放它创建的连接池。
如果任务需要稍后执行,可以在入队时增加选项:
1 | import "time" |
ProcessIn 表示相对当前时间延迟一段时间;如果业务已经确定了具体时刻,则使用 ProcessAt 指定绝对时间。两者都会让任务先进入 scheduled 状态,等到时间到达后再交给 Worker;这条路径会在后续的调度文章中详细说明。
启动 Worker
再创建 worker/main.go:
1 | package main |
在两个终端分别启动 Worker 和生产者(Worker 会一直运行,直到收到终止信号):
1 | # 终端一:启动 Worker |
Worker 会打印类似下面的日志:
1 | send email: user_id=42 template_id=welcome |
Run 会启动后台组件并等待终止信号,收到 SIGTERM 或 SIGINT 后执行优雅关闭。如果需要由测试或自己的生命周期管理代码控制,可以使用 Start 和 Shutdown,它们与 Run 使用的是同一套处理逻辑。
从 API 走到源码
示例跑通后,再回头看源码会清楚很多。Asynq 的核心路径并不长,关键在于每一层只承担一部分职责。
Task 描述任务
asynq.go 中的 Task 保存类型、payload、headers 和选项:
1 | type Task struct { |
NewTask 只是填充这些字段,并不会启动 goroutine,也不会写 Redis。这样做有两个好处:任务构造和任务投递相互独立;业务代码可以在真正入队前完成参数校验和序列化。
Client.Enqueue 负责入队
Client.Enqueue 使用后台上下文,实际工作交给 EnqueueContext。源码先合并任务创建时和入队时的选项,再根据时间和分组信息选择路径:
1 | opts = append(task.opts, opts...) |
对普通任务来说,最后会走 c.enqueue。如果指定了未来时间,就进入 scheduled;如果指定了 Group,则等待聚合。这里先记住这三个入口即可,聚合和调度的内部组件会在后文介绍。
选项有明确的覆盖规则:先使用 NewTask 里的选项,再使用 Enqueue 传入的选项;同一种选项出现多次时,最后一个生效。没有显式设置时,源码会使用默认值:最大重试次数为 25,处理超时时间为 30 分钟。
常用选项可以按业务意图理解:
| 业务意图 | 选项 | 作用 |
|---|---|---|
| 失败后再试几次 | MaxRetry(n) |
设置最大重试次数 |
| 放入指定队列 | Queue(name) |
将任务交给对应队列的 Worker |
| 限制单次处理时间 | Timeout(d) |
超时后结束本次处理并进入失败路径 |
| 指定最终截止时间 | Deadline(t) |
到达绝对时间后不再继续处理 |
| 稍后再处理 | ProcessIn / ProcessAt |
进入 scheduled,等待到期 |
| 防止重复入队 | Unique(ttl) |
在 TTL 内拒绝相同任务 |
Server 组织处理流程
NewServer 会把连接配置转换成 redis.UniversalClient,并创建多个后台组件。为了先建立整体认识,可以按职责把这些组件归为三组:
processor负责取任务、创建 Worker goroutine、调用 Handler;heartbeater、recoverer、forwarder、janitor等后台组件维护租约、恢复、延迟任务和过期数据;healthchecker、subscriber、syncer等组件分别处理健康检查、取消消息和失败后的状态同步。
调用 srv.Start(mux) 时,这些组件会依次启动;processor 的并发数由 Config.Concurrency 控制。如果没有设置有效值,Server 会使用当前进程可用的 CPU 数量。Queues 未配置时只处理名为 default 的队列。
ServeMux 路由任务
Server 不需要知道每种业务任务的细节,它只依赖 Handler 接口:
1 | type Handler interface { |
Asynq 用一个函数类型把普通函数适配成 Handler。server.go 中的定义很简单:
1 | type HandlerFunc func(context.Context, *Task) error |
函数类型本身也可以定义方法。因此,只要函数签名符合要求,就能通过 asynq.HandlerFunc(fn) 转换为 Handler,不必额外声明一个结构体。前面的 mux.HandleFunc 正是帮我们完成了这次转换。
这种“函数类型 + 方法”的适配器写法在 Go 中很常见。标准库的 http.HandlerFunc 让普通函数变成 http.Handler,sort.Interface 也经常通过自定义类型加方法来适配。它保留了接口带来的可替换性,又让简单场景只需要写一个函数。
在此基础上,ServeMux 把任务类型映射到 Handler,并且支持最长前缀匹配:注册 image 和 image:resize 时,类型为 image:resize:thumbnail 的任务会优先匹配更具体的 image:resize。这和 net/http.ServeMux 的使用体验很接近,也允许通过 Use 添加日志、鉴权或限流中间件。
Redis 保存任务状态
Asynq 并不是把任务直接序列化成一条普通字符串就结束了。一个任务既有自己的内容,也需要被快速地按状态查询、按时间调度,还要能在处理中被找到并续租。因此,Asynq 把 任务实体 和 状态索引 拆开保存:任务内容只有一份,各种队列结构只保存任务 ID。
Redis Key 统一按队列组织
Redis Key 可以理解为 Redis 中一条数据的名称,作用类似数据库表中的主键。写入、读取或删除数据时,客户端都要通过 Key 找到目标。例如,HSET some-key field value 中的 some-key 就是这组 Hash 数据的 Key。Asynq 用不同的 Key 区分不同队列、不同任务和不同状态集合,Worker、Inspector 以及后台维护组件都通过这些 Key 找到需要处理的数据。
源码在 internal/base/base.go 中集中生成这些 Key。队列名会出现在统一的前缀中,例如默认队列的前缀是:
1 | asynq:{default}: |
其中的 {default} 不只是为了好看,它是 Redis Cluster 的 hash tag。属于同一个队列的 Key 会尽量落在同一个 slot,Lua 脚本才能把这些 Key 一起作为参数执行。基于这个前缀,一个任务 ID 为 abc123 的实体 Key 是:
1 | asynq:{default}:t:abc123 |
任务实体和状态索引
以队列 default 为例,常见 Key 可以这样理解:
| Redis Key | 类型 | 保存内容 | 用途 |
|---|---|---|---|
asynq:{default}:t:<id> |
Hash | protobuf 编码的消息、state、错误信息等 |
任务实体,贯穿整个生命周期 |
asynq:{default}:pending |
List | 待处理任务 ID | Worker 从这里取任务 |
asynq:{default}:active |
List | 正在处理的任务 ID | 记录当前正在执行的任务 |
asynq:{default}:scheduled |
ZSet | 任务 ID,score 是执行时间 | 保存未来执行的任务 |
asynq:{default}:retry |
ZSet | 任务 ID,score 是下次重试时间 | 保存等待重试的任务 |
asynq:{default}:lease |
ZSet | 任务 ID,score 是租约过期时间 | 检查 Worker 是否还活着 |
asynq:{default}:completed |
ZSet | 已完成任务 ID,score 是完成时间 | 在保留期内供查询 |
asynq:{default}:archived |
ZSet | 归档任务 ID,score 是归档时间 | 保存失败且不再重试的任务 |
这些结构之间并不是互相复制任务内容。例如任务从 pending 进入 active 时,脚本只需把 ID 从一个 List 移到另一个 List,再更新任务 Hash 的 state 字段并写入 Lease。查询任务详情时,Inspector 根据 ID 找到同一个 Hash 即可。这样既避免了多份 payload,也让状态迁移的含义非常明确。
除了前面列出的状态 Key,Asynq 还维护了一些辅助 Key。它们不负责保存任务正文,而是分别服务于队列发现、统计、去重和分组:
asynq:queues是一个全局 Set,记录系统中出现过的队列名称。Inspector 先从这里发现队列,再查询每个队列的具体状态。asynq:{<queue>}:processed、failed记录累计数量;带日期后缀的processed:YYYY-MM-DD和failed:YYYY-MM-DD记录每日数量,并设置过期时间,避免统计数据无限增长。- 使用
Unique时,Asynq 会根据队列名、任务类型和 payload 计算摘要,生成类似asynq:{default}:unique:email:deliver:<hash>的 String Key。这个 Key 通过 TTL 充当“去重锁”,锁存在时相同任务不会再次入队。 - 使用 Group 时,
asynq:{<queue>}:groups保存该队列的组名,asynq:{<queue>}:g:<group>是保存组内任务 ID 的 ZSet;真正开始聚合时,还会在这个前缀下创建带聚合集合 ID 的临时 Key。
这些辅助 Key 与任务实体、状态索引使用同一套命名规则,因此一个队列的相关数据可以按前缀集中查找和维护。
真正入队时,internal/rdb 使用 Lua 脚本一次完成 检查任务 ID、写入任务内容、加入 pending 列表 这几步:
1 | if redis.call("EXISTS", KEYS[1]) == 1 then |
这段脚本的意义在于:检查和写入之间不会被另一个 Client 插入操作打断,任务不会出现 Hash 已写入但列表没有 ID 的中间状态。脚本返回 0 时,Asynq 会把它转换成任务 ID 冲突错误。
任务消息内部使用 protobuf 编码后存入 Redis;业务 payload 的编码方式则由应用决定,本文示例使用 JSON。Redis 的 List、ZSet、Hash、Pub/Sub 和 Lua 能力共同构成了 Asynq 的基础设施,go-redis/v9 负责在 Go 与 Redis 之间建立连接并执行这些操作。
项目依赖
从这个最小示例可以看出,Asynq 的运行时依赖并不多:
github.com/redis/go-redis/v9提供 Redis 单机、Sentinel 和 Cluster 的连接,以及 Lua、Pipeline、Pub/Sub 等操作;google.golang.org/protobuf用于编码保存在 Redis 中的内部任务消息;github.com/google/uuid在没有指定TaskID时生成任务 ID;- Go 标准库负责 Context、JSON 编码、TLS 和时间处理。
Cron 调度使用的 github.com/robfig/cron/v3、Prometheus 指标和分布式限流属于特定功能或扩展,并不是运行这个基本示例的前置条件。把依赖按功能拆开,后面阅读源码时就能知道某个库为什么只出现在某个模块中。
小结
到这里,可以用四句话概括 Asynq:
- Client 把“要做什么”写进 Redis,Server 决定“什么时候取、交给谁做”;
- Task 本身只是类型和数据,业务行为写在 Handler 中;
- ServeMux 解决任务类型到 Handler 的路由,Config 控制并发和队列;
- Redis 保存任务状态,Lua 脚本保证关键状态变更的一致性。
这个模型解释了 Asynq 为什么既适合单体应用中的后台任务,也适合多个服务、多个 Worker 实例共同消费。它同时也提醒我们:任务可能因为故障被再次执行,Handler 最好设计成幂等的;Redis 是必需的运行时依赖,部分 Lua 脚本在 Redis Cluster 下还存在兼容性限制。