Skip to content
Open
Show file tree
Hide file tree
Changes from 2 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions conf/defaults.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,7 @@ storm.compression.zstd.level: 3
storm.compression.zstd.max.decompressed.bytes: 104857600
storm.compression.gzip.max.decompressed.bytes: 104857600
topology.tuple.compression.max.decompressed.bytes: 10485760
topology.tuple.deserialization.strict.enable: false
storm.codedistributor.class: "org.apache.storm.codedistributor.LocalFileSystemCodeDistributor"
storm.workers.artifacts.dir: "workers-artifacts"
storm.health.check.dir: "healthchecks"
Expand Down
2 changes: 1 addition & 1 deletion docs/Metrics.md
Original file line number Diff line number Diff line change
Expand Up @@ -302,7 +302,7 @@ Be aware that the `__system` bolt is an actual bolt so regular bolt metrics desc

`dequeuedMessages` is a throwback to older code where there was an internal queue between the server and the bolts/spouts. That is no longer the case and the value can be ignored.
`enqueued` is a map between the address of the remote worker and the number of tuples that were sent from it to this worker.
`deserializationFailures` is the number of incoming messages that failed to deserialize and were dropped.
`deserializationFailures` is the number of incoming messages that failed to deserialize and were dropped. When `topology.tuple.deserialization.strict.enable` is set, deserialization failures are not dropped or counted; they propagate and terminate the worker instead.

##### Send (Netty Client)

Expand Down
1 change: 1 addition & 0 deletions docs/Serialization.md
Original file line number Diff line number Diff line change
Expand Up @@ -131,6 +131,7 @@ Be aware that the topology-wide form enables compression for *every* remote-boun
| `topology.tuple.compression.threshold` | `1460` | Minimum serialized tuple size, in bytes, before compression is attempted. Tuples at or below this size are sent uncompressed. The default matches the typical Ethernet TCP MSS, so payloads that already fit in a single network frame are never compressed. |
| `storm.compression.zstd.level` | `3` | Zstd compression level. Supported range is 1–19; levels 20–22 (ultra mode) are prohibited because of their memory requirements. |
| `topology.tuple.compression.max.decompressed.bytes` | `10485760` (10 MB) | Upper bound on the decompressed size of a single tuple. Decompression that would exceed this limit fails, guarding against malicious or corrupt payloads. |
| `topology.tuple.deserialization.strict.enable` | `false` | Makes tuple deserialization failures fatal: an incoming message that fails to decode kills the receiving worker instead of being dropped and counted. This restores the pre-3.1.0 behavior and is mainly useful for debugging, because a persistent bad frame will put the worker in a restart loop. |

#### How decompression works

Expand Down
10 changes: 10 additions & 0 deletions storm-client/src/jvm/org/apache/storm/Config.java
Original file line number Diff line number Diff line change
Expand Up @@ -1648,6 +1648,16 @@ public class Config extends HashMap<String, Object> {
*/
@IsPositiveNumber(includeZero = false)
public static final String TOPOLOGY_TUPLE_COMPRESSION_MAX_DECOMPRESSED_BYTES = "topology.tuple.compression.max.decompressed.bytes";
/**
* Topology configuration to make tuple deserialization failures fatal instead of dropping the undecodable message.
* By default a message that fails to decode on the receiving worker is dropped and counted, and the worker keeps
* running. When set to {@code true}, any deserialization failure propagates and the worker exits, restoring the
* pre-3.1.0 behavior. Be aware that a single corrupt frame from a peer then kills the worker, and the supervisor
* restarts it into the same failure, so a persistent bad frame results in a restart loop.
* Default: {@code false}.
*/
@IsBoolean
public static final String TOPOLOGY_TUPLE_DESERIALIZATION_STRICT_ENABLE = "topology.tuple.deserialization.strict.enable";
/**
* Configure the topology metrics reporters to be used on workers.
*/
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
import org.apache.storm.daemon.worker.WorkerState;
import org.apache.storm.metric.api.IMetric;
import org.apache.storm.serialization.KryoTupleDeserializer;
import org.apache.storm.serialization.TupleDeserializationException;
import org.apache.storm.task.GeneralTopologyContext;
import org.apache.storm.tuple.AddressedTuple;
import org.apache.storm.tuple.Tuple;
Expand All @@ -43,8 +44,10 @@ public class DeserializingConnectionCallback implements IConnectionCallback, IMe
private static final Logger LOG = LoggerFactory.getLogger(DeserializingConnectionCallback.class);

// A tuple that cannot be decoded is dropped instead of killing the worker; anything outside this set keeps
// the fatal handling in StormServerHandler.
// 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,

IOException.class,
KryoException.class,
IllegalArgumentException.class,
Expand All @@ -71,6 +74,8 @@ protected KryoTupleDeserializer initialValue() {
}
};

private final boolean strictMode;

// Track serialized size of messages.
private final boolean sizeMetricsEnabled;
private final ConcurrentHashMap<String, AtomicLong> byteCounts = new ConcurrentHashMap<>();
Expand All @@ -87,6 +92,7 @@ public DeserializingConnectionCallback(final Map<String, Object> conf, final Gen
this.context = context;
cb = callback;
sizeMetricsEnabled = ObjectReader.getBoolean(conf.get(Config.TOPOLOGY_SERIALIZED_MESSAGE_SIZE_METRICS), false);
strictMode = ObjectReader.getBoolean(conf.get(Config.TOPOLOGY_TUPLE_DESERIALIZATION_STRICT_ENABLE), false);

}

Expand All @@ -104,7 +110,7 @@ public void recv(List<TaskMessage> batch) {
try {
tuple = des.deserialize(message.message());
} catch (Exception e) {
if (!isToleratedDeserializationFailure(e)) {
if (strictMode || !isToleratedDeserializationFailure(e)) {
throw e;
Comment on lines +111 to 112

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.

}
deserializationFailures.incrementAndGet();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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) {

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.

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);
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
/**
* 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 {

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.

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.

public TupleDeserializationException(String message) {
super(message);
}

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);
}
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import org.apache.storm.daemon.worker.WorkerState;
import org.apache.storm.serialization.KryoTupleDeserializer;
import org.apache.storm.serialization.KryoTupleSerializer;
import org.apache.storm.serialization.TupleDeserializationException;
import org.apache.storm.task.GeneralTopologyContext;
import org.apache.storm.testing.TestWordCounter;
import org.apache.storm.testing.TestWordSpout;
Expand Down Expand Up @@ -132,14 +133,60 @@ public void testTruncatedKryoPayloadDroppedAndBatchContinues() {
@Test
public void testUnknownSourceTaskDroppedAndBatchContinues() {
Map<String, Object> conf = baseConf();

TupleDeserializationException thrown = assertThrows(TupleDeserializationException.class,
() -> new KryoTupleDeserializer(conf, context).deserialize(unknownSourceTaskTuple()));
assertTrue(thrown.getMessage().contains("9999"),
"expected the task id in the message but was: " + thrown.getMessage());

assertBatchDeliversOnlyValidMessages(conf, unknownSourceTaskTuple());
}

@Test
public void testUnknownStreamIdDroppedAndBatchContinues() {
Map<String, Object> conf = baseConf();
Output out = new Output(16, 32);
out.writeInt(9999, true); // source task that does not exist in the topology
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.

byte[] unknownStream = out.toBytes();

assertThrows(IllegalArgumentException.class, () -> new KryoTupleDeserializer(conf, context).deserialize(unknownTask));
TupleDeserializationException thrown = assertThrows(TupleDeserializationException.class,
() -> new KryoTupleDeserializer(conf, context).deserialize(unknownStream));
assertTrue(thrown.getMessage().contains(SOURCE_COMPONENT),
"expected the component name in the message but was: " + thrown.getMessage());

assertBatchDeliversOnlyValidMessages(conf, unknownTask);
assertBatchDeliversOnlyValidMessages(conf, unknownStream);
}

@Test
public void testStrictModeMakesFailuresFatal() {
Map<String, Object> conf = baseConf();
conf.put(Config.TOPOLOGY_TUPLE_DESERIALIZATION_STRICT_ENABLE, true);
byte[] full = serializedTuple(conf, new Values("a-string-long-enough-to-survive-truncation", 7));
byte[] truncated = Arrays.copyOf(full, full.length - 10);

WorkerState.ILocalTransferCallback transfer = mock(WorkerState.ILocalTransferCallback.class);
DeserializingConnectionCallback callback = new DeserializingConnectionCallback(conf, context, transfer);

assertThrows(KryoException.class, () -> callback.recv(Collections.singletonList(taskMessage(truncated))));

verify(transfer, never()).transfer(any());
assertEquals(0L, callback.getAndResetDeserializationFailures());
}

@Test
public void testStrictModeMakesUnknownTaskFailureFatal() {
Map<String, Object> conf = baseConf();
conf.put(Config.TOPOLOGY_TUPLE_DESERIALIZATION_STRICT_ENABLE, true);

WorkerState.ILocalTransferCallback transfer = mock(WorkerState.ILocalTransferCallback.class);
DeserializingConnectionCallback callback = new DeserializingConnectionCallback(conf, context, transfer);

assertThrows(TupleDeserializationException.class,
() -> callback.recv(Collections.singletonList(taskMessage(unknownSourceTaskTuple()))));

verify(transfer, never()).transfer(any());
assertEquals(0L, callback.getAndResetDeserializationFailures());
}

@Test
Expand Down Expand Up @@ -269,6 +316,13 @@ private void assertBatchDeliversOnlyValidMessages(Map<String, Object> conf, byte
assertNull(callback.getValueAndReset());
}

private static byte[] unknownSourceTaskTuple() {
Output out = new Output(16, 32);
out.writeInt(9999, true); // source task that does not exist in the topology
out.writeInt(1, true); // default stream id
return out.toBytes();
}

private Map<String, Object> baseConf() {
Map<String, Object> conf = new HashMap<>(Utils.readStormConfig());
return conf;
Expand Down
Loading