上周有个朋友在群里发了一段 Storm 拓扑的截图,问为什么明明把 Bolt 的并行度都调到 4 了,处理订单事件还是乱套,create 事件居然比 update 晚到。这个问题我在做流式计算这几年里见过太多次了,几乎每个接触 Storm 的人,都要在“并发控制”和“顺序执行”之间撞一次墙。因为 Storm 默认是分布式并行引擎,天然把数据打散到多进程多线程里跑,可业务上偏偏需要某些数据严格按照先后顺序被处理——这两件事天然互相拉扯。
这篇就围绕 Storm 并发控制与顺序执行这件事,把 worker、executor、task 之间的并发模型说清楚,再讲清楚 fieldsGrouping 到底保证了什么、不保证什么,最后给出几种真正能拿到“精准排序”结果的方案,以及我在实际项目里调参、踩坑的具体过程。无论你是刚入门分布式流处理,还是已经被乱序问题折磨了一阵子,这篇文章里应该有你能直接拿去用的东西。
1. Storm 的并行模型:worker、executor、task 三层中顺序到底在哪一层
1.1 从 Topology 到物理进程的三层映射
很多人写 Storm 代码时只关心 Spout、Bolt 和它们之间的连线,觉得画完一个有向无环图,提交上去,数据自然会按照图的方向流动。这个理解没错,但到了调度层面,一个 Topology 会被拆成三层物理结构,顺序和并发的问题就出在这三层的关系上。
第一层是 worker。一个 topology 提交到集群后,会被分配一个或者多个 worker,每个 worker 是一个独立的 JVM 进程。你可以用 Config.setNumWorkers(3) 设置 worker 数量,这个数字决定了集群里会有多少个 JVM 进程来跑你的任务。
第二层是 executor。每个 worker 进程里会启动若干线程,这些线程就是 executor。一个 executor 通常对应一个线程,它负责执行一个或者多个 task 实例的代码逻辑。
第三层是 task。task 是真正跑 Spout 或 Bolt 实例的最小单元。你在 builder.setBolt("order-bolt", new OrderBolt(), 4) 里写的这个 4,指的是这个 Bolt 的 task 数量,也是它的并行度。在新版 Storm 里,默认一个 executor 只负责一个 task,所以并行度 4 通常就意味着 4 个 executor 线程,分别跑 4 个 OrderBolt 实例。
为了好理解,我列个表对照一下。
| 层级 | 维度 | 类比 | 并行度控制方式 |
|---|---|---|---|
| worker | JVM 进程 | 一家店的独立分店 | Config.setNumWorkers() |
| executor | 线程 | 分店里的收银员 | setBolt(..., 并行度) 结合最大 executor 数 |
| task | Bolt/Spout 实例 | 收银员处理的柜台队列 | setBolt() 的并行度参数,或显式 setNumTasks() |
1.2 并行度参数到底在调什么
搞明白三层结构后,再看并行度参数就不会发怵了。worker 数是进程级并发,executor 数是线程级并发,task 数是逻辑实例数。很多人以为把 task 数调大,单线程的处理能力就上去了,其实不一定。如果 executor 数量没变,task 只是在这个线程上排队等待调度。真正让数据“同时被处理”的,是 worker 和 executor 的并发;task 数量决定的是状态切分的粒度。
举个例子。一个 OrderBolt 并行度为 8,就意味着有 8 个 task 实例,同时最多有 8 个线程在跑这个 Bolt 的 execute 方法。如果一个 executor 内配置了多个 task(老版本可以这么干),那这些 task 依然在这个线程里轮流执行,并没有真正并行。
顺序性在这个模型里的位置就很清楚了:同一个 task 内部,execute 方法是在单线程里被依次调用的,所以天然有序。但只要数据被分发到了不同的 task,甚至只是两个不同的 executor 线程,它们之间就没有任何先后约束,全看系统调度和网络延迟。
1.3 顺序位置的结论
所以你要记住一个很关键的前提:Storm 的并发是“并行处理”层面的并发,不是“单条消息内部”的并发。想在一个 Bolt 内部拿到所有消息的全局顺序,唯一的可能就是这个 Bolt 只有一个 task、只有一个 executor。一旦并行度大于 1,顺序只能存在于单条消息路径的局部位置,也就是“同一个 key 被路由到同一个 task 后,这个 task 内处理它们的先后顺序”。
这就提出了一个核心问题:如何在不牺牲并行度的情况下,让业务关心的关键消息依然有序?这才是 Storm 并发控制真正难的地方。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. fieldsGrouping 的真实边界:按 key 路由,不按时间排序
2.1 为什么同一个 key 会被送去同一个 task
fieldsGrouping 是 Storm 内置分组策略里最常用的一个,专门解决“同 key 数据去同一个 task”的需求。它的实现原理是对指定字段的值做哈希,再对 task 总数取模,把元组分发给对应的 task。
例如,订单系统里所有事件都带一个 orderId,我这样写:
java复制TopologyBuilder builder = new TopologyBuilder();
builder.setSpout("order-spout", new OrderSpout(), 2);
builder.setBolt("order-bolt", new OrderAssemblerBolt(), 8)
.fieldsGrouping("order-spout", new Fields("orderId"));
这段代码的意思是:所有 orderId 相同的订单事件,会被分到同一个 order-bolt task 里。因为同一个 task 的 execute 是单线程执行的,所以从“这个 task 内部看到的处理顺序”来说,同一个订单的事件是连续的。
这是 fieldsGrouping 能给你的全部保证:同一字段值的元组,会在同一个 task 内被顺序处理。不同字段值之间的顺序,完全没有保证。
2.2 三个最常见的执行顺序误解
我刚用 Storm 的时候,在这个问题上栽过三次跟头,分别对应三种不同程度的误解。
第一个误解是:用了 fieldsGrouping,同一个 key 的所有事件就必然是时间顺序。实际情况是,fieldsGrouping 只保证路由结果,不保证原始数据的到达顺序。上游 Spout 从 Kafka 拉数据时,如果并行度是 2,两个 Spout 实例会从不同分区拉取数据,而 Kafka 不同分区间并没有全局先后关系。一个订单的事件可能一部分在分区 0,一部分在分区 1,分区 1 的消息先被拉到,它就可能先进入 order-bolt,哪怕它的事件时间更晚。所以在源头,顺序就已经可能被打乱了。
第二个误解是:同一个 task 内就一定按业务顺序处理。单线程处理,指的是 execute 方法一个接一个执行。但如果我在 execute 里把消息丢给一个线程池或者异步 RPC 后再返回,那么业务逻辑的完成时序立刻就不再受控。很多“明明用了 fieldsGrouping 还是乱序”的问题,最后查下来都是 Bolt 内部自己异步化了。Storm 的单线程保证,只在 execute 方法同步执行时成立。
第三个误解是:全局顺序可以通过 Grouping 搞定。fieldsGrouping 根本不面向全局顺序,“全局”这个词和它没关系。你想让所有订单事件严格按时间序被下游消费,需要的是全局排序,fieldGrouping 做不到,后面我会讲怎么设计。
2.3 路由字段自身才是最大的乱序源头
这一节值得单独拿出来说。fieldsGrouping 是按字段值哈希路由的,那么路由字段值必须稳定且类型一致。我记得有一次压测环境里,上游 Kafka 的消息由两个不同版本的服务生产,一个把 orderId 写成 String 类型,另一个在某个字段为 null 时直接把整个字段丢成了 null。结果在同一时刻,同一个 orderId 的 create 事件和 update 事件被哈希到了两个不同的 task,顺序瞬间瓦解。
这种问题不是 Grouping 的锅,是数据质量的问题。但它在实际生产里特别常见,尤其是跨团队协作,多个系统往同一个 Kafka topic 里写数据时。我的做法是在 Storm 入口的 Spout 或者第一个 Bolt 里做统一清洗,把路由字段强制转成约定类型,对 null 值做默认映射,然后再往下游发。
这里顺便说一个数据库领域的老问题:Oracle 执行 WHERE id IN (3,1,2) 时,返回结果并不会按 3、1、2 的顺序给你,必须显式加 ORDER BY。流处理也一样,fieldsGrouping 只是“IN 分组”的路由逻辑,它不承载排序语义。想让数据有序,你必须自己显式地设计顺序机制,不能指望框架隐式帮你排好。
3. 真正需要全局顺序时:三种方案与代价
有些场景,比如账户余额变动、库存扣减、状态机流转,业务上就是要严格全局顺序。并行与精确排序在这里是硬冲突,但只要分清楚需求和代价,还是有路可走的。
3.1 单线程串行:最笨但最可靠
最简单的方案是放弃并行。把关键 Bolt 的并行度设成 1,或者用 globalGrouping 把上游所有消息都发到 task id 最小的那个 task。
java复制builder.setBolt("strict-order-bolt", new StrictOrderBolt())
.globalGrouping("order-spout");
globalGrouping 负责把整个 stream 的全部元组发送到同一个 task,也就是强制把所有数据汇聚到一个处理线程里,这样下游自然严格有序。代价很明显:吞吐量被单线程卡死。如果这个流本身每秒只有几百条消息,完全没问题;如果每秒钟几十万条,这个方案会直接成为整个拓扑的瓶颈。
我一般只在配置下发、字典更新、低频状态快照这类场景用这种方式。它不优雅,但足够可靠,排障时也最简单。
3.2 编号 + 乱序窗口:保留并行的排序方案
另一种思路是源头给每条消息编号,下游按编号排序。这个思路就像数据库的 MVCC 多版本并发控制:每个事务都有版本号,系统根据版本号判断先后顺序。流不也一样吗?你要判断两条消息谁先谁后,首先得给它们一个可比较的“版本”,没有版本号的流,永远说不清顺序。
具体做法分成三步。第一步,在源头 Spout 里给每条消息赋一个自增序号,或者直接利用 Kafka 的 offset 作为序号。同一分区内,Kafka offset 天然递增,这个序号是可信的。第二步,用 fieldsGrouping 按业务 key 路由,保证同一个 key 的编号消息进入同一个 task。第三步,在 Bolt 里维护一个优先级队列和最新已处理序号,只按顺序把消息吐出去。
我写过这样的排序 Bolt 骨架,核心逻辑如下:
java复制public class OrderAssemblerBolt extends BaseRichBolt {
private OutputCollector collector;
private Map<String, PriorityQueue<OrderEvent>> buffers;
private Map<String, Long> latestSeq;
private long maxOutOfOrderMs = 3000;
@Override
public void prepare(Map stormConf, TopologyContext context, OutputCollector collector) {
this.collector = collector;
this.buffers = new HashMap<>();
this.latestSeq = new HashMap<>();
}
@Override
public void execute(Tuple tuple) {
try {
OrderEvent event = (OrderEvent) tuple.getValueByField("event");
String orderId = event.getOrderId();
long seq = event.getSeq();
buffers.computeIfAbsent(orderId, k -> new PriorityQueue<>(
Comparator.comparingLong(OrderEvent::getSeq)))
.offer(event);
long next = latestSeq.getOrDefault(orderId, 0L) + 1;
PriorityQueue<OrderEvent> queue = buffers.get(orderId);
while (!queue.isEmpty() && queue.peek().getSeq() <= next) {
OrderEvent ready = queue.poll();
collector.emit(new Values(ready));
next = ready.getSeq() + 1;
}
latestSeq.put(orderId, next - 1);
collector.ack(tuple);
} catch (Exception e) {
collector.fail(tuple);
}
}
}
这段代码做的事情,是把乱序到达的事件在内存缓冲里先“攒着”,等序号连上了,再按顺序向下游放行。它保住了上游的并行能力,代价是内存缓存和额外延迟。
3.3 事件时间 + 水位线:为迟到数据留量
如果业务关心的是“事件发生时间”而不是“进入系统的时间”,那么纯编号方式可能不够,因为你不能保证 Kafka 里的消息就是按事件时间排好的。这时需要给数据打上事件时间戳,然后引入类似水位线的机制。
水位线的思想很简单:每条消息带上事件时间,系统维护一个已经安全推进的时间水位线。当某条消息的事件时间早于水位线时,就认为它的所有前置数据都已经到达,可以把缓冲里排序好的数据放行。这个时间水位线的推进速度,决定了你愿意为“乱序容忍”付出多少等待代价。
在 Storm 里没有内建的 Watermark 实体,我通常用定期的 tick 流(Config.TOPOLOGY_TICK_TUPLE_FREQ_SECS)来驱动水位线推进,或者由上游 Spout 定期广播一个水位线事件。实现不复杂,但需要自己控制推进逻辑,容易出错的地方是水位线推进过快会漏掉迟到数据,推进过慢会增加无谓的延迟。
3.4 三种方案对比
把三种方案摆在一起看,选择逻辑就清楚了。
| 方案 | 顺序严格度 | 吞吐 | 延迟 | 实现复杂度 | 适用场景 |
|---|---|---|---|---|---|
| 单线程 + globalGrouping | 严格全局顺序 | 低,受单线程限制 | 低 | 低 | 低频配置流、状态快照 |
| 编号 + 乱序窗口 | 按序号严格有序 | 中高,可并行 | 增加等待时间 | 中 | 源头有序、下游需要排序 |
| 事件时间 + 水位线 | 按事件时间有序 | 中高,可并行 | 可控,取决于水位线 | 高 | 上游多个分区的数据混排 |
实际项目里,绝大多数业务需要的不是“全局严格顺序”,而是“按业务 key 的局部顺序”。这时候根本不用上全局排序,把 fieldsGrouping 用对,再把乱序窗口的容忍度调好,就已经能解决 80% 的问题。
4. 顺序执行与可靠性机制的纠缠:ack、重发和 try-finally
4.1 重发机制是如何把顺序彻底搅乱的
聊到顺序执行,不能绕开 Storm 的可靠性机制。默认开启消息追踪后,Spout 发射的每条消息都会形成一棵 tuple 树。Bolt 每次调用 collector.ack(tuple),Acker 会更新追踪状态;全部节点都 ack 后,Spout 认为这条消息成功处理。一旦某个 Bolt 处理失败,或者整个 tuple 树超时,Spout 的 fail 方法会被触发,通常是重新发射一条一模一样的消息。
这个重发机制给顺序带来的麻烦很直接:一条旧消息在失败后重新进入流,如果它的业务时间比当前已经处理的消息更早,就会在排序窗口里插入一个“过去的版本”。如果没有幂等保护,它会覆盖掉后面已经更新过的数据,效果等于数据库里的旧事务把新事务回滚了。
所以顺序执行一定要和 ack/fail 机制一起设计。你设计的不是一条消息从进到出的单向流程,而是一个可能被重复注入的、有回溯风险的流程。
4.2 状态更新 + ack 的正确代码结构
很多人在写 Bolt 时习惯这样:
java复制@Override
public void execute(Tuple tuple) {
// 业务处理
businessLogic(tuple);
// 最后 ack
collector.ack(tuple);
}
这个结构本身没问题,但就怕你在业务处理过程中抛出异常,导致后面的 ack 没有执行。这时 tuple 会一直挂着,直到超时再被重发,于是下游就收到重复数据。
一个稳妥的写法是用 try-finally 保证 ack 或 fail 一定会被调用:
java复制@Override
public void execute(Tuple tuple) {
try {
businessLogic(tuple);
collector.ack(tuple);
} catch (Exception e) {
collector.fail(tuple);
}
}
try-finally 这个代码结构在 Java 里是最常见的保证“无论中间发生什么,收尾动作都会执行”的手段。放在 Storm 的语义里,它的作用就是确保消息生命周期有明确的终结节点,不会因为一场异常变成无人认领的悬空消息,进而导致超时重发和乱序雪崩。
有一点要提醒:fail 之后,Spout 到底重不重发,取决于你自定义 Spout 的 fail 实现。如果直接返回,消息就丢了;如果重新 emit,顺序问题又回来了。我的建议是重发必须保留原始序号,不要重新生成新序号,否则排序窗口在“下一批新数据”和“重发的旧数据”之间会彻底失去判断依据。
4.3 幂等去重:让重发不再伤害顺序
即便 ack 流程写得再规整,网络抖动、进程崩溃依旧可能造成重发。想要顺序稳,下游必须对重复数据有天然的免疫力。
最简单可靠的做法是“业务主键 + 序号”去重。比如订单事件有一个全局唯一的 eventId,或者我们可以把“orderId + seq”拼成一个唯一键。Bolt 在执行业务逻辑前先判断这个唯一键的序号是否已经处理过。处理过就丢掉,不处理过才更新状态。
判断的方式要小心,不能只靠内存 Map,因为进程重启后内存就没了。正经做法是借助外部存储,比如 Redis 或者数据库唯一索引:把这唯一键插入表里,主键冲突就说明重复,直接跳过。数据库层面再用一个版本号做乐观锁,保证即使两条线程并发更新,最终提交的也是更新的版本。
这套“去重 + 版本号”的组合,看起来不是在讲顺序,但它是顺序执行能长久稳定运行的兜底。顺序系统最怕的不是慢,也不是并发高,而是数据被无声无息地重复处理,把状态覆写成一个错误的结果。
5. 实战:订单事件流从“并发乱序”到“精准排序”
5.1 场景与整体设计
还是回到开头的订单场景。Kafka 里有一个订单事件 topic,里面是下单、支付、发货等事件。下游要维护订单状态机,必须保证同一个订单的事件按业务顺序处理,否则可能出现“已发货”又被“待支付”覆盖的荒唐结果。
整个拓扑结构分两段。第一段是入口清洗,Kafka Spout 消费原始消息,清洗字段后统一为 OrderEvent 对象,同时带上事件序号。第二段是核心排序 Bolt,按 orderId 做 fieldsGrouping,内部维护乱序窗口,每放行一条有序事件,就调用一次下游回调,将结果写入目标存储。
选型上我没有用全局排序,也没有把并行度降成 1,原因很现实:订单量每天几百万,单线程根本扛不住。用“按订单维度分区 + 每个分区内局部排序”的方式,既保住了并发吞吐,又拿到了业务真正需要的顺序。
5.2 排序 Bolt 的代码骨架
下面这段是精简过的核心代码,去掉了存储细节,保留了排序和去重的骨架。
java复制public class OrderedOrderBolt extends BaseRichBolt {
private OutputCollector collector;
private Map<String, PriorityQueue<OrderEvent>> buffer;
private Map<String, Long> seqMap;
private long maxOutOfOrderMs = 1000;
@Override
public void prepare(Map stormConf, TopologyContext context, OutputCollector collector) {
this.collector = collector;
this.buffer = new ConcurrentHashMap<>();
this.seqMap = new ConcurrentHashMap<>();
}
@Override
public void execute(Tuple tuple) {
OrderEvent event = (OrderEvent) tuple.getValueByField("event");
long seq = event.getSeq();
long current = seqMap.getOrDefault(event.getOrderId(), 0L);
// 重复或者旧消息,直接丢弃
if (seq <= current) {
collector.ack(tuple);
return;
}
buffer.computeIfAbsent(event.getOrderId(), k -> new PriorityQueue<>(
Comparator.comparingLong(OrderEvent::getSeq)))
.offer(event);
// 持续放行已经连号的连续事件
PriorityQueue<OrderEvent> queue = buffer.get(event.getOrderId());
long expect = current + 1;
while (!queue.isEmpty() && queue.peek().getSeq() == expect) {
OrderEvent ready = queue.poll();
collector.emit(new Values(ready));
expect++;
}
seqMap.put(event.getOrderId(), expect - 1);
collector.ack(tuple);
}
@Override
public void declareOutputFields(OutputFieldsDeclarer declarer) {
declarer.declare(new Fields("ordered-event"));
}
}
注意,这里我要求同一个 orderId 的所有事件必须通过 fieldsGrouping 进入同一个 task,所以 buffer 和 seqMap 虽然声明为并发容器,但实际只有一个线程在访问某个 orderId 对应的键,不会出现同时读写冲突。如果业务里有跨字段路由的需求,比如同一个用户的所有订单都要顺序处理,那就用 userId 做路由字段,这样同一用户的数据也会进同一个 task。调整路由字段就可以灵活适应不同粒度的顺序需求。
5.3 并行度与窗口参数调优实录
第一次压测我把并行度直接拉到 8,Kafka Spout 并行度 4,worker 数 3,maxPending 设了 10000。结果乱序率很高,大概有 5% 的数据在排序 Bolt 里等了很久才被放行。后来排查发现,maxPending 太大导致 Spout 一次性放太多消息进流,而 Spout 从两个分区拉取的消息顺序不一致,同一条订单事件在整条链路里的序号差被拉开。
我把 maxPending 从 10000 调到了 2000,乱序窗口从 3 秒降到 1 秒,效果立刻好了很多。核心原因在于 maxPending 控制了在途消息的数量,在途消息越少,同一个 key 的事件越容易在短时间内聚齐。如果 maxPending 过大,大量消息拥堵在传输链路中,不同分区的先后差异会被放大,排序窗口就必须等更长的时间才能把数据集齐。
最终参数落在 worker 数 3、Spout 并行度 4、排序 Bolt 并行度 8、maxPending 3000、乱序窗口 1.5 秒。这个组合在压测数据 1 万条时,排序正确率 100%,单条消息平均处理耗时在毫秒级,吞吐也能稳定在每秒 2 万条以上。调参这件事没有银弹,不同业务的消息到达模式差异很大,但 maxPending 和乱序窗口这两个参数永远是调节顺序和吞吐的核心旋钮。
5.4 验证与结果
验证排序结果时,我习惯在排序 Bolt 的输出日志里打印每个订单的事件序列,随机抽几个大流量的订单人工核对。再写一个校验任务,扫描最终落库的订单状态变化记录,检查有没有类似“已支付 -> 已下单”的反向流转。测试结束后,用幂等去重逻辑再重放一次整批 Kafka 数据,确认重复消费时状态不被二次覆盖。
这轮验证让我强烈意识到一个问题:顺序是否正确,不是看代码运行有没有报错,而是看数据最终落到存储后的状态。一定要在状态层做校验,否则表面上一路绿灯,实际的数据早已乱序。
6. 几个我踩过的暗坑:性能和顺序都要时怎么办
6.1 在 Bolt 里用线程池,顺序立刻崩溃
这个坑我前面提到过,但值得再说一次。某次我为了提升单 Bolt 的吞吐能力,在 execute 里把业务计算仍给了一个固定线程池,execute 本身很快返回。测出来的吞吐确实漂亮,结果下游对所有数据的处理顺序就乱了套。原因非常简单,Storm 保证的是 execute 方法按序被调用,不保证线程池里的任务按序跑完。
如果确实需要异步化,一定要把“排序”放在异步化之前,也就是说,先完成顺序整理,再把有序的结果交给异步线程做输出。顺序控制必须是单线程的,输出可以并发,这两者不能混在一起。
6.2 路由字段类型不一致,数据分到了两个世界
前面提过 String 和 Long 混用的问题。两个上游服务,一个发 "10023",一个发 10023L,它们在 fieldsGrouping 里的哈希值不同,于是同一个订单被拆到两个 task 里。这意味着这个订单的 create 和 update 事件可能永远无法相遇,秩序必然崩溃。
解决方法是入口统一。我的经验是,凡是作为路由字段的值,一律统一成 String,并且在清洗 Bolt 里做非空校验。流处理系统的最前端一定是数据治理的位置,Routing 字段尤其不能放任自流。
6.3 乱序窗口设置不合理,要么延迟高要么漏数据
乱序窗口设得太大,每条消息都要多等几秒才能放行,实时性全没了;设得太小,迟到消息直接不满足条件被丢弃,顺序对但数据完整度错了。窗口大小没有固定公式,需要根据上游 Kafka 消息到达的抖动范围来定。
我在实践里一般这么做:先统计同一业务 key 的消息在 Kafka 里出现的最长时间差,比如 95% 的事件在 500ms 内到达,99% 在 1 秒内到达,那就把窗口设为 1 到 1.5 秒。这个统计可以定期从线上日志里跑出来,初期拍脑袋不要紧,后面根据丢数据率再调。漏数据的代价通常比延迟高,所以窗口宁大勿小。
6.4 幂等只做内存去重,进程重启就破功
内存 Map 去重在单机单实例下没问题,但 Storm 的 worker 会重启,重启之后内存里的已处理序号全部丢失,重放数据又会重新砸向下游。真正可靠的幂等必须落在外部存储,Redis 可以靠 SETNX,DB 可以靠唯一索引,把“orderId+seq”作为主键。这个投入不算大,但能避免很多暗无天日的半夜故障。
6.5 忽略了 Kafka 分区 rebalance 造成的源头乱序
Storm 的 Kafka Spout 在 worker 数量调整、topic 分区数变化时会触发 rebalance,分区与 Spout 实例的对应关系会发生改变。如果 Spout 在 rebalance 后重新分配了分区,那么某些分区的消费位点可能会有重复或者断裂,带来源头的乱序和数据跳变。
这种问题很难在代码层面完全规避,只能靠监控位点变化来降低影响。我现在的习惯是,把 Kafka 的 offset 作为事件序号的一部分,同时在重放和数据校验时做完整性比对,确保分区变化没有造成某个订单的事件漏读或者重读。
按我现在的习惯,任何要求顺序的流,都会在源头强制加统一序号,而不是指望某一种 Grouping 策略能兜底一切。同时把重发、去重、状态更新这三件事一起考虑,而不是分散着设计。Storm 的并发控制从来不是某一处配置能解决的,它是一条贯穿消息产生、路由、排序、状态落库全链路的思路。每次看到把并行度调大后数据乱成一团的案例,我几乎都能从上面几个暗坑里找到对应答案。希望这篇记录,能帮你少走几次我正在走的弯路。
