Go 后台任务重试入门:失败后怎么重试才不添乱

本文详解后台任务重试的设计方法,涵盖最大次数、退避等待、错误分类、幂等性、死信队列和可观察性,附带最佳实践和常见陷阱。

后台任务失败后要不要重试?答案通常是"要,但不能乱重试"。发送通知、同步数据、生成报表、调用第三方支付接口等操作都可能遇到临时失败。如果完全不重试,一次网络抖动就会变成用户可见的故障;如果无限制重试,可能把一个小问题放大成队列堵塞或下游服务雪崩。

本文用发送通知任务作为贯穿示例,深入讲解一个完整的重试系统设计:错误分类、退避策略、重试循环、幂等性、死信队列、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 次200ms200ms - 300ms
第 2 次400ms400ms - 600ms
第 3 次800ms800ms - 1200ms
第 4 次1.6s1.6s - 2.4s
第 5 次3.2s3.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_totalCounter成功的任务数
job_failure_totalCounter失败的任务数(按永久/重试耗尽分类)
job_retry_totalCounter重试次数
job_queue_lengthGauge当前队列长度
job_processing_durationHistogram任务处理耗时

没有这些指标,系统可能已经在大量重试,你却只看到下游服务变慢。重试越隐蔽,越容易把真实故障拖到更晚才暴露。

日志策略:

  • 每次重试记录 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 33s + 处理时间集中爆发低并发场景
线性退避逐渐增长中等简单场景
指数退避 + jitter合理均匀分散高并发、生产环境

没有绝对最好的策略,只有最适合当前场景的权衡。

FAQ:重试常见问题

Q1: 退避延迟的最大值应该设多少?

取决于业务容忍度。用户通知可以短(几秒到几十秒),批处理任务可以长(几分钟)。一般不建议超过 5 分钟,否则延迟队列可能更合适。

Q2: 重试次数用完怎么办?

必须归档到失败记录表或死信队列,不能默默丢弃。失败任务要有告警,让运维人员能及时介入。

Q3: 怎么防止重试把下游打垮?

除了 jitter 打散,还应该使用熔断器(Circuit Breaker)。当下游连续失败超过阈值时,暂停调用一段时间。

Q4: context 超时要比重试总时间长吗?

是的。如果 context 只有 5 秒,但重试总延迟有 30 秒,实际上重试循环会在第一次 context 取消时退出,达不到预期的重试次数。建议 context 超时至少等于重试总延迟的 1.5 倍。

Q5: 随机数种子会影响测试吗?

是的。测试中使用固定种子或接口注入伪随机数源,保证测试可重复。

最佳实践总结

  1. 错误分类是重试的前提:永久错误立即失败,临时错误才重试。
  2. 最大重试次数必须有限:无限重试会把小问题变成大问题。
  3. 退避不能固定:指数退避加 jitter 是标准做法。
  4. 幂等性不能省略:没有幂等性的重试是赌博。
  5. 失败必须归档:死信队列是重试系统的最后防线。
  6. 监听 context:让重试可取消、可超时。
  7. 日志保留关键信息:job id、attempt、错误类型,不要泄露隐私。
  8. 可观测性跟上:成功、失败、重试次数都要有指标。
  9. 测试重试逻辑:用注入的延迟件测试各种边界情况。
  10. 分布式场景格外谨慎:考虑重复投递、锁和队列隔离。

Go 中实现重试本身不难,难的是业务语义的正确性。重复执行是否安全?失败日志是否可排查?队列会不会被坏任务堵住?这些问题想清楚,重试才不会添乱。

性能对比与基准测试

理解性能问题的最佳方式是通过基准测试观察实际行为。运行 go test -bench=. -benchmem 可以得到每个操作的耗时和内存分配数据。对比不同实现时,建议固定输入规模,跑多次取平均值。

常见错误与最佳实践

错误一:性能优化过早
很多初学者刚写好代码就开始担心性能,结果引入了不必要的复杂度。正确的做法是先用清晰的写法实现功能,在性能问题真实出现时再通过 profile 定位热点。

错误二:忽略边界条件
空输入、超大输入、并发场景、系统资源耗尽等边界条件往往是 bug 的来源。写代码时养成习惯:每个函数都问自己,空值怎么办?错误怎么处理?

错误三:错误处理不完整
Go 的错误处理要求显式检查。常见问题是只在最外层处理错误,中间层把 error 吞掉。使用 fmt.Errorf 配合 %w 保留原始错误链。

错误四:并发代码缺少同步
Go 的并发模型很简洁,但共享内存访问必须同步。用 go test -race 验证并发安全性。

生产环境注意事项

  1. 日志要克制:不要记录敏感信息,不要在热路径上打印大量日志。
  2. 超时和取消:所有外部调用都要有超时。
  3. 资源限制:限制请求体大小、并发连接数、内存使用。
  4. 优雅关闭:http.Server 要设置 Shutdown 超时,goroutine 要有退出机制。
  5. 可观测性:至少记录关键指标。

测试策略

好的测试应该覆盖正常路径、错误路径和边界条件。表驱动测试是推荐的方式。每次修改代码后都要跑一遍测试,CI 中集成 go test ./... 是最基本的自动化保障。

实战 FAQ

Q: 这个功能在旧版 Go 中能用吗?
A: 需要看具体功能引入的版本。建议使用最新的稳定版 Go。

Q: 第三方库更好还是标准库更好?
A: 能标准库解决先用标准库,第三方库引入依赖成本和许可证风险。

Q: 怎么判断代码算不算过度设计?
A: 问自己:这个抽象让调用方更简单了吗?减少了多少重复?维护成本是增加还是减少了?

小结

掌握这项技能的关键不是记住所有 API,而是理解背后的设计原则和适用边界。先让代码工作,再让它正确,最后才考虑让它更快。清晰的代码比聪明的代码更有价值。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「golang」更多文章

  1. 熔断、降级与限流:Go 微服务韧性设计完全指南
  2. 事件溯源与 CQRS 在 Go 中的实践:复杂业务系统的架构升级
  3. TinyGo 嵌入式开发与物联网实战:微控制器编程完全指南