reiabreu commented on code in PR #9094:
URL: https://github.com/apache/storm/pull/9094#discussion_r4167102414
##########
storm-client/src/jvm/org/apache/storm/messaging/DeserializingConnectionCallback.java:
##########
@@ -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,
Review Comment:
`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):
```suggestion
```
##########
storm-client/src/jvm/org/apache/storm/serialization/TupleDeserializationException.java:
##########
@@ -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 {
+ public TupleDeserializationException(String message) {
+ super(message);
+ }
+
+ public TupleDeserializationException(String message, Throwable cause) {
+ super(message, cause);
+ }
Review Comment:
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:
```suggestion
}
```
##########
storm-client/src/jvm/org/apache/storm/serialization/KryoTupleDeserializer.java:
##########
@@ -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) {
Review Comment:
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.
##########
storm-client/test/jvm/org/apache/storm/messaging/DeserializingConnectionCallbackTest.java:
##########
@@ -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
Review Comment:
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`.
##########
storm-client/src/jvm/org/apache/storm/serialization/TupleDeserializationException.java:
##########
@@ -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 {
Review Comment:
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.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]