业务系统性能改造手册 · 异步拆链

治理重试、积压与失败补偿

2026-08-054 min read异步拆链
摘要

将异步失败分成可重试、不可重试、业务终止和结果未知,使用有限退避、积压时间预算、可操作死信与带审计的补偿流程,让持续失败从无限循环变成明确状态。

04-02 已经保证订单与事件一起落盘,也接受生产确认和消费确认窗口中的重复。链路不会静默忘掉工作了,新的问题随之出现:下游持续超时、消息一直重投,待处理量越积越多。

最省事的配置通常是“失败就重试”。它对短暂抖动有效,对权限错误、消息结构不兼容和业务状态冲突却没有帮助。相反,同一批失败工作会持续抢线程、连接和下游配额,把原本健康的消息也拖慢。

小编更愿意把重试看成一笔有限预算。每次失败先分类,再决定等待、终止、对账还是转入人工处置;预算耗尽后必须离开热路径,不能靠无限循环假装系统仍在处理。

一、失败分类决定下一步动作

失败分类决定重试暂停死信或人工处理不能让所有错误进入同一种无限重试 异常类名不等于失败类型。TimeoutException 可能是下游短暂抖动,也可能是对方已经成功而响应丢失;HTTP 400 可能是永久无效输入,也可能来自错误的版本映射。分类器需要使用操作类型、协议结果、业务状态和目标系统契约共同判断。

失败类型识别依据允许的下一步不应采取的动作
瞬时技术失败已确认未生效,且依赖短暂不可用、连接中断或可恢复限速在预算内退避后重试立即无间隔循环
持续技术失败凭据、路由、Schema、配置或代码错误,重复执行不会自行恢复隔离事件并告警,修复后受控重放用更多次数掩盖配置错误
业务性失败订单已取消、版本过期、业务前置条件不成立记录明确终态,必要时进入业务处理当成基础设施抖动重试
结果未知04-02 的提交/远端调用可能已生效,但确认丢失按稳定 requestId/eventId 查询和对账直接重新执行有副作用操作

结果未知必须单列。它不是“也许下次会好”的瞬时失败,而是系统不知道副作用是否已经发生。先查状态;能确认未生效时才进入普通重试,确认已生效则推进状态,仍无法确认就保持 RECONCILING 并交给人工或专门核对流程。

一个可追踪的状态流可以写成:

text
1READY → RUNNING → SUCCEEDED 23 ├─ 瞬时失败且预算未耗尽 → RETRY_SCHEDULED ─→ READY 4 ├─ 结果未知 ───────────→ RECONCILING 5 ├─ 业务条件不成立 ─────→ TERMINAL_REJECTED 6 └─ 持续失败/预算耗尽 ───→ QUARANTINED

QUARANTINED 可以映射到死信队列、失败主题或数据库失败表。名称并不重要,关键是消息离开热重试路径后仍能按 eventId 查询,原始负载、失败原因和处理历史没有丢失。

二、有限退避要先定义 attempt 语义

先统一计数口径:本文中 attempt=1 表示首次真实执行,不是第一次“重试”;一次执行失败后,只有 failedAttempt < maxAttempts 才能安排下一次。MQ 的投递次数、应用调用次数、对下游的真实请求次数可能不同,落地时应选一个作为策略计数,并把另外两个作为观测值保留。

候选等待上界采用指数增长并封顶:

text
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 修改可见时间或投递延迟消息:

java
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}

maxAttemptsbaseDelaymaximumDelay 都是待实测参数。上限要同时受业务完成承诺、下游恢复时间、额外流量预算和消息保留时间约束。十次重试并不比三次更可靠;若错误长期存在,它只意味着每条原始工作最多制造更多下游调用。

失败分类也要版本化。新增一个错误码时,默认进入 UNKNOWN 或隔离通常比默认重试安全;否则一次协议升级就可能把永久失败变成重试风暴。

三、积压预算要用时间表达

积压要用预计清空时间衡量补偿则是一笔可审计可幂等的新业务动作 队列里有多少条消息,单独看没有判断力。同样的十万条工作,对每秒完成量不同、单条业务时效不同的消费者,含义完全不同。

固定窗口内至少记录:

yaml
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,在新工作仍持续到达且 μ > λ 时,理想清空时间可以粗估为:

text
1drain time ≈ B / (μ - λ)

这个式子要求工作近似同质、口径一致、完成速率稳定,且重试工作已经计入负载。分区热点、顺序键阻塞、延迟消息、长尾任务或下游限速都会破坏这些前提。实际告警应优先看 oldestUnfinishedEventAge 是否逼近业务完成时限,再用增长速率和清空估算解释风险。

重试量要单列。只看 Broker 总消费速率,可能出现“消费者很忙、业务完成很少”:大量执行都在处理同一批失败消息。为重试预留多少执行份额应由压测决定,同时确保新到的健康工作仍有完成能力;不要让历史失败独占所有 Worker。

四、死信是待办区,不是垃圾桶

进入死信或隔离区的消息至少保留:原始 eventId、事件类型与版本、业务键、首次和最后失败时间、总执行次数、规范化错误码、错误摘要、原始负载引用、来源位置、当前负责人和处置状态。敏感字段不能因为进入运维通道就失去访问控制。

不同 MQ 对延迟、投递次数和死信的定义并不相同。以 RabbitMQ 为例,拒绝且不重新入队、消息过期、超过队列长度或特定队列的投递上限都可能触发 dead-letter;TTL 与 dead-letter 的具体行为还受队列类型和版本影响。RabbitMQ Dead Letter ExchangesRabbitMQ TTL 本系列尚未确定 MQ,实施时必须把下面的抽象状态映射到目标产品并做故障测试,不能照抄 RabbitMQ 参数。

受控重放要经过检查,而不是点一下“全部重试”:

  1. 冻结目标范围,记录事件 ID 列表、筛选条件和操作者。
  2. 确认根因已经修复,并用代表性事件做只读校验。
  3. 检查事件 Schema、权限上下文、业务当前状态和幂等记录仍兼容。
  4. 选择小批量重放,保留原 eventId,记录新的操作 ID。
  5. 核对业务终态与重复副作用,再逐步扩大范围。

重放工具不能修改原始负载来“顺便修数据”。需要修改业务事实时,应产生一条经过审批的新操作,原事件和失败证据保持不变。

五、补偿是一笔新的业务动作

事务回滚发生在原事务尚未提交时;补偿面对的是已经提交并可能被外部观察到的事实。订单已经标成已受理、库存也可能变化,再把数据库字段改回旧值,并不能自动撤销通知、外部调用和用户决策。

补偿入口至少需要以下契约:

yaml
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 查询或对账,不能直接重放。人工入口也不能绕开权限校验、审批和审计;“手工执行”只改变触发方式,不降低业务约束。

一个面向值班人员的最小操作手册可以这样写:

yaml
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 / (μ - λ) 的条件化估算对照。

第二类让某一类事件持续返回确定的配置错误或业务错误,检查它是否按分类策略有限执行、离开热路径并进入可查询状态。与此同时,其他健康事件的完成延迟不应因为这批失败工作无限恶化。

yaml
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 工作也不能因为配置回退被降级成普通重试。

到这里,异步链路才算把失败变成了可以操作的状态:短暂故障有限等待,永久故障离开热路径,积压用完成时限衡量,未知结果先核对,已生效事实通过新业务动作补偿。至于下游本身变慢时如何用超时和隔离限制故障半径,留到下一章。