Skip to content

Throw a typed exception for tuples with unknown task or stream ids and add a deserialization strict mode - #9094

Open
L1nq0 wants to merge 3 commits into
apache:masterfrom
L1nq0:9077-typed-deserialization-exception
Open

L1nq0 wants to merge 3 commits into
apache:masterfrom
L1nq0:9077-typed-deserialization-exception

Conversation

@L1nq0

@L1nq0 L1nq0 commented Sep 16, 2026

Copy link
Copy Markdown
Contributor

Relates to #9077

This implements part of #9077: the unknown-task and unknown-stream checks in KryoTupleDeserializer now throw a typed exception, and a new config makes deserialization failures fatal again.

What changed

KryoTupleDeserializer throws TupleDeserializationException (a RuntimeException in org.apache.storm.serialization) when a tuple names a source task the receiving topology cannot resolve, or a stream id the source component does not declare. The unknown-task case previously threw a bare IllegalArgumentException; the unknown-stream case previously resolved to a null stream name and the tuple was delivered anyway.

DeserializingConnectionCallback treats the new exception as a tolerated deserialization failure, so both cases are dropped and counted like the other decode failures, with the task or stream id in the message.

The new config topology.tuple.deserialization.strict.enable (default false) makes any deserialization failure propagate instead of being dropped, restoring the pre-3.1.0 fail-fast behavior. The default keeps the behavior introduced by #9076.

Behavior changes

  • The exception type for an unresolvable source task changes from IllegalArgumentException to TupleDeserializationException. The message text is unchanged.
  • A tuple whose stream id does not resolve for its source component is now dropped with an error naming the component and the stream id. Previously it was delivered with a null stream name.

Tests

DeserializingConnectionCallbackTest gains a test for the unknown-stream drop and one for strict mode making a truncated payload fatal. The existing unknown-task test now asserts the typed exception and that the message names the task id. The storm-client test suite passes.

Open questions

Whether the other tolerated exception types in DeserializingConnectionCallback should also be replaced by typed exceptions is left open; this change only types the failures KryoTupleDeserializer itself detects. The config flag name follows the existing topology.*.enable convention; happy to rename if maintainers prefer another form.

@rzo1

rzo1 commented Sep 16, 2026

Copy link
Copy Markdown
Contributor

The unknown-stream check is correct and closes a real hole: a null stream name was previously passed to TupleImpl and the tuple delivered anyway. Sender and receiver build IdDictionary from the same topology and KryoTupleSerializer would fail on the sending side before emitting an unresolvable id, so a null on the receiving side does mean a corrupt or mismatched frame. No concerns there.

Three things before this can go in.

  1. TupleDeserializationException extends RuntimeException. The unknown-task case threw IllegalArgumentException in 3.1.0, so anything catching that stops matching. Please extend IllegalArgumentException instead. Existing callers keep working, and the entry you added to TOLERATED_DESERIALIZATION_FAILURES becomes redundant, though keeping it explicit is fine. While you are in there, add a (String, Throwable) constructor.

  2. The config is undocumented. topology.tuple.deserialization.strict.enable exists only in Config.java and defaults.yaml. docs/Serialization.md carries the table where topology.tuple.compression.max.decompressed.bytes is documented, and docs/Metrics.md covers deserializationFailures, which you added in Drop malformed tuple payloads instead of killing the receiving worker #9076. The flag belongs in both.

  3. The config javadoc does not say what enabling it costs. Under strict mode a single corrupt frame from a peer kills the worker, the supervisor restarts it, and the same frame kills it again. That restart loop is what Drop malformed tuple payloads instead of killing the receiving worker #9076 fixed, and users should read it in the config description rather than infer it. The "pre-3.1.0 behavior" wording is accurate, Drop malformed tuple payloads instead of killing the receiving worker #9076 is contained in v3.1.0, so keep that and add the consequence.

Minor:

  • assertTrue(thrown.getMessage().contains("id 3")) ties the test to the exact message text. The exception type plus the component name would be enough.
  • No test covers strict mode with the new exception. testStrictModeMakesFailuresFatal only exercises the truncated payload path.

On your open question: topology.tuple.deserialization.strict.enable is fine, it matches the existing keys.

@L1nq0

L1nq0 commented Sep 16, 2026

Copy link
Copy Markdown
Contributor Author

@rzo1 Thanks for the review. Commit be82785 addresses everything.

TupleDeserializationException now extends IllegalArgumentException and has a (String, Throwable) constructor, so handlers written against the 3.1.0 behavior keep matching. The entry in TOLERATED_DESERIALIZATION_FAILURES stays; it is redundant now but it documents where the exception comes from.

The config is documented in the table in docs/Serialization.md and next to deserializationFailures in docs/Metrics.md. The javadoc keeps the pre-3.1.0 wording and now states the consequence: a single corrupt frame kills the worker and the supervisor restarts it into the same failure, so a persistent bad frame results in a restart loop.

On the test points: the unknown-stream case no longer pins the message text, it checks the exception type and the component name, and testStrictModeMakesUnknownTaskFailureFatal covers the strict mode with the typed exception thrown for an unknown task id.

The config name stays topology.tuple.deserialization.strict.enable, per your confirmation.

@rzo1
rzo1 requested a review from GGraziadei September 18, 2026 17:14
@rzo1 rzo1 added this to the 3.3.0 milestone Sep 18, 2026
@rzo1
rzo1 requested a review from reiabreu September 18, 2026 17:14
@GGraziadei GGraziadei modified the milestones: 3.3.0, 3.2.0 Sep 20, 2026
@reiabreu

Copy link
Copy Markdown
Contributor

Currently traveling and unable to properly review this. Happy to do it once I'm able to

@GGraziadei

Copy link
Copy Markdown
Member

Review asap.

@GGraziadei GGraziadei left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good proposal! Some changes required. Thanks.

Comment on lines +113 to 114
if (strictMode || !isToleratedDeserializationFailure(e)) {
throw e;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

When a compressed frame is broken or too large, it throws an error. In strict mode, we expect this error to stop the worker, but Netty catches it, logs it, closes the connection, and keeps the worker running. This means one bad frame can cause all valid tuples in the same batch to be lost. The peer must reconnect, and no failure metric is increased.

This is confusing because the documentation says that strict mode should terminate the worker. We should either change the behavior or update the documentation.

wdyt?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Confirmed, and thanks for tracing it to the Netty layer.

The mechanism, as far as I can tell: StormServerHandler.exceptionCaught hands the cause to Utils.handleUncaughtException with ALLOWED_EXCEPTIONS containing only IOException, and the check walks the cause chain. A broken or oversized compressed frame fails inside decompression as an IOException, which ZstdUtils.decompress and KryoTupleDeserializer.deserializeTuple wrap in a RuntimeException, so under strict mode the rethrown failure still matches the IOException exemption. The connection is closed, the batch including valid tuples is discarded, and the worker keeps running. Tuple-level decode failures such as unknown task or stream ids carry no IOException in the chain, so they reach the Error path and the worker exits.

The documentation in this PR now describes both outcomes instead of claiming the worker always terminates: the Config javadoc, the configuration table in docs/Serialization.md, and the deserializationFailures paragraph in docs/Metrics.md.

The behavior itself is filed as #9157: whether strict mode should terminate the worker uniformly, or whether the connection-close outcome is the intended boundary for failures the messaging layer classifies as transport-level. The missing deserializationFailures increment on the strict path belongs there too.

/**
* Thrown when a serialized tuple names a source task or stream that the receiving topology cannot resolve.
*/
public class TupleDeserializationException extends IllegalArgumentException {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The new entry in the tolerated-failures list is not needed because TupleDeserializationException already extends IllegalArgumentException, which is already in the list.
The list is here https://github.com/apache/storm/pull/9094/changes#diff-e9421c851286d22152559b22fc546c160cf2ed1dce13f719783617b5f73e19dfR49
If the goal is to treat TupleDeserializationException separately, it should extend RuntimeException instead. Otherwise, the PR description and the comment should be updated.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Dropped in c3d3747. The comment above the set stays as the pointer to where TupleDeserializationException comes from; the type remains tolerated through its IllegalArgumentException supertype.

L1nq0 added 2 commits October 1, 2026 13:16
…d document the strict mode

The exception thrown for unknown task or stream ids now extends
IllegalArgumentException, so handlers written against 3.1.0 keep matching,
and it gains a (String, Throwable) constructor. The strict mode flag is
documented in Serialization.md and Metrics.md, and the config javadoc spells
out that a persistent bad frame puts the worker in a restart loop. The
unknown-stream test no longer pins the exact message text and the strict
mode is now covered for the typed exception as well.
@rzo1
rzo1 force-pushed the 9077-typed-deserialization-exception branch from be82785 to 1e4e849 Compare October 1, 2026 11:16

@reiabreu reiabreu left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This review was generated with the help of an LLM (Claude) and reviewed by me before posting.

Small, careful PR and the test coverage is good. No blocking issues — a couple of cleanups with suggestions inline, plus two points worth your call (the typed exception not yet being branched on, and an unknown-component edge the new stream null-check doesn't cover). Details inline.

// the fatal handling in StormServerHandler. TupleDeserializationException is thrown by KryoTupleDeserializer
// for unknown task or stream ids.
private static final Set<Class<?>> TOLERATED_DESERIALIZATION_FAILURES = new HashSet<>(Arrays.asList(
TupleDeserializationException.class,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

TupleDeserializationException extends IllegalArgumentException, and IllegalArgumentException.class is already in this set below, so isToleratedDeserializationFailure() already matches it via the supertype — this explicit entry never changes the result. Suggest dropping it (the explanatory comment just above can stay):

Suggested change
TupleDeserializationException.class,

Comment on lines +22 to +25

public TupleDeserializationException(String message, Throwable cause) {
super(message, cause);
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The (String, Throwable) constructor has no caller — both throw sites in KryoTupleDeserializer use the single-arg form. Suggest dropping it until a cause-carrying throw site exists:

Suggested change
public TupleDeserializationException(String message, Throwable cause) {
super(message, cause);
}
}

throw new TupleDeserializationException("Received a tuple from unknown task " + taskId);
}
String streamName = ids.getStreamName(componentName, streamId);
if (streamName == null) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nice — this closes the old "deliver with a null stream name" path. One asymmetry worth noting: this guards the case where the component exists but the stream id is unknown. If componentName itself isn't a key in IdDictionary.streamIdToName, getStreamName() does streamIdToName.get(componentName).get(streamId) and NPEs on the inner call before reaching this check. NPE isn't in TOLERATED_DESERIALIZATION_FAILURES, so that case stays fatal even in non-strict mode — the opposite of the drop-and-count behavior here. It shouldn't be reachable today (component comes from the same topology that populated the dictionary), so this is more of a robustness note than a live bug — but if you want the two unresolved-routing cases to behave consistently, worth a guard.

/**
* Thrown when a serialized tuple names a source task or stream that the receiving topology cannot resolve.
*/
public class TupleDeserializationException extends IllegalArgumentException {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Design question (not blocking): nothing currently branches on this type. In default mode it's tolerated via its IllegalArgumentException supertype exactly as the old bare IAE was; in strict mode every exception is fatal regardless of type. So the actual behavior change (unknown stream id now dropped instead of delivered) comes from the new throw, not from the new type — a bare IllegalArgumentException would behave identically today. Is the plan to branch on it later (per #9077)? If so, fine as a forward step; if not, it's currently just a message-bearing marker.

out.writeInt(1, true); // default stream id
byte[] unknownTask = out.toBytes();
out.writeInt(SOURCE_TASK_ID, true); // source task that exists in the topology
out.writeInt(3, true); // stream id the source component does not declare

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Minor test-robustness: 3 is "unknown" only while TestWordSpout declares fewer than 3 output streams (IdDictionary.idify() assigns ids 1..n to the sorted declared streams). If that spout ever gains a third stream, id 3 becomes valid, no exception is thrown, and this test fails confusingly. Deriving an id known to be out of range from the component's declared stream count would be sturdier than hardcoding 3.

@reiabreu

reiabreu commented Oct 2, 2026

Copy link
Copy Markdown
Contributor

@L1nq0 ran this PR through Claude. Made a few suggestions. Let us know what you think of them.
Happy to approve this once they are resolved (plus the ones highlighted by @GGraziadei )

…trict-mode boundary for IOException-wrapped failures

The tolerated-failures set matches TupleDeserializationException through its
IllegalArgumentException supertype, so the explicit entry and the unused
two-argument constructor are removed. IdDictionary.getStreamName returns null
when the component is unknown, so both unresolved-routing cases reach the typed
exception instead of a NullPointerException that stays fatal in default mode.
The strict-mode documentation now separates the two outcomes: decode failures
exit the worker, while failures wrapping an IOException close the connection and
drop the batch with the worker kept running. The unknown-stream test derives its
id from the declared stream count instead of a hardcoded constant.
@L1nq0

L1nq0 commented Oct 6, 2026

Copy link
Copy Markdown
Contributor Author

@reiabreu Thanks for the review. Commit c3d3747, on top of the rebased branch, addresses the points.

The redundant TupleDeserializationException entry is dropped from the tolerated set; the comment above the list stays. This also covers the second point from GGraziadei's review.

The two-argument constructor is removed. It was added at rzo1's request in the first review round; if you would rather keep it, say so and it goes back in.

IdDictionary.getStreamName now returns null for a component that is not in the dictionary, so both unresolved-routing cases reach the typed exception. A test covers the unknown-component lookup. Previously the inner map lookup threw a NullPointerException, which is not in the tolerated set and stayed fatal even in default mode.

On the design question: yes, the type is a forward step for #9077. In default mode it is tolerated through its IllegalArgumentException supertype exactly like the bare IAE it replaces; the behavioral change in this PR comes from the new throw for unknown stream ids. Branching on the type is the follow-up work.

The unknown-stream test now derives the id one past the component's declared stream count instead of the hardcoded 3.

The points GGraziadei raised are answered on their threads. The strict-mode documentation now separates the two outcomes, and the behavior question is filed as #9157.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants