找回密码
立即注册
搜索
热搜: Java Python Linux Go
发回帖 发新帖
Claude、GPT 海外模型 API 接入Claude skills 从入门到精通 吴恩达亲授 AI Agent 核心技能2026 瞪哥公务员考试全攻略 行测申论一站式系统备考
Agent 文心智能蒸馏模型实战 90G 课程智泊 AI 大模型训练营 基于 LangChain 的 RAG 与提示工程实战构建企业级 AI 大脑:大模型微调与 RAG / Agent 全栈实战

6028

积分

0

好友

735

主题
发表于 2 小时前 | 查看: 3| 回复: 0

假如一个已发货的订单突然变回“待付款”,我会先翻消费代码。

创建、付款、发货,三条消息按顺序到达。消费者收到就扔进线程池,创建事件碰上慢查询,等它落库,付款和发货早处理完了。旧状态直接覆盖新状态。

消息按顺序到了,业务未必按顺序生效。 “打开顺序消息就行”,这句话省掉的解释,迟早得用排障时间补回来。

同一订单有序处理与不同顺序单元并行执行的流程图

同一订单串行;不同顺序单元并行。版本缺口只阻塞受影响的范围。

2. 需要有序的,通常只是同一个订单

订单 A 先付款再发货,订单 B 没必要跟着等。用订单 ID 做顺序键,同一订单串行,不同订单落在不同处理单元时可以并行。全局排一条队,慢一条就全员陪跑,我一般不会这么选。

这里有三个顺序:业务产生顺序、Broker 写入顺序、数据库提交顺序。 同一个 Key 只能帮助路由,不能替多个生产者判断业务先后。

我的做法是让订单的权威写入方生成事件版本,在修改订单的同一个事务里写入 MQ,再按订单版本串行发送。不让几个服务各自加版本,也别靠机器时间戳猜顺序。

方案 需要怎么配合 保证到哪儿
RocketMQ 4.x 同键选同队列,使用顺序消费监听器 队列内顺序
RocketMQ 5.x FIFO FIFO Topic,设置消息组,同组串行发送 消息组内顺序
Kafka 常规分区消费 稳定键路由到同分区,开启生产幂等,串行处理 分区日志顺序

RocketMQ 5.x 同组有序发送有单生产者、串行发送的前提。Kafka 生产幂等能处理协议重试,但不会给多个生产者的业务事件重新排序。消费端再开无约束线程池,两边都救不了。这是 RocketMQ 4.x、5.x FIFO、Kafka 生产配置的差异。

3. 版本号拦缺口,状态机拦乱改

订单当前版本是 7,收到 8 才能继续;收到 9,先找缺失的 8。已处理的事件 ID 可以幂等返回,陌生的旧版本则需要核对,不能一律当重复扔掉。

下面是抽象伪代码,假设订阅的是完整连续事件流,初始版本为 0。过滤掉中间事件的订阅,不能直接套用加一规则。

ounter(lineounter(lineounter(lineounter(lineounter(lineounter(lineounter(lineounter(lineounter(lineounter(lineounter(lineounter(lineounter(lineounter(lineounter(lineounter(lineounter(line
void onMessage(Event e) {
    try {
        db.transaction(() -> {
            Order o = orders.lockOrCreateInitial(e.orderId);
            if (inbox.exists(handler, e.id)) return;
            if (e.version != o.version + 1) throw new SequenceGap();
            stateMachine.requireAllowed(o.state, e.type);
            orders.applyAndAdvanceVersion(o, e);
            inbox.insertUnique(handler, e.id);
        });
        delivery.confirm(e); // 提交业务后确认,或推进连续完成位点
    } catch (SequenceGap gap) {
        delivery.pauseAndRetainUnprocessed();
        repair.fetchMissingEvents(e.orderId); // 不确认,不跳过
    }
    // 其他异常交给框架重试,不能推进消费进度。
}

业务更新、版本推进、去重记录必须在同一个事务里。行锁或等效并发控制,防止两个消费者同时读到旧版本。外部动作还需要自己的可靠投递与幂等,数据库事务包不住远程接口。

9 堵在队头,8 却在后面,原地重试 9 没用。修复程序要能从权威事件日志补回缺失事件,再恢复消费。只有暂停按钮,没有补齐路径,叫卡死。

不过实际业务流程通常不会真的卡死,这个版本号方案落地成本偏高。我们更多还是靠状态机扭转来兜底。

4. 两个特别容易把顺序拆散的操作

失败消息挪去重试队列,后面的照常跑。 第 2 条没成功,第 3 条先生效,严格顺序已经没了。要么阻塞受影响的订单,要么业务明确允许跳过并补偿,别两头都想要。

关于消费位点:Kafka 的 pause/resume 不会倒退读取位置。失败记录和本批未处理记录要保留,或者重新定位到最早未完成处;提交位点不能跨过未完成记录。

扩分区后,同一个 Key 换了地方。 使用取模路由时,分区数变化可能改变映射,历史消息却仍留在旧分区。Kafka 分区变更时,新旧两边并行消费,顺序就断了。RocketMQ 按队列数量选择队列,也要防这件事。

可以暂停相关生产,等旧路由消费到切换边界,再启用新路由;不能停,就得管理路由版本和业务归属迁移。加机器和加顺序单元,不是一回事。

5. 别为了有序,把所有业务都堵住

订单生命周期、账户内事件适合局部有序。完整快照刷新缓存,如果业务允许高版本覆盖低版本,可以简化;逐笔资金和库存变化,不能随便丢掉中间步骤。

上线前,我会故意让发货先于付款到达,制造重复投递,再带着积压扩一次分区。验收就看三件事:不越级生效,缺口能告警,修复后能继续。平时除了总积压,还要看最老消息等待时间和阻塞的业务键。

让 MQ 尽量按序送,让业务在顺序出错时拒绝做错。

最后实际生产远远比这个复杂:重复推送幂等、幽灵事件隔离、业务补偿事件等等。




上一篇:Agent 为什么需要自己的 Git?大模型状态版本控制全解析
下一篇:数据漂移、概念漂移怎么分?Python 漂移检测与 Evidently 实战
您需要登录后才可以回帖 登录 | 立即注册

手机版|小黑屋|网站地图|云栈社区 ( 苏ICP备2022046150号-2 )

GMT+8, 2026-9-25 03:09 , Processed in 0.852613 second(s), 40 queries , Gzip On.

Powered by Discuz! X3.5

© 2025-2026 云栈社区.

快速回复 返回顶部 返回列表