diff --git a/cdc-service/src/main/java/com/xtrmetl/cdc/config/KafkaConfig.java b/cdc-service/src/main/java/com/xtrmetl/cdc/config/KafkaConfig.java index 36388b64..16e6df9f 100644 --- a/cdc-service/src/main/java/com/xtrmetl/cdc/config/KafkaConfig.java +++ b/cdc-service/src/main/java/com/xtrmetl/cdc/config/KafkaConfig.java @@ -62,6 +62,14 @@ public DefaultErrorHandler kafkaListenerErrorHandler( return errorHandler; } + /** + * Builds the replica Kafka listener factory with record-level acknowledgement semantics. + * + * @param consumerFactory Kafka consumer factory for replica records + * @param kafkaListenerErrorHandler bounded retry and dead-letter error handler + * @param concurrency requested listener concurrency + * @return configured listener container factory + */ @Bean public ConcurrentKafkaListenerContainerFactory kafkaListenerContainerFactory( @NonNull ConsumerFactory consumerFactory, diff --git a/cdc-service/src/test/java/com/xtrmetl/cdc/config/KafkaConfigTest.java b/cdc-service/src/test/java/com/xtrmetl/cdc/config/KafkaConfigTest.java index 0efd35f8..4c116b14 100644 --- a/cdc-service/src/test/java/com/xtrmetl/cdc/config/KafkaConfigTest.java +++ b/cdc-service/src/test/java/com/xtrmetl/cdc/config/KafkaConfigTest.java @@ -54,6 +54,21 @@ void configuresErrorHandlerWithDeadLetterRecovererAndBackOff() { assertFalse(classifier.classify(new IllegalStateException("test"))); } + @Test + void preservesZeroRetrySettings() { + KafkaConfig config = new KafkaConfig(); + KafkaTemplate kafkaTemplate = mock(KafkaTemplate.class); + + DefaultErrorHandler errorHandler = config.kafkaListenerErrorHandler(kafkaTemplate, 0L, 0L); + + Object failureTracker = ReflectionTestUtils.getField(errorHandler, "failureTracker"); + assertNotNull(failureTracker); + FixedBackOff fixedBackOff = (FixedBackOff) ReflectionTestUtils.getField(failureTracker, "backOff"); + assertNotNull(fixedBackOff); + assertEquals(0L, fixedBackOff.getInterval()); + assertEquals(0L, fixedBackOff.getMaxAttempts()); + } + @Test void rejectsNegativeRetryBackoffBeforeBuildingErrorHandler() { KafkaConfig config = new KafkaConfig();