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]

Reply via email to