This is an automated email from the ASF dual-hosted git repository.
rzo1 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/storm.git
The following commit(s) were added to refs/heads/master by this push:
new 42147bda4 Log at error level when the kafka spout gives up on a tuple;
document max-retry give-up semantics (#9093)
42147bda4 is described below
commit 42147bda41ad65d705656122c1c47230e905b28c
Author: L1nq <[email protected]>
AuthorDate: Sat Oct 3 00:56:53 2026 +0800
Log at error level when the kafka spout gives up on a tuple; document
max-retry give-up semantics (#9093)
* Log at error level when the kafka spout gives up on a tuple; document
max-retry give-up semantics
* Reword the give-up error to hold for any retry service and pin the
max-retry ack ordering
KafkaSpoutRetryService#schedule may decline a message for reasons other than
reaching the retry limit, so the give-up log no longer names the reason.
KafkaSpoutRetryLimitTest now verifies that onMaxRetryReached runs before the
tuple is acked, matching the documented contract.
---
.../src/main/java/org/apache/storm/kafka/spout/KafkaSpout.java | 3 ++-
.../storm/kafka/spout/KafkaSpoutRetryExponentialBackoff.java | 6 +++++-
.../java/org/apache/storm/kafka/spout/KafkaTupleListener.java | 4 +++-
.../org/apache/storm/kafka/spout/KafkaSpoutRetryLimitTest.java | 9 +++++++++
4 files changed, 19 insertions(+), 3 deletions(-)
diff --git
a/external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/KafkaSpout.java
b/external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/KafkaSpout.java
index 73582dd4c..b4dca55af 100644
---
a/external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/KafkaSpout.java
+++
b/external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/KafkaSpout.java
@@ -605,7 +605,8 @@ public class KafkaSpout<K, V> extends BaseRichSpout {
msgId.incrementNumFails();
if (!retryService.schedule(msgId)) {
- LOG.debug("Reached maximum number of retries. Message [{}] being
marked as acked.", msgId);
+ LOG.error("The retry service will not retry message [{}]: the
tuple will be acked and its offset and later "
+ + "offsets may be committed even though the record was not
processed.", msgId);
// this tuple should be removed from emitted only inside the ack()
method. This is to ensure
// that the OffsetManager for that TopicPartition is updated and
allows commit progression
tupleListener.onMaxRetryReached(msgId);
diff --git
a/external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/KafkaSpoutRetryExponentialBackoff.java
b/external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/KafkaSpoutRetryExponentialBackoff.java
index 8a2f54303..76b7e3921 100644
---
a/external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/KafkaSpoutRetryExponentialBackoff.java
+++
b/external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/KafkaSpoutRetryExponentialBackoff.java
@@ -156,7 +156,11 @@ public class KafkaSpoutRetryExponentialBackoff implements
KafkaSpoutRetryService
* where failCount = 1, 2, 3, ... nextRetry = Min(nextRetry, currentTime +
maxDelay).
* <p/>
* By specifying a value for maxRetries lower than Integer.MAX_VALUE, the
user decides to sacrifice guarantee of delivery for the
- * previous polled records in favor of processing more records.
+ * previous polled records in favor of processing more records. Setting a
finite limit also stops the spout from retrying forever
+ * tuples that fail every time they are emitted, e.g. tuples the receiving
worker drops because they cannot be deserialized:
+ * once the limit is reached, the tuple is acked and offsets past it can
be committed. These tuples are reported to
+ * {@link KafkaTupleListener#onMaxRetryReached(KafkaSpoutMessageId)},
which receives the topic, partition and offset needed
+ * to fetch the record from Kafka again.
*
* @param initialDelay initial delay of the first retry
* @param delayPeriod the time interval that is the ratio of the
exponential backoff formula (geometric progression)
diff --git
a/external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/KafkaTupleListener.java
b/external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/KafkaTupleListener.java
index 100148312..e9e5cb7a7 100644
---
a/external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/KafkaTupleListener.java
+++
b/external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/KafkaTupleListener.java
@@ -74,7 +74,9 @@ public interface KafkaTupleListener extends Serializable {
void onRetry(KafkaSpoutMessageId msgId);
/**
- * Called when the maximum number of retries have been reached.
+ * Called when the maximum number of retries have been reached. The tuple
is acked right after this callback and
+ * commits can then move past its offset; this is the last point at which
the record can be retained. The msgId
+ * identifies the record by topic, partition and offset.
*
* @param msgId The id of the tuple in the spout.
*/
diff --git
a/external/storm-kafka-client/src/test/java/org/apache/storm/kafka/spout/KafkaSpoutRetryLimitTest.java
b/external/storm-kafka-client/src/test/java/org/apache/storm/kafka/spout/KafkaSpoutRetryLimitTest.java
index 381796b73..6d4a9e803 100644
---
a/external/storm-kafka-client/src/test/java/org/apache/storm/kafka/spout/KafkaSpoutRetryLimitTest.java
+++
b/external/storm-kafka-client/src/test/java/org/apache/storm/kafka/spout/KafkaSpoutRetryLimitTest.java
@@ -54,6 +54,7 @@ public class KafkaSpoutRetryLimitTest {
private final long offsetCommitPeriodMs = 2_000;
private final TopologyContext contextMock = mock(TopologyContext.class);
private final SpoutOutputCollector collectorMock =
mock(SpoutOutputCollector.class);
+ private final KafkaTupleListener tupleListener =
mock(KafkaTupleListener.class);
private final Map<String, Object> conf = new HashMap<>();
private final TopicPartition partition = new
TopicPartition(SingleTopicKafkaSpoutConfiguration.TOPIC, 1);
@Mock
@@ -72,6 +73,7 @@ public class KafkaSpoutRetryLimitTest {
spoutConfig = createKafkaSpoutConfigBuilder(mock(TopicFilter.class),
mock(ManualPartitioner.class), -1)
.setOffsetCommitPeriodMs(offsetCommitPeriodMs)
.setRetry(ZERO_RETRIES_RETRY_SERVICE)
+ .setTupleListener(tupleListener)
.build();
}
@@ -99,6 +101,13 @@ public class KafkaSpoutRetryLimitTest {
spout.fail(messageId);
}
+ // The give-up contract: onMaxRetryReached runs before the tuple
is acked.
+ InOrder giveUpOrder = inOrder(tupleListener);
+ for (KafkaSpoutMessageId messageId : messageIds.getAllValues()) {
+ giveUpOrder.verify(tupleListener).onMaxRetryReached(messageId);
+ giveUpOrder.verify(tupleListener).onAck(messageId);
+ }
+
// Advance time and then trigger call to kafka consumer commit
Time.advanceTime(KafkaSpout.TIMER_DELAY_MS + offsetCommitPeriodMs);
spout.nextTuple();