Go 基础体系 · 第 75/113 篇。示例统一基于 Go 1.26.4;核心片段可能省略 package 与 import,完整程序可直接按文中结构运行。

Go Asynq 异步任务:Redis 队列、重试、定时与唯一任务

本文以 Go 1.26.4 和稳定版 github.com/hibiken/asynq v0.25.1 为基准。Asynq 把任务持久化到 Redis,由独立 worker 拉取执行,适合邮件、图片处理、报表生成和外部系统同步。它提供的是后台任务队列而非事件日志:任务通常面向“完成一次业务动作”,Kafka 一类日志则面向可重放事件流、分区顺序和多消费者订阅。

Asynq 的交付语义接近至少一次。worker 可能在副作用完成后、确认成功前崩溃,因此同一任务可能再次执行。重试、Unique 和 Redis 原子迁移能降低丢失或重复概率,却不能替业务完成端到端幂等。

1. 组件模型与数据路径

生产者用 ClientTask 写入 Redis;Server 的 fetcher 从一个或多个队列取任务,processor 调用 Handler,成功后删除或保留结果,失败后计算下一次重试时间。ServeMux 只负责把任务类型映射到处理器。SchedulerPeriodicTaskManager 负责产生任务,本身不执行任务。

producer -> Client -> Redis pending/scheduled/retry
                         | atomic lease/dequeue
                         v
                    Server workers -> Handler -> database/API
                         | success       | error
                         v               v
                  completed/delete   retry/archived

ClientServerScheduler 都持有连接或后台 goroutine,应在进程装配层创建并明确关闭。不要在每个 HTTP 请求中创建 Client;也不要在 init 中启动 Server。Redis 是控制面和任务状态存储,但业务事实仍应写入业务数据库。

2. Task、类型与 payload 契约

asynq.NewTask(typename, payload) 创建不可变任务描述。类型名是生产者与消费者之间的路由契约,宜使用 领域:动作:v版本;payload 只放标识、幂等键和执行所需的小型快照,不放大文件、访问令牌或完整用户资料。

const TypeDeliverReport = "report:deliver:v1"

type DeliverPayload struct {
	ReportID      string `json:"report_id"`
	IdempotencyID string `json:"idempotency_id"`
}

func NewDeliverTask(p DeliverPayload) (*asynq.Task, error) {
	if p.ReportID == "" || p.IdempotencyID == "" {
		return nil, errors.New("report_id and idempotency_id are required")
	}
	body, err := json.Marshal(p)
	if err != nil {
		return nil, fmt.Errorf("marshal deliver payload: %w", err)
	}
	return asynq.NewTask(TypeDeliverReport, body), nil
}

JSON 字段显式写 tag。变更字段含义时新增任务版本并在迁移期同时注册旧处理器;不能假设队列在发布瞬间为空。payload 大会增加 Redis 内存、复制、持久化和网络成本,大对象应进入对象存储,任务只传受权限保护的对象 ID。

3. Client 生命周期与入队 API

NewClient(RedisClientOpt) 创建可并发复用的生产者。EnqueueContext 接收调用方 context,返回的 TaskInfo 包含 ID、队列和状态。入队失败必须返回业务入口;若“业务数据提交”和“任务入队”必须原子一致,应采用事务 Outbox,而不是数据库提交后直接调用 Redis。

client := asynq.NewClient(asynq.RedisClientOpt{
	Addr:     "redis:6379",
	Username: os.Getenv("REDIS_USER"),
	Password: os.Getenv("REDIS_PASSWORD"),
	DB:       2,
})
defer client.Close()

task, err := NewDeliverTask(DeliverPayload{
	ReportID: "r-42", IdempotencyID: "report:r-42:deliver:v1",
})
if err != nil {
	return err
}
info, err := client.EnqueueContext(ctx, task,
	asynq.Queue("critical"),
	asynq.MaxRetry(8),
	asynq.Timeout(45*time.Second),
	asynq.Retention(24*time.Hour),
)
if err != nil {
	return fmt.Errorf("enqueue report delivery: %w", err)
}
logger.InfoContext(ctx, "task enqueued", "task_id", info.ID, "queue", info.Queue)

ProcessInProcessAt 让任务先进入 scheduled;Deadline 给绝对截止时刻,Timeout 给单次执行预算。两者是执行控制,不保证下游事务被强制终止。TaskID 可由调用者提供,但冲突会导致入队错误,应先定义 ID 的业务稳定性。

4. 状态机与至少一次语义

任务通常经历 pending、active、scheduled、retry、archived 或 completed。pending 等待消费;active 表示已被 worker 领取;scheduled 等待首次执行;retry 等待重试;超过重试上限进入 archived;设置 retention 的成功任务可在 completed 暂存,否则成功后删除。

这些状态不是业务事务。worker 在发送邮件后、标记成功前退出,租约恢复会让任务重做;网络超时也只表示“不知道对方是否完成”。因此不能把 active 当成业务锁,不能把 completed 当成业务账本,更不能通过“查询队列里有没有同 ID”证明动作未发生。

服务异常退出后,Asynq 会恢复未完成任务,这是可靠性的来源,也是重复执行的来源。状态监控应关注数量与驻留时间,而非只看瞬时长度:retry 数量稳定但最老任务不断变老,同样意味着故障。

5. Server、队列权重与并发

NewServer 配置并发、队列、错误处理和重试策略。Concurrency 是当前 Server 进程的最大并行处理数,多副本总并发约为各实例之和。队列 map 的值是权重,不是硬配额;启用 StrictPriority 后会优先耗尽高优先级队列,可能让低优先级饥饿。

server := asynq.NewServer(redisOpt, asynq.Config{
	Concurrency: 24,
	Queues: map[string]int{
		"critical": 6,
		"default":  3,
		"bulk":     1,
	},
	StrictPriority: false,
	ShutdownTimeout: 30 * time.Second,
	RetryDelayFunc: func(n int, err error, task *asynq.Task) time.Duration {
		if errors.Is(err, errRateLimited) {
			return time.Minute
		}
		return asynq.DefaultRetryDelayFunc(n, err, task)
	},
})

并发度由下游容量决定,不由 CPU 数量机械决定。一个任务同时占数据库连接、HTTP 连接和内存;24 个 worker 加 5 个副本可能形成 120 个并发。为不同资源等级拆队列或拆 worker 部署,比让大任务堵住短任务更可控。权重不能替代租户限流。

6. ServeMux、Handler 与错误边界

处理器实现 ProcessTask(context.Context, *asynq.Task) error。先解码、校验和检查任务版本,再调用业务服务。可重试故障返回带上下文的 error;永久输入错误包装 asynq.SkipRetry。错误在 Handler 返回,由 Server 的统一错误边界记录一次,避免底层又记录又返回。

type DeliverHandler struct{ service *DeliveryService }

func (h *DeliverHandler) ProcessTask(ctx context.Context, task *asynq.Task) error {
	var p DeliverPayload
	if err := json.Unmarshal(task.Payload(), &p); err != nil {
		return fmt.Errorf("decode deliver payload: %w: %v", asynq.SkipRetry, err)
	}
	if p.ReportID == "" || p.IdempotencyID == "" {
		return fmt.Errorf("validate deliver payload: %w", asynq.SkipRetry)
	}
	if err := h.service.Deliver(ctx, p.ReportID, p.IdempotencyID); err != nil {
		return fmt.Errorf("deliver report %q: %w", p.ReportID, err)
	}
	return nil
}

mux := asynq.NewServeMux()
mux.Handle(TypeDeliverReport, &DeliverHandler{service: service})
if err := server.Run(mux); err != nil {
	return fmt.Errorf("run asynq server: %w", err)
}

HandleFunc 适合短处理器,结构体 Handler 更容易注入依赖。注册相同 pattern 要在启动测试中发现。panic 应由恢复边界转成失败,但 panic 往往代表程序缺陷;恢复后仍需告警和修复,不能把它当普通重试机制。

7. 超时、取消与优雅关闭

Asynq 为一次处理提供 context;超时、deadline 或关停会发出取消信号。context 不会杀死 goroutine,Handler 的 SQL、HTTP 和阻塞 channel 必须使用 ctx,CPU 循环也要周期检查 ctx.Err()。另开 goroutine 后立刻返回会让 Server 错误地确认任务成功。

Shutdown 停止取新任务并等待 active 任务到 ShutdownTimeoutStopShutdown 的职责要按库 API 组织。进程收到信号后,先摘除 readiness,再停生产流量和 scheduler,随后关闭 worker,最后关闭业务连接、Client 和日志。过短退出宽限会增加重复执行,过长又会拖慢滚动发布。

ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()

server.Start(mux)
<-ctx.Done()
server.Shutdown()
if err := client.Close(); err != nil {
	return fmt.Errorf("close asynq client: %w", err)
}

长任务应分片并持久化进度,取消后让当前小单元回滚或安全提交。不要用 context.Background() 绕过关停;需要跨请求可靠延续的工作,本来就应先入队。

8. 重试分类、退避与归档

重试只对瞬时故障有效:连接重置、明确的 429/503、短期锁冲突可重试;字段非法、资源明确不存在、权限拒绝通常永久失败。未知提交状态要先查询或依赖幂等键,不能直接重复有副作用请求。MaxRetry 表示最大重试次数,不等于总尝试次数,运维文档应统一口径。

退避应指数增长并带抖动,避免依赖恢复时所有任务同时冲击。每次尝试还受总业务时限约束:三天后发送“注册验证码”即使技术上成功也没有意义。永久错误使用 SkipRetry,达到上限的任务进入 archived;归档不是垃圾桶,需要告警、查看、修复、重放和审计流程。

重放前确认处理器仍兼容旧 payload,并为批量重放设置速率。直接把十万归档任务全部重新排队,可能把已经恢复的数据库再次压垮。

9. Unique 与真正的业务幂等

asynq.Unique(ttl) 在指定窗口内拒绝等价任务,适合抑制按钮连点和短期重复调度。唯一锁有 TTL,会过期;任务完成或状态变化也可能释放。payload 的字节差异、类型或队列差异会影响唯一性。它不是永久业务约束。

_, err := client.EnqueueContext(ctx, task,
	asynq.Queue("critical"),
	asynq.Unique(10*time.Minute),
)
if errors.Is(err, asynq.ErrDuplicateTask) {
	return nil // 产品契约允许把窗口内重复请求视为已接受
}
if err != nil {
	return fmt.Errorf("enqueue unique delivery: %w", err)
}

真正幂等应在副作用所在系统建立。例如数据库表对 idempotency_key 加唯一索引,在同一事务中写处理结果;重复任务读取并返回原结果。外部 HTTP API 若支持 idempotency key,应把稳定键传过去。发送邮件这类外部动作很难与本地 DB 原子提交,可使用提供方幂等能力、发送账本和可接受重复的产品设计共同降低风险。

10. 定时、延迟与 Scheduler 高可用

一次延迟使用 ProcessIn/ProcessAt;重复计划使用 Scheduler.Register。Scheduler 到点只是创建普通任务,后续执行仍遵守队列、重试和幂等语义。

scheduler := asynq.NewScheduler(redisOpt, &asynq.SchedulerOpts{
	Location: time.FixedZone("CST", 8*60*60),
})
entryID, err := scheduler.Register(
	"0 3 * * *",
	asynq.NewTask("report:daily:v1", nil),
	asynq.Queue("bulk"),
	asynq.Unique(23*time.Hour),
)
if err != nil {
	return fmt.Errorf("register daily report: %w", err)
}
logger.Info("schedule registered", "entry_id", entryID)

显式使用 IANA 时区优于固定偏移;有夏令时的地区会出现不存在或重复的当地时间。多副本同时运行静态 Scheduler 可能重复入队,是否部署单副本、使用 leader election,或依赖 Unique 与业务幂等必须明确。动态计划适合 PeriodicTaskManager,配置同步失败时要保留最后成功快照并告警。

11. Redis 持久性、拓扑与故障边界

Redis 断开时生产者无法入队、worker 无法迁移状态。选择单机、Sentinel 或 Cluster 前先核对 Asynq 版本支持的连接选项和操作约束。配置连接超时、TLS、认证和专用 DB/实例;不要与可随意淘汰的缓存共享 maxmemory-policy。任务键被 eviction 等同数据丢失。

RDB/AOF、复制和故障转移仍可能存在已确认写入的丢失窗口。需要“业务提交后任务绝不遗漏”时,事务 Outbox 才是可靠边界:业务事务写 outbox,relay 可重复投递 Asynq,消费者幂等。Redis 备份恢复会让旧任务重现,因此恢复演练必须验证重复处理。

Redis 凭据使用 Secret 挂载或受控环境注入,不写 payload、日志或镜像。启用网络隔离和最小权限;管理 UI 不能直接暴露公网,重放、删除和归档操作需要认证、授权与审计。

12. Inspector、指标与诊断路径

Inspector 可查看队列、任务状态和归档任务,适合管理工具而非请求热路径。诊断“任务没执行”按固定顺序:确认入队返回和 queue;查看 pending/scheduled;确认 worker 订阅相同队列;检查 active 是否卡住;再看 retry/archived 的最后错误。不要一开始就删除或重放任务。

inspector := asynq.NewInspector(redisOpt)
defer inspector.Close()

info, err := inspector.GetTaskInfo("critical", taskID)
if err != nil {
	return fmt.Errorf("inspect task %q: %w", taskID, err)
}
fmt.Printf("type=%s state=%s retried=%d last_error=%q\n",
	info.Type, info.State, info.Retried, info.LastErr)

指标至少包括各状态数量、最老 pending/retry 年龄、入队与处理速率、成功/重试/归档率、处理 P95/P99、Redis 延迟和连接错误。label 使用任务类型、队列和稳定错误码,不能使用 task ID、用户 ID 等无界值。日志记录 task ID、类型、尝试次数、幂等键哈希和 trace 关联,不记录完整 payload。

13. 测试策略与可重复验证

纯单元测试直接构造 Task 调用 Handler,覆盖合法 payload、坏 JSON、取消、永久错误、瞬时错误和重复幂等键。集成测试启动隔离 Redis,验证入队、消费、重试、scheduled、Unique 和关停。时间相关测试尽量用短而有余量的 deadline,不依赖秒级 Sleep 猜测调度。

func TestDeliveryIsIdempotent(t *testing.T) {
	task, err := NewDeliverTask(DeliverPayload{
		ReportID: "r-42", IdempotencyID: "delivery-42",
	})
	if err != nil {
		t.Fatal(err)
	}
	store := NewDeliveryStore()
	var group sync.WaitGroup
	for range 20 {
		group.Go(func() {
			if err := store.ProcessTask(t.Context(), task); err != nil {
				t.Error(err)
			}
		})
	}
	group.Wait()
	if got := store.Count(); got != 1 {
		t.Errorf("delivery count = %d, want 1", got)
	}
}
go test ./...
go test -race -count=20 ./...
go vet ./...
redis-cli -u "$REDIS_URL" --scan --pattern 'asynq:*'

race detector 只能发现被执行路径的数据竞争,不能证明 Redis 与数据库的分布式原子性。故障测试应在副作用提交前后分别终止 worker,验证重复任务不会重复扣款或破坏状态。

14. 性能、容量与反压

先测单任务的 CPU、内存、下游连接和耗时,再算并发。Little 定律可用于粗估 active 数:吞吐每秒 100、平均耗时 0.2 秒,稳态约需 20 个并发;还要为尾延迟和故障留余量。无限提高 Concurrency 只会把排队从 Redis 搬到数据库连接池。

大任务拆成有界批次,但不要为每条记录生成数百万微任务而没有生产速率限制。生产者需要配额和背压:队列年龄超过阈值时拒绝非关键请求、降采样或转入批处理。压缩 payload 会消耗 CPU 且妨碍诊断,优先缩小契约。批量查询、连接复用和避免重复 JSON 编解码应通过 benchmark/profile 验证。

优先级是业务策略:critical 持续满载时 bulk 是否允许饥饿、每个租户能占多少 worker、下游降级时哪些任务应暂停,都要形成配置和告警,而不是只写一个队列权重 map。

15. 安全与生产部署清单

任务是内部输入也必须验证,因为旧生产者、重放工具和被攻陷服务都可能写坏 payload。对对象 ID 再做租户授权,不信任 payload 中的角色。下载 URL 要限制 scheme、host、重定向、响应大小和内网地址,防止 SSRF;文件处理放隔离容器并限制 CPU、内存、临时空间。

部署时固定 Go 与 Asynq 版本,先在测试 Redis 验证升级兼容性。worker 使用独立 Deployment,readiness 表示能否接新任务,liveness 不应因下游短抖动反复杀进程。设置终止宽限大于 ShutdownTimeout,用 PodDisruptionBudget 避免同时驱逐全部 worker。迁移任务 payload 采用先部署兼容消费者、再切生产者、最后清理旧处理器的顺序。

Asynq 的合理生产边界是:用 Redis 状态机承接可重试后台工作,用受控并发和 context 管理执行生命周期,用归档与指标运营失败;用业务数据库唯一约束、事务 Outbox 和下游幂等键解决端到端可靠性。 把这两层职责分开,才能在重启、超时和重放发生时仍得到可解释结果。


系列导航与关联阅读

官方资料

本文依据 Go 官方规范、标准库文档和 Go 官方博客重新梳理;正文与示例由 WR BLOG 编写。