Skip to content

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 包含必须单独入账的增量效果,不能仅比较序列就丢弃。应补取事件或事实来源,并用独立业务键保护效果。设计前先确认消息语义是状态还是增量。

继续阅读:状态机、业务接入与投影。

理解原理 · 分析取舍 · 用真实代码检验