I´ve been working on a Storm Topology and I´m facing some tuple failures. I suspect that one of the bolts is not acking in a specific case causing these failures by timeout. Is there a way in Apache Storm API (0.10.0) to identify which bolt is not acking as expected?
Let´s suppose we have MySpout, BoltA and BoltB as components of this topology and MySpout emits tuples for both bolts expecting they´re going to ack after processing tuples. BoltA in executeTuple() method is always acking, but BoltB is acking only for even values it receives. All tuples with odd values are going to fail after 10 minutes they were emmitted.
In this small sample, it´s easy to identify the failed flow. But in a complex system that we can follow dozens of different flows having multiple bolts it´s like looking for a needle in a haystack. Is there any smart way to find this failure?
public class MySpout extends BaseRichSpout {
protected SpoutOutputCollector collector;
//...
@Override
public void nextTuple() {
Integer msgId = new Integer((int)(Math.random() * 5000 + 1));
collector.emit(new Values(msgId), msgId);
}
@Override
public void fail(Object msgId) {
new Exception("Failed tuple. msgId="+msgId).printStackTrace();
}
}
public class BoltA extends BaseRichBolt {
private OutputCollector outputCollector;
//...
@Override
protected void executeTuple(Tuple input) {
Integer n = (Integer) input.getValues().get(0);
outputCollector.ack(input);
}
}
public class BoltB extends BaseRichBolt {
private OutputCollector outputCollector;
//...
@Override
protected void executeTuple(Tuple input) {
Integer n = (Integer) input.getValues().get(0);
if (n%2==0) {
outputCollector.ack(input);
}
}
}
A timeout value of 10 minutes was configured for this Storm.
<!-- storm config -->
<property>
<name>topology.enable.message.timeouts</name>
<value>true</value>
</property>
<property>
<!-- 10 mins -->
<name>topology.message.timeout.secs</name>
<value>600</value>
</property>
This is the stack trace we see when a tuple fails. I was trying to extract some information of the tuple failure flow. It didn´t say anything that helped me.
java.lang.Exception: Failed tuple. msgId=1234
at MySpout.fail(MySpout.java:127) [myJar.jar:?]
at backtype.storm.daemon.executor$fail_spout_msg.invoke(executor.clj:401) [storm-core-0.10.0-beta1.jar:0.10.0-beta1]
at backtype.storm.daemon.executor$fn$reify__4467.expire(executor.clj:461) [storm-core-0.10.0-beta1.jar:0.10.0-beta1]
at backtype.storm.utils.RotatingMap.rotate(RotatingMap.java:73) [storm-core-0.10.0-beta1.jar:0.10.0-beta1]
at backtype.storm.daemon.executor$fn__4464$tuple_action_fn__4470.invoke(executor.clj:466) [storm-core-0.10.0-beta1.jar:0.10.0-beta1]
at backtype.storm.daemon.executor$mk_task_receiver$fn__4455.invoke(executor.clj:433) [storm-core-0.10.0-beta1.jar:0.10.0-beta1]
at backtype.storm.disruptor$clojure_handler$reify__4029.onEvent(disruptor.clj:58) [storm-core-0.10.0-beta1.jar:0.10.0-beta1]
at backtype.storm.utils.DisruptorQueue.consumeBatchToCursor(DisruptorQueue.java:125) [storm-core-0.10.0-beta1.jar:0.10.0-beta1]
at backtype.storm.utils.DisruptorQueue.consumeBatch(DisruptorQueue.java:87) [storm-core-0.10.0-beta1.jar:0.10.0-beta1]
at backtype.storm.disruptor$consume_batch.invoke(disruptor.clj:76) [storm-core-0.10.0-beta1.jar:0.10.0-beta1]
at backtype.storm.daemon.executor$fn__4464$fn__4479$fn__4510.invoke(executor.clj:578) [storm-core-0.10.0-beta1.jar:0.10.0-beta1]
at backtype.storm.util$async_loop$fn__543.invoke(util.clj:475) [storm-core-0.10.0-beta1.jar:0.10.0-beta1]
at clojure.lang.AFn.run(AFn.java:22) [clojure-1.6.0.jar:?]
at java.lang.Thread.run(Thread.java:745) [?:1.8.0_45]