Ohhnews

分类导航

$ cd ..
foojay原文

幂等写入路径:在消息重复时保持集成管道正确性

#幂等性#集成管道#消息去重#数据一致性#分布式系统

一笔对不平的账

先讲一个真实案例。一次大促结束后对账时,财务发现一批订单状态不对:我们这边显示“已支付”,但财务系统并没有入账。

顺着链路追查后我们发现,支付网关在高峰期重试了一次,同一条支付通知被推送了两次。第一次被正常处理;第二次到达时,消费者刚好处于滚动重启过程中,消息被重新入队并再次被消费。

中间状态表被写了两遍。下游做了自己的去重所以没受影响,但我们这边的对账报表逻辑被多出来的一条脏数据搞崩了。

这件事本身不算严重,但它让我们清楚地认识到:在分布式环境下,消息就是会重复。网络重传、队列重投、消费者重启、上游超时重发——你根本阻止不了。

唯一能做的,就是让系统在收到一次消息和收到五次消息时,行为保持一致。这就是幂等性。

幂等键选什么,才是关键

很多人以为幂等就是“加个唯一键”。对,也不对。关键在于你究竟以什么为键。

一开始我们直接拿消息队列的消息 ID 作为幂等键。看起来合理:每条消息不都有 ID 吗?

但坑就在这里。如果上游手动重推一次,或者它自己的重试机制触发了,就会生成一个全新的消息 ID。在消息层面看,这像是两条“不同的消息”,但在业务层面它们是同一件事,所以去重直接失效。我们在这个坑里栽了两次,才彻底修好。

后来我们改成要求所有上游在事件体里带上业务级别的幂等键。幂等键怎么构造,不同场景不一样,但核心原则只有一个:必须能唯一标识“这个业务事件发生过一次”。举几个例子:

$ cat
// payment event -> payment id + status
{
  "idempotent_key": "PAY-20240315-0042:PAID",
  "event_type": "payment.completed",
  "payment_id": "PAY-20240315-0042",
  "status": "PAID",
  "amount": 1299.00
}

// inventory change -> SKU + warehouse + batch
{
  "idempotent_key": "SKU-8823:WH-SZ-01:BATCH-20240315-003",
  "event_type": "inventory.adjusted",
  "sku": "SKU-8823",
  "warehouse": "WH-SZ-01",
  "batch": "BATCH-20240315-003",
  "delta": -5
}

// generic entity sync -> entity type + id + version
{
  "idempotent_key": "customer:C-10042:v17",
  "event_type": "entity.updated",
  "entity_type": "customer",
  "entity_id": "C-10042",
  "version": 17
}

注意,幂等键往往是复合的。同一个订单会触发多个状态变更,所以不能只用订单 ID:“已支付”和“已发货”是两件不同的事,必须加上状态或版本号来区分。

“唯一标识一个实体”和“唯一标识一次事件发生过”是两码事,很多人会搞混。

推行这个标准也不是一帆风顺。有些上游团队会觉得:“凭什么要我改数据格式来配合你的管道?”遇到这种情况,我们只能请架构师出面对齐,解释为什么这件事必须在源头做。

如果你想在中间层打补丁,那是补不完的,因为你根本不知道上游的业务语义。

去重检查到底放在哪

一开始我们在业务代码里做去重:处理消息前先查一下这个键是否已经处理过。功能上没毛病,但存在竞态条件。

当消费者的处理逻辑比较复杂时(一个事件要写三张表、调用两个下游服务),从去重检查到真正写入之间有一段处理时间。如果重复消息恰好在这个窗口内到达,就会出现“检查时没有重复,写入时却有重复”。在高并发下,这种情况比你想象的更常见。

后来我们把去重下沉到数据库层,用唯一约束兜底。先建一张去重表:

[LOADING...]

处理逻辑大致如下:去重记录和业务数据在同一个事务里写入,要么都成功,要么都回滚:

$ java
@Transactional
public void handleEvent(IntegrationEvent event) {
    try {
        dedupRepository.insert(DedupRecord.builder()
            .idempotentKey(event.getIdempotentKey())
            .eventType(event.getEventType())
            .sourceSystem(event.getSource())
            .correlationId(event.getCorrelationId())
            .build());

        // no exception means it's not a duplicate; carry on
        businessService.process(event);

    } catch (DuplicateKeyException e) {
        // primary-key collision = duplicate event; skip it and move on
        log.info("Duplicate event skipped: key={}", event.getIdempotentKey());
    }
}

这样,去重和业务写入就是原子的,不再有竞态窗口。

顺便说一句,这张去重表会一直增长,所以需要定期清理。我们搞了一个定时任务,每天晚上清理 30 天前的记录。30 天这个数字是根据我们业务中实际的消息重复窗口定的:绝大多数重复投递都发生在几分钟内,30 天绰绰有余。

如果业务写入跨多个数据源(写完数据库还要调外部 API),单靠数据库事务就覆盖不了,需要补偿机制。这部分后面讲韧性的时候再展开。

重复消息来了:忽略、覆盖还是合并

怎么处理取决于业务语义,没有标准答案。实践中我们遇到过三种策略。

忽略用得最多:检测到重复就直接丢掉,什么都不做。它适合天然幂等的操作,比如“把用户状态设为已激活”,做一次和做十次效果一样。在我们系统里大约 70% 的事件类型走这条路。它实现最简单,也最不容易出 bug。

覆盖适合最终状态同步,比如从 CRM 同步客户的最新联系方式。但这里有个坑:如果事件乱序到达(先到的那条反而是新数据,后到的是旧数据),盲目覆盖会把新数据回滚掉。这类 bug 排查起来特别恶心,因为数据“看起来是正常的”,只是不是最新的,你可能很久都不会察觉。所以我们给覆盖加了一个版本号检查:

$ java
public void upsertWithVersionCheck(EntitySync sync) {
    int updated = jdbcTemplate.update(
        "UPDATE entity_store SET data = ?, version = ?, updated_at = NOW() " +
        "WHERE entity_id = ? AND entity_type = ? AND version < ?",
        sync.getData(), sync.getVersion(),
        sync.getEntityId(), sync.getEntityType(), sync.getVersion()
    );
    if (updated == 0) {
        // either a new row to INSERT, or a stale version to drop
        try {
            jdbcTemplate.update(
                "INSERT INTO entity_store (entity_id, entity_type, data, version) " +
                "VALUES (?, ?, ?, ?)",
                sync.getEntityId(), sync.getEntityType(),
                sync.getData(), sync.getVersion());
        } catch (DuplicateKeyException e) {
            // a newer version already landed; dropping this one is correct
            log.debug("Stale version discarded: entity={}, ver={}",
                sync.getEntityId(), sync.getVersion());
        }
    }
}

这本质上是一个简化的 Last-Write-Wins:只有更高版本才能覆盖。这个版本号必须由源系统生成。你不能在中间层自己造,因为中间层并不知道源系统在什么时候更新了什么。

合并是最复杂的:把旧数据和新数据合到一起。它适合增量数据(给订单追加备注、往购物车加商品这类场景)。它的实现复杂度最高,也最容易出 bug。就拿“追加备注”来说:如果同一条备注因为重复投递被追加了两次怎么办?那你就得在合并逻辑里再做一层去重。套娃,基本就是这样。所以能不用就不用。

跨系统把身份串起来

做集成的朋友大概都遇到过:同一个实体,不同系统有不同的 ID。一个客户在 CRM 里叫 CRM-00123,到 ERP 里变成 ERP-C456,到 WMS 里又是 WMS-CUST-789,我们平台上还有一套自己的内部 ID。

当事件从 CRM 流向 ERP 时,你得知道 CRM-00123 对应 ERP 里的哪条记录。我们在集成平台里维护了一张 ID 映射表,用平台内部的 internal_id 把各个系统的 external_id 串起来:

[LOADING...]

事件进入管道后做的第一件事,就是查映射表拿到目标系统的 ID。如果查不到(比如创建事件,下游还没有对应记录),就标记为“待创建”,等下游创建完成后,我们再把映射关系写回去。

这里有个细节:映射回写本身也可能失败或重复,所以对这表的写入也得是幂等的。我们用 INSERT ... ON CONFLICT DO NOTHING

另一个非常重要的东西是关联 ID(correlation ID)。每个业务流程(比如“创建一个订单并同步到所有系统”)都会有一个全局唯一 ID,贯穿该流程的每一个事件和每一行日志。

排查问题时,一个 grep 就能把散落在几十个系统、几百行日志里的信息串起来。用过的人都知道,出故障时这东西有多重要。没有它,你只能瞪着几十个系统各自独立的日志,完全不知道谁和谁有关联,定位一个问题从半小时变成半天。

怎么证明你没丢数据

每隔几天老板就会问同一个问题:“你怎么证明管道没有丢数据?”分布式系统确实没法从数学上证明,但你可以用工程手段把置信度推得很高。

我们做了端到端的事件追踪。每个事件从进入管道到离开,每经过一个处理节点都会留一条追踪日志:

$ java
public class EventTracer {

    private final KafkaTemplate<String, TraceRecord> traceProducer;

    public void trace(IntegrationEvent event, String node, TraceStatus status) {
        TraceRecord record = TraceRecord.builder()
            .eventId(event.getEventId())
            .idempotentKey(event.getIdempotentKey())
            .correlationId(event.getCorrelationId())
            .processingNode(node)
            .status(status)  // RECEIVED / PROCESSING / COMPLETED / FAILED
            .timestamp(Instant.now())
            .build();
        traceProducer.send("event-trace-log", event.getEventId(), record);
    }
}

这些追踪日志写入 Kafka,并长期归档到对象存储。在此基础上,我们跑两种对账。

实时对账每五分钟扫描一次,找出“已接收但未完成”的事件。如果某个事件在接收后超过十分钟还没有 COMPLETED 记录,就触发告警。这能很快暴露“卡住”的事件,比如消费者死锁,或者下游接口一直不返回。

离线对账每天晚上做一次全量比对。拿上游声称已发送的事件列表,和我们实际处理完成的事件列表做 diff,找出“上游发了但我们没处理完”以及“我们处理了但上游根本没发”。前者意味着丢数据;后者可能是重复或幽灵记录。两种对账互补:实时对账快但覆盖有限,离线对账慢但能兜住边界情况。

有了这套追踪系统,对老板的答复就不再是“不可能丢数据”,而是“如果丢了一条,我们五分钟内发现,十分钟内定位,半小时内恢复”。

关于作者

Yuelin Ou 是一名数据与 AI 工程师,工作重点是幂等写入路径、分布式管道韧性,以及在不破坏正确性保证的前提下扩展企业集成系统。她拥有罗切斯特大学数学学士学位,辅修计算机科学。个人网站:yuelinou.com。

本文《幂等写入路径:当消息重复时如何保证集成管道正确》首发于 foojay