Repository navigation
Throw a typed exception for tuples with unknown task or stream ids and add a deserialization strict mode #9094
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: master
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -79,9 +79,12 @@ private TupleImpl deserializeTuple(byte[] data) { | |
| int streamId = kryoInput.readInt(true); | ||
| String componentName = context.getComponentId(taskId); | ||
| if (componentName == null) { | ||
| throw new IllegalArgumentException("Received a tuple from unknown task " + taskId); | ||
| throw new TupleDeserializationException("Received a tuple from unknown task " + taskId); | ||
| } | ||
| String streamName = ids.getStreamName(componentName, streamId); | ||
| if (streamName == null) { | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 |
||
| throw new TupleDeserializationException("Component " + componentName + " has no stream with id " + streamId); | ||
| } | ||
| MessageId id = MessageId.deserialize(kryoInput); | ||
| List<Object> values = kryo.deserializeFrom(kryoInput); | ||
| return new TupleImpl(context, values, componentName, taskId, streamName, id); | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,22 @@ | ||
| /** | ||
| * Licensed to the Apache Software Foundation (ASF) under one or more contributor license agreements. See the NOTICE file distributed with | ||
| * this work for additional information regarding copyright ownership. The ASF licenses this file to you under the Apache License, Version | ||
| * 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at | ||
| * | ||
| * http://www.apache.org/licenses/LICENSE-2.0 | ||
| * | ||
| * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS, | ||
| * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions | ||
| * and limitations under the License. | ||
| */ | ||
|
|
||
| package org.apache.storm.serialization; | ||
|
|
||
| /** | ||
| * Thrown when a serialized tuple names a source task or stream that the receiving topology cannot resolve. | ||
| */ | ||
| public class TupleDeserializationException extends IllegalArgumentException { | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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.
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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.
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 |
||
| public TupleDeserializationException(String message) { | ||
| super(message); | ||
| } | ||
| } | ||
There was a problem hiding this comment.
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?
There was a problem hiding this comment.
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.