Repository navigation
[STORM-2359] Revising Message Timeouts #6141
Description
Activity
Subtask of parent task STORM-2284
srdo:
This is a great idea, we've also seen topologies exhibit bad performance when under pressure because a lot of the queued tuples are already considered timed out by the spout.
It seems like part of the reason the tuple timeout is currently a set value that is not automatically reset is to avoid flooding the spout instances with ack messages. The ackers can be scaled out to handle a large number of acks, and can then notify the spout of acks only once per tuple tree. I'm assuming we want to keep the management of tuple tree ack/fail inside the spout( ? ), so have you considered a way to deduplicate or "bundle up" the ack stream in the acker before notifying the spout that the timeout should be reset? Maybe only notify once a percentage of the timeout has elapsed?
The bolt heartbeating may be able to reuse some code from https://issues.apache.org/jira/browse/STORM-1549
srdo:
I thought a bit more about this.
Here are the situations I could think of where tuples are currently being expired without being lost:
- The tuple tree is still making progress, and tuples are getting acked, but the time to process the entire tree is larger than the tuple timeout. This is very likely to happen if there's congestion somewhere in the topology.
- The tuple tree is not making progress, because a bolt is currently processing a tuple in the tree, and processing is taking longer than expected.
- The tuple tree is not making progress, because the tuple(s) are stuck in queues behind other slow tuples. This is also very likely to happen if there's congestion in the topology.
The situation where there's still progress being made can be solved by resetting the tuple timeout whenever an ack is received. In order to reduce load on the spout, we should try to "bundle up" these resets in the acker bolts before sending them to the spout. I think a decent way to do this bundling is to make the acker bolt keep track of which tuples they've received acks from since the last time timeouts were reset. When a configured interval expires, the bolt empties out the list of live tuples, and sends timeout resets for all of them to the spout. The interval should probably be specified as a percentage of the tuple timeout.
If a bolt is taking longer to process a tuple than expected, it can be solved in the concrete bolt implementation by using OutputCollector.resetTimeout at an appropriate interval (e.g. the tuple timeout minus a few seconds).
When tuples are stuck in queues behind other tuples, the topology can have a hard time recovering. This is because the expiration timer starts ticking for a tuple as soon as it's emitted, so if the bolt queues are congested, the bolts may be spending all their time processing tuples that belong to expired tuple trees. In order to solve this, we need to reset timeouts for queued tuples from time to time. It should be possible to add a thread that peeks at the available messages in the DisruptorQueue with some interval, and resets the timeout for any messages that were also queued last time the thread was run. Only sending the resets once a tuple has been queued for the entire interval should help decrease the number of unnecessary resets sent to the spout. We should be able to reuse the interval configuration also added to the acker bolt.
sorry just noticed your comments.... will go though it and respond next week.
In order to reduce load on the spout, we should try to "bundle up" these resets in the acker bolts before sending them to the spout
The resets are managed internally by the ACKer bolt. The spout only gets notified if the timeout expires or if tuple-tree is fully processed.
When tuples are stuck in queues behind other tuples
That case would be more accurately classified as "progress is being made" ... but slower than expected.
The case of 'progress is not being made' is when a worker that is processing part of the tuple tree dies.If a bolt is taking longer to process a tuple than expected, it can be solved in the concrete bolt implementation by using OutputCollector.resetTimeout at an appropriate interval (e.g. the tuple timeout minus a few seconds).
I think you are trying to mitigate the number of resets being sent to Spout... but like i mentioned before reset are never sent to spout.
srdo:
The resets are managed internally by the ACKer bolt. The spout only gets notified if the timeout expires or if tuple-tree is fully processed.
How will this work? The current implementation has a pending map in both the acker and the spout, which rotate every topology.message.timeout.secs. If the acker doesn't forward reset requests to the spout, the spout will just expire the tuple tree on its own when the message timeout has passed.
That case would be more accurately classified as "progress is being made" ... but slower than expected.
The case of 'progress is not being made' is when a worker that is processing part of the tuple tree dies.Yes, you are right. But it is currently possible that the topology may degrade to no progress being made even if each individual tuple could be processed under the message timeout, because tuples can expire while queued and get reemitted, where they can then be delayed by their own duplicates which are ahead in the queue. For IRichBolts, this can be mitigated by the bolt being written to accept and queue tuples internally, where the bolt can then reset their timeouts manually if necessary, but for IBasicBolt this is not possible.
Just to give a concrete example, we had an IBasicBolt enrich tuples with some database data. Most tuples were processed very quickly, but a few were slow. Even the slow tuples never took longer than our message timeout individually. We then had an instance where a bunch of slow tuples happened to come in on the stream close to each other. The first few were processed before they expired, but the rest expired while queued. The spout then reemitted the expired tuples, and they got into the queue behind their own expired instances. Since the bolt won't skip expired tuples, the freshly emitted tuples also expired, which caused another reemit. This repeated until the topology was restarted so the queues could be cleared.
The current implementation has a pending map in both the acker and the spout, which rotate every topology.message.timeout.secs.
Need to see if we can eliminate the timeout logic from the spout and have it only the ACKer (i can think of some issues). If we must retain that logic in the spouts, the timeout value that it operates on (full tuple tree processing) would have to be separated from the timeout value that the ACKER uses to track progress between stages.
The spout then reemitted the expired tuples, and they got into the queue behind their own expired instances.
Perfect example indeed. The motivation of this jira is to try to eliminate/mitigate triggering of timeouts for queued/inflight tuples that are not lost. The only time we need timeouts/remits to be triggered is when one/more tuples in the tuple tree are truly lost. I think that can only happen if a worker/bolt/spout died. So the case your describing should not happen if we solve this problem correctly.
IMO, the ideal solution would have the spouts remit only the specific tuples whose tuple trees had some loss due to a worker going down. I am not yet certain whether/not this initial idea described in the doc is the optimal solution. Perhaps a better way is to trigger such re-emits only if a worker/bolt/spout went down.
srdo:
I think the reason it would be tough to move timeouts entirely to the ackers is that we'd need to figure out how to deal with message loss between the spout and acker when the acker sends a timeout message to the spout. The current implementation errs on the side of caution by always reemitting if it can't positively say that a tuple has been acked. I'm not sure how we could do the same when the acker has to notify the spout to reemit after the timeout, because that message could be lost.
It might be a good idea as you mention to instead have two timeouts, a short one for the acker and a much longer one for the spout. It would probably mean that messages where the acker init message is lost will take much longer to retry than messages that are lost elsewhere, but it might allow us to keep timeout resets out of the spout.
Tuples can be lost if a worker died, but what if there's a network issue? Can't messages also be lost then?
srdo:
I've been taking a look at this, and have a proposal for how we could move the timeout entirely to the acker.
I'm going to assume in the following that tuples sent from one task to another are received in the order they were sent, ignoring message loss. I think this is the case, but please correct me if it isn't.
The problem
Both the acker and spout executor currently have rotating pending maps, which time out tuples based on received ticks. The acker pending map just discards the tuple, while the spout pending map fails them if they get rotated out. The reason this currently happens in both acker and spout is that we need to ensure that the spout fails tuples if they time out, even in the presence of message loss.
If we were to move timeouts entirely to the acker, there would be a potential for "lost" tuples (ones that end up not being reemitted) if messages between the acker and spout get lost.
- The spout may emit a new tuple, and try to notify the acker about it. If this message is lost, the acker doesn't know about the tuple, and thus can't time it out.
- The acker may send an ack or fail message back to the spout, which may get lost. Since the acker currently deletes acked/failed messages as soon as they are acked/failed, this would prevent the message from being replayed.
Suggested solution
We can move the timeout logic to the acker, and it will work out of the box as long as there are no lost messages. I think we can account for lost messages by having the spout executor periodically update its view of pending tuples based on the state in the acker.
Say the spout pending root ids are A, and the acker pending root ids are B. The spout can periodically (e.g. once per second) send to each acker the root ids in A. The acker should respond to this tuple by sending back A - B (it can optionally also delete anything in B that is not in A). The spout can safely fail any tuple in A - B which is also in the spout pending when the sync tuple response is received.
- If a tuple is in A and B, it is still pending, and shouldn't be removed. Returning only A - B ensures pending tuples remain.
- If a tuple is in A but not B, the ack_init message was lost, the acker may have crashed and restarted or the tuple has simply been acked or failed.
- If the ack_init was lost, or the acker crashed, then the spout should replay the message. Since A - B contains the tuple, the spout will fail it when it receives the response.
- If the tuple was acked, then the acker is guaranteed to have sent the ack message on to the spout before handling the sync tuple. Due to message ordering, the ack will be processed before the sync tuple response, making the presence of the tuple in A - B irrelevant.
- If a tuple is not in A but in B, then the spout may have crashed. The acker can optionally just discard the pending tuple without notifying the spout, since notifying the spout about the state of a tuple emitted by a different instance is a no op.
Benefit
Moving the timeout logic to the acker makes it much cheaper to reset tuple timeouts, since the spout no longer needs to be notified directly. We could likely make the acker reset the timeout automatically any time it receives an ack.
Depending on overhead, we might be able to add an option to let Storm reset timeouts for tuples that are still being actively processed (i.e. received by a bolt but not yet acked/failed). This would be beneficial to avoid the event tsunami problem described in the linked design doc. It could help prevent the type of degenerate case described at https://issues.apache.org/jira/browse/STORM-2359?focusedCommentId=16043409&page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel#comment-16043409.
Cost
There is a bit of extra overhead to doing this. The spout needs to keep track of which acker tasks are responsible for which root ids. There is also the (fairly minor) overhead of sending the sync tuples back and forth. In terms of processing reliability, a tuple the acker considers failed/timed out will be failed in the spout once one of the sync tuples make it through.
srdo:
There's a branch implementing this at https://github.com/srdo/storm/tree/STORM-2359-experimentation
I did a couple of runs of the ConstSpoutNullBoltTopology from storm-perf, as well as the ThroughputVsLatency topology.
I'm seeing an overhead of about 10% for the ConstSpoutNullBoltTopology, since the spout and ackers are the only components involved. For the ThroughputVsLatency topology, the overhead looks to be negligible. I've attached the raw numbers and some charts.
The tests were run locally, with the default configuration for Storm as well as the default configuration for ConstSpoutNullBoltTopology. For TVL I reused the setup from the Netty benchmarks: ./storm jar .\storm-loadgen-2.0.1-SNAPSHOT.jar org.apache.storm.loadgen.ThroughputVsLatency --rate 90000 --spouts 4 --splitters 4 --counters 4 --reporter 'tsv:test.csv' -c topology.workers=4 --test-time 5
srdo:
I will check what impact it has to reinsert tuples in the acker's pending when acks are received.
Edit: It seems to make no difference to throughput for ConstSpoutNullBoltTopology, so I'll just roll it in to the current branch.
If we were to move timeouts entirely to the acker, there would be a potential for "lost" tuples (ones that end up not being reemitted) if messages between the acker and spout get lost.
- The spout may emit a new tuple, and try to notify the acker about it. If this message is lost, the acker doesn't know about the tuple, and thus can't time it out.
Acker will be totally clueless if all acks from downstream bolts are also lost for the same tuple tree.
We can move the timeout logic to the acker,
Alternatives :
1) Handle timeouts in SpoutExecutor. It sends a timeout msg to ACKer. Benefit: May not be any better than doing in ACKer. Do you have any thoughts about this ?2) Eliminate timeout communication between Spout & ACK. Let each does its own timeout. Benefit: Eliminates need for:
- Periodic sync up msgs (also avoid possibility of latency spikes in multi worker mode when these msgs get large).
- A - B calculation, as well as
- timeout msg exchanges.
TimeoutReset msgs can still be supported.
srdo:
Acker will be totally clueless if all acks from downstream bolts are also lost for the same tuple tree.
Yes. This is why the sync tuple is needed. If the spout executor doesn't time out the tree on its own, then we need a mechanism to deal with trees where the acker isn't aware of the tree due to message loss of the init and ack/fail messages. If we don't do a sync, the spout will never know that the acker isn't aware of the tree, and the tree will never fail.
1) Handle timeouts in SpoutExecutor. It sends a timeout msg to ACKer. Benefit: May not be any better than doing in ACKer. Do you have any thoughts about this ?
I don't think there's any benefit to doing this over what we're doing now. Currently the spout executor times out the tuple, and the acker doesn't try to fail tuples on timeout. Instead it just quietly discards whatever information it has about a tree once the pending map rotates a few times. I'm not sure what we'd gain from the acker not rotating pending and relying on a timeout tuple instead. The benefit as I see it of moving the timeouts to the acker would be the ability to reset timeouts more frequently (e.g. on every ack) without increasing load on the spout executor, which we can't do if the spout is still handling the timeout.
2) Eliminate timeout communication between Spout & ACK. Let each does its own timeout.
This is how it works now. As you note, there are some benefits to doing it this way. An additional benefit is that we can easily reason about the max time to fail a tuple on timeout, since the tuple will fail as soon as the spout rotates it out of pending. The drawbacks are:
- Any time the timeout needs to be reset, a message needs to go to both the acker and the spout (current path is bolt -> acker -> spout)
- Since resetting timeouts is reasonably expensive, we don't do it as part of regular acks, the user has to manually call collector.resetTimeout().
The benefit of moving timeouts entirely to the acker is that we can reset timeouts automatically on ack. This means that the tuple timeout becomes somewhat easier to work with, since tuples won't time out as long as an edge in the tree is getting acked occasionally.
This is currently more of a loose idea, but in the long run I'd like to try adding a feature toggle for more aggressive automatic timeout resets. I'm not sure what the performance penalty would be, but it would be nice if the bolts could periodically reset the timeout for all queued/in progress tuples (tuples received by the worker, but not yet acked/failed by a collector). Doing this would eliminate the degenerate case you mention in the design doc, where some bolt is taking slightly too long to process a tuple, and the queued tuples time out, which causes the spout to reemit, which puts the reemits in queue behind the already timed out tuples. In some cases this can cause the topology to spend all its time processing timed out tuples, preventing it from making progress, even though each individual tuple could have been processed within the timeout.
If this idea turns out to be possible to implement without adding a large overhead, it would add an extra drawback to timeouts in the spout:
- We can't do automatic timeout resetting, since flooding the spout with reset messages is a bad idea. In particular we can't reset timeouts for messages that are sitting in bolt queues. The user can't reset these timeouts manually either.
srdo:
I made a few more changes, mainly to remove the pending map from the spout. It looks like there is no real difference in performance between timeouts in the acker and timeouts in the spout now. I've updated the document with benchmarks, please take a look.
I see your point about having the timeout logic in acker. It’s a good one!.
Sorry for delay in my response. I am traveling so unable to look into code
Could u elaborate on the contents of the message being exhanged for Comoutung A-B?
10 remaining items
Can we revive this? @rzo1
If you feel like implementing it, feel free.
Reacted by Gowtham SIf you feel like implementing it, feel free.
Thanks Richard — I'll take this up. Planning to implement topology.message.progress.timeout.secs using the lastProgressTime approach (additive field on AckObject, no protocol changes, off by default). Will share a draft PR once I have the core logic and tests in place.
- addedjavaPull requests that update Java codePull requests that update Java code
on Jun 11, 2026 Following up — I've implemented the progress-based timeout approach and
opened a PR:Richard, rzo1 — would appreciate your review when you have a moment.
- added 2 commits that reference this issue
on Jun 11, 2026 - Happy to review it as soon as possible. Cheers…On Thu, 11 Jun 2026 at 12:22, Gowtham S ***@***.***> wrote: *Gowtham-Gowts* left a comment (apache/storm#6141) <#6141 (comment)> Following up — I've implemented the progress-based timeout approach and opened a PR: #8787 <#8787> Richard, rzo1 — would appreciate your review when you have a moment. — Reply to this email directly, view it on GitHub <#6141?email_source=notifications&email_token=AAG5GIWTAMUAUUFNTUHGH2L47KJANA5CNFSNUABFM5UWIORPF5TWS5BNNB2WEL2JONZXKZKDN5WW2ZLOOQXTINRXHE4TSOJWGE2KM4TFMFZW63VKON2WE43DOJUWEZLEUVSXMZLOOSWGM33PORSXEX3DNRUWG2Y#issuecomment-4679999614>, or unsubscribe <https://github.com/notifications/unsubscribe-auth/AAG5GIQ7E234YUOET3PEWJL47KJANAVCNFSNUABEKJSXA33TNF2G64TZHMYTIMJTGU2DOMB3JFZXG5LFHMZDQMBZHEYDCOBQHGQXMAQ> . You are receiving this because you are subscribed to this thread.Message ID: ***@***.***>
- added a commit that references this issue
on Oct 1, 2026
A revised strategy for message timeouts is proposed here.
Design Doc:
https://docs.google.com/document/d/1am1kO7Wmf17U_Vz5_uyBB2OuSsc4TZQWRvbRhX52n5w/edit?usp=sharing
Originally reported by roshan_naik, imported from: Revising Message Timeouts