后台任务失败后要不要重试?答案通常是"要,但不能乱重试"。发送通知、同步数据、生成报表、调用第三方支付接口等操作都可能遇到临时失败。如果完全不重试,一次网络抖动就会变成用户可见的故障;如果无限制重试,可能把一个小问题放大成队列堵塞或下游服务雪崩。
本文用发送通知任务作为贯穿示例,深入讲解一个完整的重试系统设计:错误分类、退避策略、重试循环、幂等性、死信队列、jitter 打散和可观测性。
任务模型与核心接口
首先定义重试任务的基本结构:
package main
import "context"
type NotifyJob struct {
ID string
UserID int64
Message string
Attempt int // 已尝试次数
}
type Notifier interface {
Send(ctx context.Context, userID int64, message string) error
}
// 模拟第三方推送服务
type PushNotifier struct{}
func (n *PushNotifier) Send(ctx context.Context, userID int64, message string) error {
// 调用第三方 API
return nil
}
Attempt 字段表示已经尝试了几次。真实队列系统(如 RabbitMQ、AWS SQS)通常会在消息元数据中维护这个计数,或者从数据库中读取。
错误分类:区分临时错误与永久错误
重试的前提是错误分类正确。不是所有错误都值得重试。
不应重试的永久错误:
- 参数格式错误(空消息、非法用户 ID)
- 权限错误(无权发送给该用户)
- 业务规则错误(消息长度超过限制)
- 资源不存在(用户已注销)
可以重试的临时错误:
- 网络超时
- 连接被拒绝(下游服务重启中)
- HTTP 503/502(服务过载或暂时不可用)
- 限流(429 Too Many Requests)
在 Go 中实现错误分类:
package main
import (
"errors"
"fmt"
)
var ErrPermanent = errors.New("permanent error")
func permanent(err error) error {
return fmt.Errorf("%w: %v", ErrPermanent, err)
}
func IsPermanent(err error) bool {
return errors.Is(err, ErrPermanent)
}
业务处理中的应用:
func HandleNotifyJob(ctx context.Context, notifier Notifier, job NotifyJob) error {
if strings.TrimSpace(job.Message) == "" {
return permanent(errors.New("empty message"))
}
if job.UserID <= 0 {
return permanent(errors.New("invalid user id"))
}
if err := notifier.Send(ctx, job.UserID, job.Message); err != nil {
return fmt.Errorf("send notification: %w", err)
}
return nil
}
空消息重试一百次也不会成功,所以直接标记为永久错误并进入失败处理流程。
重试循环:最大次数与退避等待
package main
import (
"context"
"fmt"
"time"
)
func RunWithRetry(ctx context.Context, job NotifyJob, fn func(context.Context, NotifyJob) error) error {
// 退避延迟:200ms, 1s, 3s, 5s
delays := []time.Duration{
200 * time.Millisecond,
1 * time.Second,
3 * time.Second,
5 * time.Second,
}
var err error
for attempt := 0; attempt <= len(delays); attempt++ {
job.Attempt = attempt + 1
err = fn(ctx, job)
if err == nil {
return nil
}
if IsPermanent(err) {
return err
}
if attempt == len(delays) {
break
}
select {
case <-time.After(delays[attempt]):
// 等待后重试
case <-ctx.Done():
return ctx.Err()
}
}
return fmt.Errorf("job failed after %d retries: %w", len(delays), err)
}
这段代码的关键设计决策:
- 最大次数有限:总尝试次数是
len(delays) + 1(初始 1 次 +len(delays)次重试) - 永久错误立即终止:不会浪费重试次数在不可能成功的任务上
- 监听 context:等待期间可以被取消,避免 goroutine 泄漏
- 退避逐渐增长:不是固定间隔,给下游恢复的时间
指数退避与 Jitter 策略
固定退避在大量任务同时失败时会形成"重试风暴":所有任务都在 1 秒后重试,给刚刚恢复的服务造成新的峰值压力。
更稳健的退避策略是指数退避加 jitter 打散:
package main
import (
"math"
"math/rand"
"time"
)
type BackoffConfig struct {
Initial time.Duration
Multiplier float64
MaxDelay time.Duration
MaxRetries int
}
func (c *BackoffConfig) Delay(attempt int) time.Duration {
if attempt < 0 {
return 0
}
// 指数计算
delay := float64(c.Initial) * math.Pow(c.Multiplier, float64(attempt))
if delay > float64(c.MaxDelay) {
delay = float64(c.MaxDelay)
}
// 添加 jitter(最多 +50% 的随机偏移)
jitter := delay * 0.5 * rand.Float64()
return time.Duration(delay + jitter)
}
var defaultBackoff = &BackoffConfig{
Initial: 200 * time.Millisecond,
Multiplier: 2.0,
MaxDelay: 30 * time.Second,
MaxRetries: 5,
}
典型的退避序列:
| 尝试次数 | 无 Jitter | 有 Jitter(范围) |
|---|---|---|
| 第 1 次 | 200ms | 200ms - 300ms |
| 第 2 次 | 400ms | 400ms - 600ms |
| 第 3 次 | 800ms | 800ms - 1200ms |
| 第 4 次 | 1.6s | 1.6s - 2.4s |
| 第 5 次 | 3.2s | 3.2s - 4.8s |
jitter 的目标不是精确,而是把重试时间打散,避免所有 worker 在同一时刻再次冲向下游。
幂等性:重试的安全基石
重试之前必须先问:重复执行会不会造成副作用?
有副作用的操作:
- 发送通知 -> 用户可能收到重复消息
- 扣款操作 -> 可能重复扣款
- 创建订单 -> 可能产生重复订单
- 积分发放 -> 用户积分异常增加
幂等性不是重试库能自动提供的,需要业务设计。常见方案:
方案一:发送日志表
package main
import "context"
type SendLog interface {
AlreadySent(ctx context.Context, jobID string) (bool, error)
MarkSent(ctx context.Context, jobID string) error
}
处理流程:
func IdempotentSend(ctx context.Context, log SendLog, notifier Notifier, job NotifyJob) error {
sent, err := log.AlreadySent(ctx, job.ID)
if err != nil {
return fmt.Errorf("check sent status: %w", err)
}
if sent {
return nil // 已处理,跳过
}
if err := notifier.Send(ctx, job.UserID, job.Message); err != nil {
return err
}
if err := log.MarkSent(ctx, job.ID); err != nil {
// 这里需要更严格的处理,可能需要用事务保证
return fmt.Errorf("mark sent: %w", err)
}
return nil
}
更严格的实现要用数据库的唯一约束保证并发安全:
CREATE TABLE notification_logs (
job_id VARCHAR(64) PRIMARY KEY,
user_id BIGINT NOT NULL,
sent_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP
);
用 job_id 作为主键,即使并发执行也不会重复插入。
方案二:外部系统自带幂等
有些系统(如 Stripe、支付宝)本身就提供了幂等性支持。你可以在请求头中传入幂等键:
func (n *PushNotifier) Send(ctx context.Context, userID int64, message string) error {
req, _ := http.NewRequestWithContext(ctx, http.MethodPost, n.URL, body)
req.Header.Set("Idempotency-Key", fmt.Sprintf("notify-%d-%s", userID, hash(message)))
// ...
}
死信队列:不要默默丢弃失败任务
当任务超过最大重试次数后,不应该悄悄丢失。常见做法是放入死信队列(Dead Letter Queue)或标记为 failed:
package main
import "context"
type FailedJobStore interface {
SaveFailed(ctx context.Context, job NotifyJob, reason string) error
}
type FailedJob struct {
JobID string
JobType string
UserID int64
Attempt int
LastError string
CreatedAt time.Time
FailedAt time.Time
}
失败记录至少包含以下信息:
- job id 和任务类型
- 尝试次数和最后错误信息
- 失败时间和创建时间
这样值班人员可以判断是临时依赖故障、数据错误还是代码 bug。如果没有失败归档,后台任务会变成黑洞:用户说没收到通知,你只能翻日志猜。
完整的重试任务处理器
将上述组件组合成一个完整的处理器:
package main
import (
"context"
"fmt"
"log/slog"
"time"
)
type JobProcessor struct {
Notifier Notifier
SendLog SendLog
FailedStore FailedJobStore
Metrics JobMetrics
}
func (p *JobProcessor) Process(ctx context.Context, job NotifyJob) error {
p.Metrics.SetQueueLength(1)
defer p.Metrics.SetQueueLength(0)
err := p.processWithRetry(ctx, job)
if err != nil {
p.Metrics.IncFailure()
if IsPermanent(err) {
p.saveFailed(ctx, job, fmt.Sprintf("permanent: %v", err))
} else {
p.saveFailed(ctx, job, fmt.Sprintf("exhausted retries: %v", err))
}
return err
}
p.Metrics.IncSuccess()
return nil
}
func (p *JobProcessor) processWithRetry(ctx context.Context, job NotifyJob) error {
backoff := defaultBackoff
for attempt := 0; attempt <= backoff.MaxRetries; attempt++ {
job.Attempt = attempt + 1
// 幂等性检查
err := IdempotentSend(ctx, p.SendLog, p.Notifier, job)
if err == nil {
return nil
}
if IsPermanent(err) {
slog.Warn("notify job permanent failure",
"job_id", job.ID,
"user_id", job.UserID,
"attempt", job.Attempt,
"err", err,
)
return err
}
if attempt == backoff.MaxRetries {
break
}
p.Metrics.IncRetry()
delay := backoff.Delay(attempt)
slog.Info("notify job retrying",
"job_id", job.ID,
"attempt", job.Attempt,
"delay_ms", delay.Milliseconds(),
)
select {
case <-time.After(delay):
case <-ctx.Done():
return ctx.Err()
}
}
return fmt.Errorf("max retries exceeded")
}
func (p *JobProcessor) saveFailed(ctx context.Context, job NotifyJob, reason string) {
if err := p.FailedStore.SaveFailed(ctx, job, reason); err != nil {
slog.Error("failed to save failed job",
"job_id", job.ID,
"err", err,
)
}
}
可观测性:重试系统的眼睛
重试系统至少需要追踪这些指标:
type JobMetrics interface {
IncSuccess()
IncFailure()
IncRetry()
SetQueueLength(n int)
}
接口只是示意,具体实现可以接 Prometheus、Datadog 或 expvar。关键指标清单:
| 指标 | 类型 | 说明 |
|---|---|---|
| job_success_total | Counter | 成功的任务数 |
| job_failure_total | Counter | 失败的任务数(按永久/重试耗尽分类) |
| job_retry_total | Counter | 重试次数 |
| job_queue_length | Gauge | 当前队列长度 |
| job_processing_duration | Histogram | 任务处理耗时 |
没有这些指标,系统可能已经在大量重试,你却只看到下游服务变慢。重试越隐蔽,越容易把真实故障拖到更晚才暴露。
日志策略:
- 每次重试记录 job id、attempt 和错误
- 永久失败用 Warn 级别
- 耗尽重试用 Error 级别
- 不要把完整消息内容、token、用户隐私塞进日志
slog.Warn("notify job failed",
"job_id", job.ID,
"attempt", job.Attempt,
"err", err,
)
测试重试逻辑
重试逻辑必须通过单元测试验证。关键是注入可控制的退避延迟:
package main
import (
"context"
"errors"
"testing"
"time"
)
func TestRunWithRetryEventuallySucceeds(t *testing.T) {
var calls int
err := RunWithRetry(context.Background(), NotifyJob{ID: "j1"}, func(ctx context.Context, job NotifyJob) error {
calls++
if calls < 2 {
return errors.New("temporary")
}
return nil
})
if err != nil {
t.Fatal(err)
}
if calls != 2 {
t.Fatalf("expected 2 calls, got %d", calls)
}
}
func TestRunWithRetryPermanentError(t *testing.T) {
err := RunWithRetry(context.Background(), NotifyJob{ID: "j2"}, func(ctx context.Context, job NotifyJob) error {
return permanent(errors.New("bad request"))
})
if err == nil {
t.Fatal("expected error")
}
if !IsPermanent(err) {
t.Fatal("expected permanent error")
}
}
func TestRunWithRetryContextCancel(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
done := make(chan struct{})
go func() {
RunWithRetry(ctx, NotifyJob{ID: "j3"}, func(ctx context.Context, job NotifyJob) error {
return errors.New("temporary")
})
close(done)
}()
cancel()
select {
case <-done:
// 正确:context 取消后重试循环应该退出
case <-time.After(2 * time.Second):
t.Fatal("retry did not respect context cancellation")
}
}
关键测试用例:
- 临时错误在重试后成功
- 永久错误立即失败不重试
- context 取消时正确退出
- 最大重试次数到达后返回错误
分布式重试的额外挑战
在分布式系统中,重试设计还要考虑更多因素:
集群中的重复重试:如果多个 worker 消费同一条消息,可能同时对失败任务重试。解决方案是用锁(分布式锁或数据库行锁)确保同一时间只有一个 worker 处理。
消息可见性超时:队列系统(SQS、RabbitMQ)通常有 visibility timeout。如果任务处理时间超过超时时间,消息会被重新投递。重试逻辑需要处理重复投递。
重试队列隔离:不要让重试任务和普通任务在同一个队列中竞争资源。可以设置专门的延迟队列,让重试任务在低峰期执行。
跨服务调用的链路重试:如果服务 A 调用服务 B,服务 B 内部又重试,可能造成级联放大。建议只在调用链的最外层做重试,内层服务返回明确的错误码。
性能对比:不同重试策略的影响
| 策略 | 平均完成时间 | 下游压力 | 适用场景 |
|---|---|---|---|
| 不重试 | 最快 | 最低 | 可容忍偶尔失败 |
| 固定间隔 1s x 3 | 3s + 处理时间 | 集中爆发 | 低并发场景 |
| 线性退避 | 逐渐增长 | 中等 | 简单场景 |
| 指数退避 + jitter | 合理 | 均匀分散 | 高并发、生产环境 |
没有绝对最好的策略,只有最适合当前场景的权衡。
FAQ:重试常见问题
Q1: 退避延迟的最大值应该设多少?
取决于业务容忍度。用户通知可以短(几秒到几十秒),批处理任务可以长(几分钟)。一般不建议超过 5 分钟,否则延迟队列可能更合适。
Q2: 重试次数用完怎么办?
必须归档到失败记录表或死信队列,不能默默丢弃。失败任务要有告警,让运维人员能及时介入。
Q3: 怎么防止重试把下游打垮?
除了 jitter 打散,还应该使用熔断器(Circuit Breaker)。当下游连续失败超过阈值时,暂停调用一段时间。
Q4: context 超时要比重试总时间长吗?
是的。如果 context 只有 5 秒,但重试总延迟有 30 秒,实际上重试循环会在第一次 context 取消时退出,达不到预期的重试次数。建议 context 超时至少等于重试总延迟的 1.5 倍。
Q5: 随机数种子会影响测试吗?
是的。测试中使用固定种子或接口注入伪随机数源,保证测试可重复。
最佳实践总结
- 错误分类是重试的前提:永久错误立即失败,临时错误才重试。
- 最大重试次数必须有限:无限重试会把小问题变成大问题。
- 退避不能固定:指数退避加 jitter 是标准做法。
- 幂等性不能省略:没有幂等性的重试是赌博。
- 失败必须归档:死信队列是重试系统的最后防线。
- 监听 context:让重试可取消、可超时。
- 日志保留关键信息:job id、attempt、错误类型,不要泄露隐私。
- 可观测性跟上:成功、失败、重试次数都要有指标。
- 测试重试逻辑:用注入的延迟件测试各种边界情况。
- 分布式场景格外谨慎:考虑重复投递、锁和队列隔离。
Go 中实现重试本身不难,难的是业务语义的正确性。重复执行是否安全?失败日志是否可排查?队列会不会被坏任务堵住?这些问题想清楚,重试才不会添乱。
性能对比与基准测试
理解性能问题的最佳方式是通过基准测试观察实际行为。运行 go test -bench=. -benchmem 可以得到每个操作的耗时和内存分配数据。对比不同实现时,建议固定输入规模,跑多次取平均值。
常见错误与最佳实践
错误一:性能优化过早
很多初学者刚写好代码就开始担心性能,结果引入了不必要的复杂度。正确的做法是先用清晰的写法实现功能,在性能问题真实出现时再通过 profile 定位热点。
错误二:忽略边界条件
空输入、超大输入、并发场景、系统资源耗尽等边界条件往往是 bug 的来源。写代码时养成习惯:每个函数都问自己,空值怎么办?错误怎么处理?
错误三:错误处理不完整
Go 的错误处理要求显式检查。常见问题是只在最外层处理错误,中间层把 error 吞掉。使用 fmt.Errorf 配合 %w 保留原始错误链。
错误四:并发代码缺少同步
Go 的并发模型很简洁,但共享内存访问必须同步。用 go test -race 验证并发安全性。
生产环境注意事项
- 日志要克制:不要记录敏感信息,不要在热路径上打印大量日志。
- 超时和取消:所有外部调用都要有超时。
- 资源限制:限制请求体大小、并发连接数、内存使用。
- 优雅关闭:http.Server 要设置 Shutdown 超时,goroutine 要有退出机制。
- 可观测性:至少记录关键指标。
测试策略
好的测试应该覆盖正常路径、错误路径和边界条件。表驱动测试是推荐的方式。每次修改代码后都要跑一遍测试,CI 中集成 go test ./... 是最基本的自动化保障。
实战 FAQ
Q: 这个功能在旧版 Go 中能用吗?
A: 需要看具体功能引入的版本。建议使用最新的稳定版 Go。
Q: 第三方库更好还是标准库更好?
A: 能标准库解决先用标准库,第三方库引入依赖成本和许可证风险。
Q: 怎么判断代码算不算过度设计?
A: 问自己:这个抽象让调用方更简单了吗?减少了多少重复?维护成本是增加还是减少了?
小结
掌握这项技能的关键不是记住所有 API,而是理解背后的设计原则和适用边界。先让代码工作,再让它正确,最后才考虑让它更快。清晰的代码比聪明的代码更有价值。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。