切换主题
Inbox、幂等与乱序
问题:超时后的重试会不会创建两次
客户端提交导出请求后等待超时。服务可能已保存任务,只是响应丢失。再次调用如果生成新任务,会发生重复扣费。另一方面,同一事件第二次送达不是一个新业务动作;这两类重复需要不同标识。
请求幂等键识别同一次操作,事件 ID 识别同一条消息,业务唯一键识别某种业务效果。三个键不能随意互换。
原理:重复是允许的,效果要受约束
请求幂等记录一般由调用方、操作类型、幂等键和请求指纹组成。同键同内容返回原结果;同键不同内容应拒绝,不能默默返回第一次结果让用户误以为新参数生效。
Inbox 保存已经接收的事件。去重记录与业务修改应在同一事务:若提前写“已处理”然后崩溃,效果丢失;若先改业务后写去重,效果重复。
最小例子:原子去重
sql
-- 教学 SQL:应用检查 INSERT 的返回值。
BEGIN;
INSERT INTO inbox (event_id, received_at)
VALUES ($1, now())
ON CONFLICT (event_id) DO NOTHING
RETURNING event_id;
-- 仅首次插入成功才执行以下业务修改。
UPDATE accounts SET credited = credited + $2 WHERE id = $3;
COMMIT;如果业务动作是外部调用,这个事务只能保护本地状态。通常先持久化新的出站意图,把远端调用交给下一个可恢复步骤。
方案比较:去重与顺序独立设计
| 情形 | 所需机制 | 不足以解决问题的做法 |
|---|---|---|
| 重复创建请求 | 调用者范围+操作+键+指纹 | 每次重试新建随机键 |
| 重复事件 | 稳定事件 ID+原子 Inbox | 仅客户端防连点 |
| 旧状态晚到 | 按对象序列比较或条件转换 | 收到就覆盖 |
| 序列有缺口 | 查询补偿或等待缺失事件 | 认为较新事件永远包含全部事实 |
全量状态快照可按版本吸收较旧更新;增量“加 10”事件缺失后不能直接跳过。序列通常属于同一对象或执行,不应把不同执行的数字混比。
真实案例:请求指纹与 Callback 投影
TDP Repository 契约包含 PrincipalID、Operation、Key 和 RequestHash。App Demo 的 Callback 入口校验 CloudEvent 版本、ID、序列和归属,再交给 Repository.AcceptCallback 去重与更新投影;缺口用于日志诊断和后续补偿。
已核对的实现 · 本地代码快照
核心定义协作契约 backend · c6dbe05c
server/internal/core/execution/ports/outbound/repository.go · 第 9–24 行
符号:Repository · 核对日期 2026-10-02
来源与提交版本一致
// Idempotency 描述一次写操作的调用方、键和请求指纹。
type Idempotency struct {
PrincipalID string // 隔离不同调用方的同名幂等键。
Operation string // 隔离同一调用方的不同业务动作。
Key string // 调用方提供的幂等键。
RequestHash []byte // 用于拒绝复用同一键但内容不同的请求。
}
// Repository 原子持久化 Task、Execution、首个事件和幂等记录。
type Repository interface {
// Create 创建聚合;重复且请求指纹一致时返回原结果。
Create(context.Context, *domain.Task, *domain.Execution, domain.Event, Idempotency) (*domain.Task, *domain.Execution, error)
// Get 在应用范围内读取任务。
Get(context.Context, Scope, domain.ID) (*domain.Task, error)
}
片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。
Callback 校验、去重与缺口 app-demo · c383860e
internal/core/task/application/callback.go · 第 26–109 行
符号:ReceiveCallback / VerifyCallbackSignature · 核对日期 2026-10-02
来源与提交版本一致
func (s *Service) ReceiveCallback(ctx context.Context, input taskdomain.CallbackInput, applicationID string, discard bool, now time.Time) error {
if input.SpecVersion != "1.0" || input.Sequence < 1 || input.EventID == "" || input.Type == "" {
return ErrInvalidCallback
}
if parsed, err := uuid.Parse(input.EventID); err != nil || parsed == uuid.Nil {
return ErrInvalidCallback
}
if input.ApplicationID != "" && input.ApplicationID != applicationID {
return ErrInvalidCallback
}
if parsed, err := uuid.Parse(input.ExecutionID); err != nil || parsed == uuid.Nil {
return ErrInvalidCallback
}
if input.TaskID != "" {
if parsed, err := uuid.Parse(input.TaskID); err != nil || parsed == uuid.Nil {
return ErrInvalidCallback
}
}
if input.ExternalRef == "" {
order, err := s.repository.GetOrderByExecution(ctx, input.ExecutionID)
if err != nil {
return ErrInvalidCallback
}
input.ExternalRef = order.ID
}
if parsed, err := uuid.Parse(input.ExternalRef); err != nil || parsed == uuid.Nil {
return ErrInvalidCallback
}
if input.Status != "" && !map[string]bool{"CREATED": true, "QUEUED": true, "RUNNING": true, "SUCCEEDED": true, "FAILED": true, "TIMED_OUT": true, "CANCELING": true, "CANCELED": true}[input.Status] {
return ErrInvalidCallback
}
if discard {
s.logger.Warn("demo discarded callback after validation", "event_id", input.EventID, "execution_id", input.ExecutionID, "sequence", input.Sequence)
return nil
}
progress := 0
if input.Progress != nil {
progress = *input.Progress
}
_, gap, err := s.repository.AcceptCallback(ctx, taskports.CallbackEvent{
ID: input.EventID,
TaskID: input.TaskID,
ExecutionID: input.ExecutionID,
Type: input.Type,
Sequence: input.Sequence,
Body: append([]byte(nil), input.RawBody...),
ReceivedAt: now,
}, input.ExternalRef, input.Status, int32(progress), input.Data, now)
if gap {
s.logger.Warn("callback sequence gap", "order_id", input.ExternalRef, "execution_id", input.ExecutionID, "sequence", input.Sequence)
}
return err
}
func VerifyCallbackSignature(header string, body, secret []byte, now time.Time) error {
timestampPart, signaturePart, ok := strings.Cut(header, ",")
if !ok {
return ErrInvalidCallbackSignature
}
rawTimestamp, okTimestamp := strings.CutPrefix(timestampPart, "t=")
rawSignature, okSignature := strings.CutPrefix(signaturePart, "v1=")
if !okTimestamp || !okSignature {
return ErrInvalidCallbackSignature
}
timestamp, err := strconv.ParseInt(rawTimestamp, 10, 64)
if err != nil {
return ErrInvalidCallbackSignature
}
if delta := now.Sub(time.Unix(timestamp, 0)); delta < -5*time.Minute || delta > 5*time.Minute {
return ErrExpiredCallback
}
received, err := hex.DecodeString(rawSignature)
if err != nil {
return ErrInvalidCallbackSignature
}
mac := hmac.New(sha256.New, secret)
_, _ = fmt.Fprintf(mac, "%d.", timestamp)
_, _ = mac.Write(body)
if !hmac.Equal(received, mac.Sum(nil)) {
return ErrInvalidCallbackSignature
}
return nil
}
片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。
Callback 收到某状态不等于业务已经消费了完整过程。业务应关联当前 Execution,防止旧执行晚到结果覆盖重跑后的状态。签名与时间窗确认来源,Inbox 保护重复,两者职责不同。
失败与边界:幂等也有保存期限
幂等记录清理后同一个键可能再次被视为新请求。接口契约要说明重放期限;对支付等长期业务唯一事实,还需要独立唯一业务键,不能完全依赖短期缓存。
消费者 ACK 也有窗口:效果提交后 ACK 丢失,会再次收到消息。只要去重与效果原子,重复可以返回成功 ACK;遇到无效载荷则按明确策略死信,避免无休止重试。
迁移练习与参考答案
练习:订单先收到 sequence=8, SHIPPED,后收到 sequence=7, PAID。可以忽略 7 吗?
参考答案:如果事件携带足够的状态快照,投影可拒绝倒退并记录缺口;但如果 7 包含必须单独入账的增量效果,不能仅比较序列就丢弃。应补取事件或事实来源,并用独立业务键保护效果。设计前先确认消息语义是状态还是增量。