ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Storm消息可靠性核心:Anchoring血缘链原理与实战

Storm消息可靠性核心:Anchoring血缘链原理与实战 刚接触Storm那会儿我对“可靠消息处理”的理解就是简单的“不丢数据”。但真到了生产环境消息还是会莫名其妙地重复或者丢失尤其当一个Tuple经过多个Bolt、产生大量分支处理后问题更难定位。后来我才真正搞明白Storm的可靠性核心不在Spout也不在Bolt而在于那套叫Anchoring的机制——它用一条看不见的“血缘链”把每个消息从生到死的处理过程串起来靠一个异或算法判断整棵处理树是否完成。这篇就把Anchoring机制从设计原理到实战坑点完整拆开讲清楚适合正在用Storm做实时计算、被消息确认问题折磨过的开发者也适合那些想深入理解Storm内部机制的读者。1. 为什么需要Anchoring先理解Storm的可靠性承诺1.1 从“消息不丢”说起at-least-once的真实含义很多教程一上来就告诉你Storm提供“至少一次”的处理保证但很少有人讲清楚这个保证到底是怎么实现的。所谓at-least-once翻译成人话就是每条消息至少被处理一次但可能被处理多次。Storm能保证的是“不会因为进程崩溃、网络超时而导致数据凭空消失”但它做不到“每条消息恰好被处理一次”。之所以是这个语义是因为Storm的容错策略是“超时重发”。一个Spout把消息发出去之后如果在一段时间内没有收到确认它就会重新发送这条消息。问题来了Storm怎么知道一条消息“处理完成”了一条消息从Spout发出来经过Bolt A做清洗又经过Bolt B做聚合再到Bolt C写数据库中间可能产生了上百个派生TupleStorm怎么区分“还在处理中”和“已经处理完”这就是Anchoring机制要解决的核心问题。简单说Anchoring让Storm能够追踪每一个Tuple的“来龙去脉”知道哪些新的Tuple是由哪些旧的Tuple衍生出来的从而构建出一棵以Spout Tuple为根节点的血缘树。1.2 Spout、Bolt、Tuple血缘链里的三个核心角色在讲具体机制之前先把角色理清楚。Spout是消息的入口它从外部数据源读取数据并发射Tuple。Spout发射Tuple时可以携带一个messageId一旦携带了这个IDStorm就会认为这是一个“需要确认的消息”进入可靠性追踪流程。对应的Spout会收到ack或者fail回调告诉你这条消息整条链路处理成功或失败。Bolt负责处理逻辑。它接收上游Tuple处理后可能发射新的Tuple到下游。这里最关键的动作是Bolt在发射新Tuple时可以选择声明“这个新Tuple是从哪个旧Tuple来的”。一旦声明新Tuple和旧Tuple之间就建立了父子关系血缘链就延伸了一步。Tuple是Storm里的数据单元本质上是一个字段列表。但在血缘链的语境下每个被追踪的Tuple都有自己的身份。Spout发射的根Tuple会生成一个rootId和一个uniqueId而每个派生Tuple也都有自己独立的uniqueId。这些ID是64位随机数是整个追踪机制的基础。打个比方Spout Tuple像是一个公司的原始订单Bolt处理的每一步像是在订单上追加了子任务。Anchoring机制就是那张订单追踪表记录着这个订单派生出了哪些子任务每个子任务目前是完成、失败还是仍在执行。1.3 为什么不能只确认Spout Tuple本身看到这里你可能会想那我让每个Bolt处理完都往Spout回个确认不就行了比如Bolt C处理完了直接告诉Spout说“这条消息搞定了”这个思路在单链路上可行但一旦遇到分叉和聚合就乱了套。举个实际情况Spout发出一条数据Bolt A同时把它发给Bolt B和Bolt CB和C各自处理后发给D。假设D最后把结果写入数据库并确认成功。如果D直接把“成功”告诉Spout但这时候B还没处理完呢再如果B处理的是“清洗逻辑”它的结果会直接影响D写入的内容B没跑完D的结果其实是错的可Spout已经收到了成功信号这条消息就算“被处理了”数据错误就这么产生了。所以Storm的确认逻辑必须覆盖整条血缘链。只有当这棵血缘树上的所有节点都处理完成才允许通知Spout这条消息ack成功。而Anchoring和它背后的异或算法正是实现“全树确认”的关键。2. Anchoring核心原理一次异或运算撑起的消息血统2.1 一个Tuple就是一次“未完成的约定”Tuple Tree与全树确认Storm把每条被追踪的Spout消息定义成一棵Tuple Tree根节点是Spout发射的那个原始Tuple子节点是所有由它直接或间接派生出的Tuple。注意这里的“派生”不是数据内容的复制而是通过Anchoring建立的血缘引用。只要树中任何一个节点没有被成功处理整棵树的最终状态就是“未完成”。只要一个节点处理失败整棵树就会被判定为失败。这个设计保证了下游任何一个环节出错都能追溯到根节点触发Spout的fail回调以及后续的重发。那么问题来了Storm是怎么在分布式环境下记录这棵树的要知道一个Tuple可能被分发到不同机器的不同Worker上执行如果采用“每个Tuple处理完就向Spout汇报”的方式网络通信量会爆炸。Storm采用了一种相当巧妙的做法它不显式保存整棵树的结构只保存一个校验值通过异或运算增量更新最终通过校验值是否为0来判断整棵树是否完成。2.2 ack value的异或算法从1到0的追踪过程异或运算这里值得单独拿出来讲因为理解了它你就理解了Anchoring的所有后续逻辑。先明确两个概念。每个被追踪的Tuple在emit时都会获得一个系统生成的64位uniqueId。同时每个Spout Tuple会有一个rootId这个rootId是这棵血缘树的“种子”。acker任务为每个Spout Tuple保存一个ack value初始值就是该Tuple的rootId。处理过程遵循三条规则当一个Spout Tuple被发射时系统在acker里登记一条记录ack value初始为rootId。当一个Bolt用Anchoring方式发射新Tuple时会把这个新Tuple的uniqueId“做进”当前记录里。具体做法是在Bolt处理完并调用ack(inputTuple)时系统会把该新Tuple的uniqueId和当前ack value进行异或。顺便说一句新Tuple的uniqueId在emit时就会与当前ack value异或一次而在其对应的处理完成被ack时又会异或一次两次异或等于抵消这个细节后面细说。当一个Tuple处理失败被fail时对应的uniqueId同样会被异或进去但因为失败会导致整棵树失败所以最终ack value不会归零spout会收到fail回调。核心技巧在于同一个uniqueId被异或两次之后结果会回到原始状态。类比现实中的开关每个Tuple的uniqueId就像是一个“灯开关的编号”你拨动一次开关Tuple产生再拨动一次Tuple处理完确认开关就恢复原样。整棵树的完成状态就是所有开关都合上了系统最终算出的ack value变回初始值rootId与rootId异或为0。来看一个具体的数值例子。假设Spout发射根Tuple系统分配rootId 11这里用十进制方便阅读实际是64位随机整数同时在acker里创建记录ack value 11。接着Bolt A收到根Tuple后emit了子Tuple A1系统给它分配uniqueId 25。Bolt A处理完毕后调用ack(rootTuple)此时系统执行第一步发生在emit时ack value 11 XOR 25第二步发生在Bolt A ack根Tuple时ack value (11 XOR 25) XOR 25 11看到了吗A1产生时异或一次确认时又异或一次两次之后A1的ID从ack value里消失了。也就是说一个子Tuple的加入和完成在数学上被完美抵消。但如果A1还没完成它的ID就会一直留在ack value里导致整体数值偏离rootId。继续往下。假设Bolt A还把根Tuple发射给了Bolt B生成了子Tuple A2uniqueId 37。那么当Bolt A ack根Tuple时ack value 11 XOR 25 XOR 37。此时Bolt B正在处理A2。A1已经完成抵消了A2还“挂”在账号里。等Bolt B也完成并ack(A2)时ack value 11 XOR 25 XOR 37 XOR 37 11 XOR 25。等等这样好像不对这里要说清楚A1虽然被Bolt A ack了但它是在A1这个子Tuple处理完成时被ack的所以A1的ID在那一刻被异或消除了留下了A2的ID。当Bolt B完成A2并ack时A2的ID也被消除最终ack value回到11。然后是Bolt A对根Tuple的ack以及根Tuple自身在Spout侧的处理……实际流程中还有一个细节Spout在emit根Tuple时rootId会参与一次异或而Spout收到最终成功时会再次异或rootId最终归零。简化而准确地说整个追踪机制就是用异或把树中尚未完成的Tuple的ID累加记账完成的Tuple逐步抵消自己的ID所有节点完成则最终归零。2.3 Acker是如何“记住”整棵树的Task与内存状态这个“记账”工作由Storm的Acker组件完成。Acker不是独立的服务而是拓扑里一类特殊的Bolt由框架自动创建。它为每个正在追踪的Spout Tuple维护一条记录记录内容为字段作用rootId标识是哪棵血缘树ack value当前校验值通过异或增量更新taskId记录是哪个Spout Task发射了这条消息用于回调超时时间超过该时间仍未归零则判定失败Acker是内存态的它会定时扫描超时记录。如果你的拓扑设置了较短的message timeoutAcker会频繁触发失败。如果Acker任务的并发度不够大量的Tuple确认消息会排队造成人为延迟甚至误判超时。所以配置Acker并发度是个常见调优点后面会详细说。3. 实操篇怎么正确构建和维护血缘链3.1 第一关emit时一定要“认亲”Anchoring的完整语义说白了就一句话你在Bolt里发射新Tuple的时候到底有没有告诉系统这个新Tuple的“父亲”是谁。看两段代码对比。第一段是错误示范public void execute(Tuple input) { String word input.getStringByField(word); // 错误没有传入input作为anchor collector.emit(new Values(word)); // 直接ack collector.ack(input); }这个Bolt处理完输入Tuple后确实发射了一个新Tuple也正确ack了输入Tuple。但问题大了新发射的Tuple没有声明父Tuple它和输入Tuple之间没有任何血缘关系。Storm根本不知道这个新Tuple的存在也不会追踪它下游发生了什么。如果下游Bolt处理失败这个失败根本回溯不到源头Spout不会收到fail回调数据就在这个“断了血缘”的地方悄悄丢了。正确写法是public void execute(Tuple input) { String word input.getStringByField(word); // 正确传入input作为anchor建立血缘关系 ListTuple anchors Arrays.asList(input); collector.emit(anchors, new Values(word)); // 处理完输入Tuple后确认 collector.ack(input); }这里emit(anchors, values)就是Anchoring的核心API。第一个参数是一个Tuple列表表示“新发射的Tuple是由这些父Tuple衍生出来的”。在这个例子里word的父就是input。什么时候用单个Tuple版本当你的新Tuple只由一个输入产生时可以直接写collector.emit(input, new Values(word))效果等同于传入单元素列表。但要养成传列表的习惯因为后面遇到多输入合并场景时你肯定得用列表形式。3.2 第二关多父元组与分支聚合的锚定处理真实业务里比链式处理更常见的是“分叉-聚合”。比如一个订单处理流程Spout发射原始订单数据Bolt A做风控判断Bolt B做库存预占这两个Bolt的结果汇总到Bolt C做最终决策。这个场景下Bolt C的输入来自A和B两个不同的父Tuple它发射结果时必须同时锚定两个父Tuple。public void execute(Tuple input) { // input来自A或B的其中一个 // 假设我们按订单ID做字段分组所以同一个订单的A结果和B结果会进入同一个Bolt C实例 String orderId input.getStringByField(orderId); String type input.getStringByField(type); decisions.put(orderId, type, input); // 只有当同一订单的两个结果都到齐时才向下游发射 if (decisions.isReady(orderId)) { Tuple riskTuple decisions.getRiskResult(orderId); Tuple stockTuple decisions.getStockResult(orderId); String finalDecision compute(riskTuple, stockTuple); // 关键同时锚定两个父Tuple ListTuple parents Arrays.asList(riskTuple, stockTuple); collector.emit(parents, new Values(orderId, finalDecision)); // 两个父Tuple都要ack collector.ack(riskTuple); collector.ack(stockTuple); } }注意两点。第一发射新Tuple时锚定了两个父Tuple那么新Tuple的血缘树就同时挂着两颗分支的ID。在新Tuple最终被确认之前这个订单的ack value始终不会归零因为两个父分支的ID都还“挂”在账上。第二你必须对每个输入Tuple都调用一次ack或者fail。漏掉任何一个即使计算逻辑完全正确这条血缘链也会一直悬着直到超时触发fail。这种“等到所有分支到齐再处理”的场景在Storm里常用fieldsGrouping把同一业务ID的数据路由到同一个Bolt实例上保证每个Bolt实例看到的都是完整的分支集合。如果没有这个路由保证你的Bolt C可能在一个实例收到风控结果在另一个实例收到库存结果那你根本无法聚合。3.3 第三关异步操作与ack时机别在子线程里提前ack这是我在生产环境踩过最深的一个坑。很多业务都会在Bolt里做异步调用比如发HTTP请求、读写外部缓存。为了不阻塞Worker线程有人会把异步操作丢到线程池里执行然后在execute方法里立刻调用collector.ack(input)。这是致命错误。看这个错误代码public void execute(Tuple input) { asyncClient.call(input.getStringByField(data), new Callback() { public void onSuccess(String result) { // 异步回调里才真正处理完了数据 collector.emit(input, new Values(result)); collector.ack(input); } }); // 错误这里就ack了异步操作可能还没执行完 collector.ack(input); }这段代码在外面那层execute里就ack了输入Tuple可异步回调还没执行呢。Storm以为这条消息处理完了血缘链的账被错误结清。实际业务里异步回调可能在几秒后才落地如果这时候Bolt挂掉了或者外部服务超时了这条消息已经不可能再被重试了——Storm认为它已经成功了。正确的做法是把ack移到真正处理完成的地方比如异步回调中还要做好超时处理public void execute(Tuple input) { asyncExecutor.submit(() - { try { String result asyncClient.call(input.getStringByField(data)); collector.emit(input, new Values(result)); collector.ack(input); } catch (Exception e) { collector.fail(input); } }); }这里有个隐含的代价execute方法返回后Bolt实例并没有等待异步结果而是立即处理下一条消息。这会导致你的处理并发度变高但代价是如果异步操作积压过多内存中等待的超时Tuple会增多。更稳妥的做法是用固定大小线程池并配合带超时机制的Future在超时后主动fail。3.4 完整示例一个带血缘链的WordCount拓扑下面给一个相对完整但极简的示例把Anchoring的关键写法都串一遍。这个拓扑从Socket读取句子拆分单词并计数。// Spout这里只展示关键部分 public class SentenceSpout extends BaseRichSpout { private SpoutOutputCollector collector; private int index 0; private String[] sentences { hello world, storm anchoring, hello storm }; public void open(Map conf, TopologyContext context, SpoutOutputCollector collector) { this.collector collector; } public void nextTuple() { if (index sentences.length) { return; } String sentence sentences[index]; // 关键传入messageId让Spout Tuple进入追踪流程 collector.emit(new Values(sentence), index); index; } public void ack(Object msgId) { System.out.println(消息处理成功: msgId); } public void fail(Object msgId) { System.out.println(消息处理失败准备重发: msgId); } } // 拆分单词的Bolt public class SplitBolt extends BaseRichBolt { private OutputCollector collector; public void prepare(Map conf, TopologyContext context, OutputCollector collector) { this.collector collector; } public void execute(Tuple input) { String sentence input.getStringByField(sentence); for (String word : sentence.split( )) { // 关键anchor到输入Tuple collector.emit(input, new Values(word)); } // 不要忘了确认输入 collector.ack(input); } } // 计数的Bolt public class CountBolt extends BaseRichBolt { private OutputCollector collector; private MapString, Integer counts new HashMap(); public void prepare(Map conf, TopologyContext context, OutputCollector collector) { this.collector collector; } public void execute(Tuple input) { String word input.getStringByField(word); counts.put(word, counts.getOrDefault(word, 0) 1); int count counts.get(word); // 这里是叶子节点没有下游了但还是要通过emit方式向下游发送 // 如果这个Bolt没有下游则可以只ack不emit collector.emit(input, new Values(word, count)); collector.ack(input); } } // 拓扑配置 public class AnchoringTopology { public static void main(String[] args) throws Exception { TopologyBuilder builder new TopologyBuilder(); builder.setSpout(sentence-spout, new SentenceSpout(), 2); builder.setBolt(split-bolt, new SplitBolt(), 2) .shuffleGrouping(sentence-spout); builder.setBolt(count-bolt, new CountBolt(), 2) .fieldsGrouping(split-bolt, new Fields(word)); Config conf new Config(); // 关键配置acker并发度 conf.setNumAckerExecutors(2); // 关键配置超时时间默认30秒 conf.setMessageTimeoutSecs(30); // 关键配置最多允许多少条spout消息同时处理 conf.setMaxSpoutPending(5000); StormSubmitter.submitTopology(anchoring-word-count, conf, builder.createTopology()); } }这个示例里最值得注意的配置是setNumAckerExecutors。很多人不设置用默认值。但如果你发现消息处理正常、外部系统也正常却频繁出现超时很可能是Acker并发度不够导致Tuple确认事件排队。一般建议设置成与Topology整体Worker数同量级或者直接设为2到4起步观察。4. 血缘链失效的常见场景与排查思路4.1 典型症状消息重复、超时、漏ackAnchoring机制一旦出问题症状通常有以下几种。症状一消息重复。Spout的fail回调频繁触发外部数据源这边不断重发。这通常是血缘链断裂导致的比如某个Bolt在处理中抛了异常但你没有捕获异常调用fail或者你用错误的方式emit导致新Tuple没被追踪。当fail发生后Spout重发原始消息可上游某些分支已经处理完了于是出现了一部分重复计算。症状二消息静默丢失。你看不到fail也没有超时但数据最终就是少了一些。这种最隐蔽。常见原因是Bolt里对输入Tuple调用了不止一次ack。比如你在循环里对同一个input每次emit后都调用ack一次导致同一ID被异或了多次数学账本被搞乱了。或者是你在某个分支里完全没有调用ack或fail但因为其他分支的ack恰好让整体数值归零Storm误认为处理完成了。症状三整条链路无限超时。一个Tuple的处理时间明明很短但拓扑UI上显示spout tuple长时间处于pending状态直到超时。这往往发生在异步场景或者Acker任务本身成了瓶颈处理不过来。另外如果Bolt的execute方法里发生了阻塞比如同步调用外部接口且没设置超时整条链路就会卡住等接口挂了才超时。4.2 排查工具与手段UI界面、日志、ack-value验证排查Anchoring问题我的第一建议永远是先看Storm UI。UI里每个Spout都会显示一条消息的整个“生命周期”状态包括pending的数量、failed的数量、平均完成时间等。如果你看到大量pending并且时间持续增长说明有一批Tuple处理时间异常或者未被正确确认。具体来说可以从三个维度看数据维度关注点判断方向Spout的failed计数fail回调在增长吗血缘链上有节点处理失败Bolt的execute latency哪个Bolt耗时异常高问题在哪一段逻辑单个Tuple的rootId日志能否在日志里找到rootId对应的处理轨迹如果日志里只看到emit没看到ack/fail说明确认丢了如果你的Bolt里打了日志建议主动把rootId打出来。在Bolt中可以通过input.getSourceStreamId()配合input.getValue拿到业务字段但拿不到rootId本身需要从Spout侧把messageId传下来。所以我通常会在每个Bolt里打印输入Tuple的关键业务ID再结合Spout侧的消息ID在日志里串出这条消息经历过的所有节点。这样排查血缘链断裂比对着UI猜测效率高得多。4.3 一个排查实例我以前遇到过这么一个问题一个Kafka实时入仓拓扑每天凌晨数据量大时failed消息暴涨但早上又恢复。表面看像是系统过载可日志显示Kafka消费端持久化成功了数据重复问题也只在高峰期出现。后来我在Spout的fail回调里加了详细日志逐个打印失败消息的msgId又去对应Bolt的日志里查这些msgId的走向。最终发现是某个Bolt在处理批量数据时用了多线程把一批数据拆成多个子任务并发执行而每个子任务完成后都会调用collector.ack(input)。同一个父Tuple被ack了多次导致数学账本错乱最终ack value提前归零或者进入错误状态。在低峰期这个问题被隐藏了因为并发量低、处理顺序固定高峰期并发上来之后多次ack的时序变化才导致一批消息被误判失败。修复方案很简单把多个并发子任务的结果聚合到回调里只ack一次父Tuple。这个案例让我意识到Anchoring的账本逻辑对时序极其敏感每个Tuple的emit和ack必须严格成对出现一旦多配对或者少配对就会产生不可预期的行为。4.4 速查表血缘链常见问题与修复整理一张速查表方便大家直接对照排查。问题可能原因处理方式spout tuple大量超时Acker并发度不足调大setNumAckerExecutorsspout tuple大量超时Bolt中有阻塞IO给外部调用加超时改用异步failed计数持续增长业务处理异常未调用fail在catch块中显式调用collector.fail(input)消息重复但无异常血缘链断裂导致超时重发检查是否漏anchor检查emit是否传了父Tuple消息静默丢失对同一Tuple多次ack确保每个输入Tuple只ack一次数据最终一致但重复外部系统未做幂等配合Kafka offset或数据库唯一键做幂等5. 血缘链之上的可靠性设计从at-least-once到业务幂等5.1 为什么Anchoring不等于exactly-once把Anchoring搞透彻之后你会意识到一个残酷的事实即使你的血缘链构建得再完美Storm也只能保证at-least-once不能做到exactly-once。原因在于Anchoring只能保证“整棵树处理完成后Spout成功确认”但它不能保证“每个节点恰好处理一次”。假设Bolt A处理完Tuple并通过血缘链确认了但结果发送到Bolt B时网络抖动Bolt A的重发机制又触发了一遍那Bolt B就会处理两次。Storm的容错模型是“至少一次”系统为了不丢数据宁可重复处理。另外即使没有网络问题消息重发也可能来自Spout层。比如Spout发出去后Acker迟迟没确认超时触发重发但此时原始消息可能还在某个Bolt的队列里排队最终两个版本都会被执行。这些都是at-least-once语义的天然代价。所以不要把Anchoring理解为“可靠传输协议”它更像是“缺陷检测协议”。它能让你快速发现哪些消息没处理完但它不能保证处理过程只发生一次。要想业务上不产生重复数据必须在应用层做幂等。5.2 与消息队列/数据库搭配的可复用兜底方案生产环境里最常用的方案是“Anchoring 外部幂等”。具体做法分两步走。第一步在Spout里配合Kafka手动offset管理。KafkaSpout本身已经支持把offset提交和tuple ack绑定在一起。只有当前的tuple被完整ack后KafkaSpout才会提交offset。这样即使进程崩溃导致消息重放也能保证从最后一个成功处理的位置继续消费避免大范围重复。第二步在写入外部存储时建立天然幂等。比如MySQL表设置唯一键业务ID写入时用INSERT ... ON DUPLICATE KEY UPDATE或者REPLACE INTO。这意味着无论那条消息被处理多少次最终落库的结果都是一样的。Redis里用SETNX或者带过期时间的set也能起到类似作用。这里分享一个我常用的模式。实时统计场景下如果结果允许近似就把重复导致的偏差控制在增量更新上。如果结果必须精确就把聚合结果按业务主键存储在外部中间表写入用“比较版本号”的方式。无论哪种思路核心都是让重复执行的结果和单次执行的结果一致。Anchoring负责告诉你哪些消息需要重试而幂等负责让重试不产生副作用。5.3 maxSpoutPending、超时参数与Acker并发度的调优注意点最后说几个Anchoring相关的关键参数这些参数直接决定了血缘链机制和你的业务吞吐能否和平共处。maxSpoutPending控制的是Spout同时能有多少条消息处于“未确认”状态。调大这个值可以提高吞吐因为Spout可以源源不断发射消息不必等待前面的消息全部确认。但调太大会带来一个副作用如果某个消息超时失败它后面排队的消息也会跟着被处理但因为前面的消息已经失败重发了拓扑里可能出现大量的重复计算。经验值是先设成5000起步观察内存和时延再调整。messageTimeoutSecs是消息超时时间默认30秒。这个参数要和你的真实处理时延匹配。如果业务平均耗时2秒设置30秒会宽裕很多但如果你想快速发现失败可以把超时缩到5到10秒。缩太短也有风险比如GC停顿或者流量洪峰时正常处理的消息会被误判超时。setNumAckerExecutors是Acker并发数。这个参数是最容易被忽视的。很多人把注意力放在Spout和Bolt的并发度上却忘了Acker本身也有处理上限。当Topology吞吐较大时Acker会成为瓶颈。一般建议Acker并发度不要小于Spout并发度或者按Worker数量的1到2倍来配。调整之后观察UI上Spout的failed和pending指标是否下降。还有一个小细节如果你的拓扑里存在不需要可靠处理的流比如纯粹的日志采集或者数据源本来就有重放机制可以显式关闭Anchoring。做法是不传messageId给Spout的emit或者在Bolt emit时不传anchor Tuple。这样Storm就不会跟踪这些消息减少了Acker的IO和内存开销。但我还是建议除非你能明确承受丢失否则别省这个成本。风控、计费、订单这类场景一点点丢数据都可能是生产事故。另外有个容易被忽略的点BaseBasicBolt这个类。如果Bolt继承的是BaseBasicBolt框架会自动帮你处理ack和anchor。它会自动把输入Tuple作为anchor传给emit也会在execute方法返回后自动ack输入Tuple。听起来很省事但它也有一个陷阱你的execute里如果有异步操作返回后框架已经自动ack了异步结果还没出来又回到了那个线程安全问题的老坑。所以用BaseBasicBolt时要确保你的处理逻辑是纯同步的否则老老实实用BaseRichBolt自己管理ack。再补充一个关于多Stream的细节。如果Bolt通过collector.emit(streamId, anchors, values)向不同Stream发射数据这些Stream上的Tuple同样属于同一棵血缘树。不要因为Stream不同就以为血缘链会断开。只要anchor传了同一个父Tuple无论发往哪个Stream最终都会被追踪。还有一点是directGrouping和Anchoring的关系。使用direct grouping时你需要通过emitDirect指定下游TaskId此时依然要传anchor。不要因为直连就跳过Anchoring血缘链和路由方式没有直接关系。这些年和Storm打交道下来我最大的体会是Anchoring是Storm这套可靠性体系里最容易学、也最容易出错的地方。网上大部分教程都在讲API怎么调用但真正决定系统能不能扛住事故的是你对它背后那套“异或账本”理解得有多深。一个Tuple的失联可能不是因为网络而是因为某个Bolt里少传了一个anchor。当你学会把每条消息看作一个需要跟进的项目把emit和ack看作项目进度的打卡你就能在复杂的分布式环境里稳稳抓住那条看不见的血缘链。希望这篇内容能帮你在自己的拓扑里少踩几个坑把可靠消息处理这杆秤真正端平。
返回列表