业务系统性能改造手册 · 异步拆链
治理重试、积压与失败补偿
将异步失败分成可重试、不可重试、业务终止和结果未知,使用有限退避、积压时间预算、可操作死信与带审计的补偿流程,让持续失败从无限循环变成明确状态。
04-02 已经保证订单与事件一起落盘,也接受生产确认和消费确认窗口中的重复。链路不会静默忘掉工作了,新的问题随之出现:下游持续超时、消息一直重投,待处理量越积越多。
最省事的配置通常是“失败就重试”。它对短暂抖动有效,对权限错误、消息结构不兼容和业务状态冲突却没有帮助。相反,同一批失败工作会持续抢线程、连接和下游配额,把原本健康的消息也拖慢。
小编更愿意把重试看成一笔有限预算。每次失败先分类,再决定等待、终止、对账还是转入人工处置;预算耗尽后必须离开热路径,不能靠无限循环假装系统仍在处理。
一、失败分类决定下一步动作
异常类名不等于失败类型。TimeoutException 可能是下游短暂抖动,也可能是对方已经成功而响应丢失;HTTP 400 可能是永久无效输入,也可能来自错误的版本映射。分类器需要使用操作类型、协议结果、业务状态和目标系统契约共同判断。
| 失败类型 | 识别依据 | 允许的下一步 | 不应采取的动作 |
|---|---|---|---|
| 瞬时技术失败 | 已确认未生效,且依赖短暂不可用、连接中断或可恢复限速 | 在预算内退避后重试 | 立即无间隔循环 |
| 持续技术失败 | 凭据、路由、Schema、配置或代码错误,重复执行不会自行恢复 | 隔离事件并告警,修复后受控重放 | 用更多次数掩盖配置错误 |
| 业务性失败 | 订单已取消、版本过期、业务前置条件不成立 | 记录明确终态,必要时进入业务处理 | 当成基础设施抖动重试 |
| 结果未知 | 04-02 的提交/远端调用可能已生效,但确认丢失 | 按稳定 requestId/eventId 查询和对账 | 直接重新执行有副作用操作 |
结果未知必须单列。它不是“也许下次会好”的瞬时失败,而是系统不知道副作用是否已经发生。先查状态;能确认未生效时才进入普通重试,确认已生效则推进状态,仍无法确认就保持 RECONCILING 并交给人工或专门核对流程。
一个可追踪的状态流可以写成:
1READY → RUNNING → SUCCEEDED
2 │
3 ├─ 瞬时失败且预算未耗尽 → RETRY_SCHEDULED ─→ READY
4 ├─ 结果未知 ───────────→ RECONCILING
5 ├─ 业务条件不成立 ─────→ TERMINAL_REJECTED
6 └─ 持续失败/预算耗尽 ───→ QUARANTINEDQUARANTINED 可以映射到死信队列、失败主题或数据库失败表。名称并不重要,关键是消息离开热重试路径后仍能按 eventId 查询,原始负载、失败原因和处理历史没有丢失。
二、有限退避要先定义 attempt 语义
先统一计数口径:本文中 attempt=1 表示首次真实执行,不是第一次“重试”;一次执行失败后,只有 failedAttempt < maxAttempts 才能安排下一次。MQ 的投递次数、应用调用次数、对下游的真实请求次数可能不同,落地时应选一个作为策略计数,并把另外两个作为观测值保留。
候选等待上界采用指数增长并封顶:
1upper(attempt) = min(cap, base × 2^(attempt - 1))
2delay = random(0, upper(attempt))这里使用 full jitter,让同一时刻失败的工作分散到一个时间区间。AWS Builders' Library 也指出,没有随机扰动时,相关失败会在相同等待后再次同时冲击下游;重试还会放大已经过载的依赖。Timeouts, retries and backoff with jitter
下面的 Java 17 示例只计算决策,不负责向某个 MQ 修改可见时间或投递延迟消息:
1import java.time.Duration;
2import java.time.Instant;
3import java.util.Objects;
4import java.util.random.RandomGenerator;
5
6enum FailureKind {
7 TRANSIENT,
8 PERSISTENT,
9 BUSINESS,
10 UNKNOWN
11}
12
13record RetryPolicy(Duration baseDelay, Duration maximumDelay, int maxAttempts) {
14 RetryPolicy {
15 Objects.requireNonNull(baseDelay);
16 Objects.requireNonNull(maximumDelay);
17 long baseMillis = baseDelay.toMillis();
18 long maximumMillis = maximumDelay.toMillis();
19 if (baseMillis <= 0
20 || maximumMillis < baseMillis
21 || maximumMillis == Long.MAX_VALUE
22 || maxAttempts <= 0) {
23 throw new IllegalArgumentException("invalid retry policy");
24 }
25 }
26}
27
28sealed interface FailureDecision
29 permits RetryAt, Quarantine, RejectBusiness, Reconcile {}
30
31record RetryAt(int nextAttempt, Instant notBefore) implements FailureDecision {}
32
33record Quarantine(String reason) implements FailureDecision {}
34
35record RejectBusiness(String reason) implements FailureDecision {}
36
37record Reconcile(String reason) implements FailureDecision {}
38
39final class RetryPlanner {
40 private final RetryPolicy policy;
41 private final RandomGenerator random;
42
43 RetryPlanner(RetryPolicy policy, RandomGenerator random) {
44 this.policy = Objects.requireNonNull(policy);
45 this.random = Objects.requireNonNull(random);
46 }
47
48 FailureDecision afterFailure(
49 FailureKind kind,
50 int failedAttempt,
51 Instant now,
52 String reason) {
53 if (failedAttempt <= 0) {
54 throw new IllegalArgumentException("failedAttempt starts at one");
55 }
56 return switch (kind) {
57 case UNKNOWN -> new Reconcile(reason);
58 case BUSINESS -> new RejectBusiness(reason);
59 case PERSISTENT -> new Quarantine(reason);
60 case TRANSIENT -> retryOrQuarantine(failedAttempt, now, reason);
61 };
62 }
63
64 private FailureDecision retryOrQuarantine(
65 int failedAttempt, Instant now, String reason) {
66 if (failedAttempt >= policy.maxAttempts()) {
67 return new Quarantine("attempt budget exhausted: " + reason);
68 }
69 long upperMillis = upperDelayMillis(failedAttempt);
70 long delayMillis = random.nextLong(upperMillis + 1);
71 return new RetryAt(failedAttempt + 1, now.plusMillis(delayMillis));
72 }
73
74 long upperDelayMillis(int failedAttempt) {
75 long value = policy.baseDelay().toMillis();
76 long cap = policy.maximumDelay().toMillis();
77 for (int i = 1; i < failedAttempt && value < cap; i++) {
78 value = value > cap / 2 ? cap : value * 2;
79 }
80 return Math.min(value, cap);
81 }
82}maxAttempts、baseDelay 和 maximumDelay 都是待实测参数。上限要同时受业务完成承诺、下游恢复时间、额外流量预算和消息保留时间约束。十次重试并不比三次更可靠;若错误长期存在,它只意味着每条原始工作最多制造更多下游调用。
失败分类也要版本化。新增一个错误码时,默认进入 UNKNOWN 或隔离通常比默认重试安全;否则一次协议升级就可能把永久失败变成重试风暴。
三、积压预算要用时间表达
队列里有多少条消息,单独看没有判断力。同样的十万条工作,对每秒完成量不同、单条业务时效不同的消费者,含义完全不同。
固定窗口内至少记录:
1backlogBudget:
2 workload: ${EVENT_TYPE_CONSUMER_AND_PRIORITY}
3 businessCompletionDeadline: ${DURATION_FROM_EVENT_CREATION}
4 window: ${START_END_AND_LENGTH}
5 rates:
6 producedPrimaryPerSecond: ${RATE}
7 completedPrimaryPerSecond: ${RATE}
8 retryExecutionsPerSecond: ${RATE}
9 quarantinedPerSecond: ${RATE}
10 backlog:
11 readyCount: ${COUNT}
12 delayedRetryCount: ${COUNT}
13 runningCount: ${COUNT}
14 reconcilingCount: ${COUNT}
15 quarantinedCount: ${COUNT}
16 oldestReadyAge: ${DURATION}
17 oldestUnfinishedEventAge: ${DURATION}
18 capacity:
19 measuredSustainableCompletionRate: ${RATE}
20 reservedRetryShare: ${RATE_OR_PERCENT_WITH_REASON}
21 alertAndAction:
22 warningCondition: ${AGE_GROWTH_AND_DURATION}
23 criticalCondition: ${DEADLINE_OR_DRAINABILITY_CONDITION}
24 ownerAndRunbook: ${REFERENCE}当主工作到达率为 λ、完成率为 μ,且两者在窗口内稳定时,λ > μ 意味着积压会以约 λ - μ 的速度增长。现有积压量为 B,在新工作仍持续到达且 μ > λ 时,理想清空时间可以粗估为:
1drain time ≈ B / (μ - λ)这个式子要求工作近似同质、口径一致、完成速率稳定,且重试工作已经计入负载。分区热点、顺序键阻塞、延迟消息、长尾任务或下游限速都会破坏这些前提。实际告警应优先看 oldestUnfinishedEventAge 是否逼近业务完成时限,再用增长速率和清空估算解释风险。
重试量要单列。只看 Broker 总消费速率,可能出现“消费者很忙、业务完成很少”:大量执行都在处理同一批失败消息。为重试预留多少执行份额应由压测决定,同时确保新到的健康工作仍有完成能力;不要让历史失败独占所有 Worker。
四、死信是待办区,不是垃圾桶
进入死信或隔离区的消息至少保留:原始 eventId、事件类型与版本、业务键、首次和最后失败时间、总执行次数、规范化错误码、错误摘要、原始负载引用、来源位置、当前负责人和处置状态。敏感字段不能因为进入运维通道就失去访问控制。
不同 MQ 对延迟、投递次数和死信的定义并不相同。以 RabbitMQ 为例,拒绝且不重新入队、消息过期、超过队列长度或特定队列的投递上限都可能触发 dead-letter;TTL 与 dead-letter 的具体行为还受队列类型和版本影响。RabbitMQ Dead Letter Exchanges、RabbitMQ TTL 本系列尚未确定 MQ,实施时必须把下面的抽象状态映射到目标产品并做故障测试,不能照抄 RabbitMQ 参数。
受控重放要经过检查,而不是点一下“全部重试”:
- 冻结目标范围,记录事件 ID 列表、筛选条件和操作者。
- 确认根因已经修复,并用代表性事件做只读校验。
- 检查事件 Schema、权限上下文、业务当前状态和幂等记录仍兼容。
- 选择小批量重放,保留原
eventId,记录新的操作 ID。 - 核对业务终态与重复副作用,再逐步扩大范围。
重放工具不能修改原始负载来“顺便修数据”。需要修改业务事实时,应产生一条经过审批的新操作,原事件和失败证据保持不变。
五、补偿是一笔新的业务动作
事务回滚发生在原事务尚未提交时;补偿面对的是已经提交并可能被外部观察到的事实。订单已经标成已受理、库存也可能变化,再把数据库字段改回旧值,并不能自动撤销通知、外部调用和用户决策。
补偿入口至少需要以下契约:
1compensationRequest:
2 compensationId: ${STABLE_UNIQUE_ID}
3 sourceEventId: ${EVENT_ID}
4 businessKey: ${ORDER_ID_OR_OTHER}
5 reasonCode: ${CONTROLLED_VALUE}
6 requestedBy: ${IDENTITY}
7 approvedBy: ${IDENTITY_OR_POLICY}
8 expectedCurrentStateAndVersion: ${PRECONDITION}
9 proposedAction: ${CANCEL_RELEASE_CORRECT_OR_OTHER}
10 affectedResources: ${DATABASE_AND_EXTERNAL_EFFECTS}
11 dryRunEvidence: ${REFERENCE}
12 createdAt: ${TIMESTAMP}执行时先读取当前状态与版本。前置条件不再成立,就停止并重新评估,不能拿事故发生时的快照覆盖后来已经合法推进的订单。compensationId 必须有唯一约束;同一补偿请求重复提交时返回原结果,不产生第二份副作用。
补偿记录应保存每个子动作的开始、结果和未知状态。外部调用结果未知时沿用 04-02 的原则:用稳定操作 ID 查询或对账,不能直接重放。人工入口也不能绕开权限校验、审批和审计;“手工执行”只改变触发方式,不降低业务约束。
一个面向值班人员的最小操作手册可以这样写:
1failureRunbook:
2 trigger: ${ALERT_OR_INCIDENT_ID}
3 scopeQuery: ${EVENT_TYPE_TIME_RANGE_TENANT_AND_IDS}
4 immediateContainment: ${PAUSE_CLASS_REDUCE_CONCURRENCY_OR_NONE}
5 diagnosis:
6 failureClass: ${TRANSIENT_PERSISTENT_BUSINESS_OR_UNKNOWN}
7 evidence: ${LOG_METRIC_DATABASE_AND_DOWNSTREAM_REFERENCES}
8 rootCauseFixed: ${PASS_FAIL_WITH_PROOF}
9 decision: ${RESUME_REPLAY_COMPENSATE_REJECT_OR_KEEP_RECONCILING}
10 execution:
11 operator: ${IDENTITY}
12 approval: ${REFERENCE}
13 batchAndRate: ${VALUES_FROM_RECOVERY_TEST}
14 operationIds: ${REFERENCES}
15 verification:
16 terminalBusinessState: ${QUERY_AND_RESULT}
17 duplicateSideEffects: ${QUERY_AND_RESULT}
18 backlogAndOldestAgeRecovered: ${EVIDENCE}
19 auditClosedAt: ${TIMESTAMP_OR_OPEN}六、恢复实验要制造慢消费和持续失败
正常压测之外,至少执行两类实验。第一类保持生产速率不变,受控降低消费者完成能力,观察积压增长、最老事件年龄和业务完成时限;恢复消费者后,测量实际清空时间,并与 B / (μ - λ) 的条件化估算对照。
第二类让某一类事件持续返回确定的配置错误或业务错误,检查它是否按分类策略有限执行、离开热路径并进入可查询状态。与此同时,其他健康事件的完成延迟不应因为这批失败工作无限恶化。
1asyncRecoveryExperiment:
2 experimentId: ${EXPERIMENT_ID}
3 fixedConditions:
4 applicationDatabaseBrokerClientVersions: ${VERSIONS}
5 eventMixAndKeyDistribution: ${REFERENCE}
6 primaryProductionRate: ${RATE}
7 normalConsumerCapacity: ${RATE_AND_WORKERS}
8 retryPolicyVersion: ${VERSION_OR_HASH}
9 slowConsumerPhase:
10 inducedCapacity: ${RATE_AND_METHOD}
11 duration: ${DURATION}
12 backlogGrowthPerSecond: ${RATE}
13 oldestAgeAtEnd: ${DURATION}
14 businessDeadlineBreaches: ${COUNT}
15 recoveryPhase:
16 restoredCapacity: ${RATE_AND_WORKERS}
17 predictedDrainTimeAndAssumptions: ${VALUE}
18 observedDrainTime: ${DURATION_OR_NOT_DRAINED}
19 downstreamAndDatabaseSaturation: ${EVIDENCE}
20 persistentFailurePhase:
21 injectedFailureClass: ${CLASS_AND_CODE}
22 affectedEvents: ${COUNT_AND_IDS}
23 executionsPerEvent: ${DISTRIBUTION}
24 quarantineOrTerminalCounts: ${COUNTS}
25 healthyEventP95P99: ${DURATIONS}
26 duplicateBusinessEffects: ${COUNT}
27 replayOrCompensation:
28 scopeAndApproval: ${REFERENCES}
29 terminalResults: ${COUNTS}
30 unresolvedUnknowns: ${COUNT_AND_IDS}
31 decision:
32 selectedParameters: ${NO_REAL_VALUES_UNTIL_MEASURED}
33 rollbackTrigger: ${CONDITION}测试记录要区分“消息被执行了多少次”和“业务完成了多少次”。重试策略的验收也不是死信为零,而是持续失败不会无限占用容量,所有未完成状态有负责人,恢复过程不产生重复业务结果。
上线时先启用状态记录和观测,再启用自动重试与死信路由;否则参数一旦不合适,现场只剩一串重复异常。回退策略时保留失败记录、原 eventId 和补偿审计,不清空去重表。正在 RECONCILING 的 04-02 UNKNOWN 工作也不能因为配置回退被降级成普通重试。
到这里,异步链路才算把失败变成了可以操作的状态:短暂故障有限等待,永久故障离开热路径,积压用完成时限衡量,未知结果先核对,已生效事实通过新业务动作补偿。至于下游本身变慢时如何用超时和隔离限制故障半径,留到下一章。