From 18e509ec01bd33d24575368ce0cea77fd24d756a Mon Sep 17 00:00:00 2001 From: Dan Blackney Date: Fri, 3 May 2024 22:28:22 +0800 Subject: [PATCH] GH-3227: Implement handleOne() in CDEH Fixes: #3227 * Implement handleOne() in `CommonDelegatingErrorHandler` * Add tests for handle methods in `CommonDelegatingErrorHandler` * Add tests for the new `handleOne()` method, as well as a test for `handleOtherException()` * Checkstyle fixes (cherry picked from commit 4e06c2c16554159179eea316c9089f1e95d85a54) --- .../CommonDelegatingErrorHandler.java | 15 +++++ .../CommonDelegatingErrorHandlerTests.java | 56 ++++++++++++++++++- 2 files changed, 69 insertions(+), 2 deletions(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/CommonDelegatingErrorHandler.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/CommonDelegatingErrorHandler.java index 003f1245..c6344b1d 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/CommonDelegatingErrorHandler.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/CommonDelegatingErrorHandler.java @@ -38,6 +38,8 @@ import org.springframework.util.Assert; * * @author Gary Russell * @author Adrian Chlebosz + * @author Antonin Arquey + * @author Dan Blackney * @since 2.8 * */ @@ -188,6 +190,19 @@ public class CommonDelegatingErrorHandler implements CommonErrorHandler { } } + @Override + public boolean handleOne(Exception thrownException, ConsumerRecord record, Consumer consumer, + MessageListenerContainer container) { + + CommonErrorHandler handler = findDelegate(thrownException); + if (handler != null) { + return handler.handleOne(thrownException, record, consumer, container); + } + else { + return this.defaultErrorHandler.handleOne(thrownException, record, consumer, container); + } + } + @Nullable private CommonErrorHandler findDelegate(Throwable thrownException) { Throwable cause = findCause(thrownException); diff --git a/spring-kafka/src/test/java/org/springframework/kafka/listener/CommonDelegatingErrorHandlerTests.java b/spring-kafka/src/test/java/org/springframework/kafka/listener/CommonDelegatingErrorHandlerTests.java index 1a25253a..73b1ab4d 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/listener/CommonDelegatingErrorHandlerTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/listener/CommonDelegatingErrorHandlerTests.java @@ -18,6 +18,7 @@ package org.springframework.kafka.listener; import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; @@ -27,6 +28,7 @@ import java.util.Collections; import java.util.Map; import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.junit.jupiter.api.Test; @@ -39,13 +41,15 @@ import org.springframework.kafka.test.utils.KafkaTestUtils; * * @author Gary Russell * @author Adrian Chlebosz + * @author Antonin Arquey + * @author Dan Blackney * @since 2.8 * */ public class CommonDelegatingErrorHandlerTests { @Test - void testRecordDelegates() { + void testHandleRemainingDelegates() { var def = mock(CommonErrorHandler.class); var one = mock(CommonErrorHandler.class); var two = mock(CommonErrorHandler.class); @@ -69,7 +73,7 @@ public class CommonDelegatingErrorHandlerTests { } @Test - void testBatchDelegates() { + void testHandleBatchDelegates() { var def = mock(CommonErrorHandler.class); var one = mock(CommonErrorHandler.class); var two = mock(CommonErrorHandler.class); @@ -92,6 +96,54 @@ public class CommonDelegatingErrorHandlerTests { verify(one).handleBatch(any(), any(), any(), any(), any()); } + @Test + void testHandleOtherExceptionDelegates() { + var def = mock(CommonErrorHandler.class); + var one = mock(CommonErrorHandler.class); + var two = mock(CommonErrorHandler.class); + var three = mock(CommonErrorHandler.class); + var eh = new CommonDelegatingErrorHandler(def); + eh.setErrorHandlers(Map.of(IllegalStateException.class, one, IllegalArgumentException.class, two)); + eh.addDelegate(RuntimeException.class, three); + + eh.handleOtherException(wrap(new IOException()), mock(Consumer.class), + mock(MessageListenerContainer.class), true); + verify(def).handleOtherException(any(), any(), any(), anyBoolean()); + eh.handleOtherException(wrap(new KafkaException("test")), mock(Consumer.class), + mock(MessageListenerContainer.class), true); + verify(three).handleOtherException(any(), any(), any(), anyBoolean()); + eh.handleOtherException(wrap(new IllegalArgumentException()), mock(Consumer.class), + mock(MessageListenerContainer.class), true); + verify(two).handleOtherException(any(), any(), any(), anyBoolean()); + eh.handleOtherException(wrap(new IllegalStateException()), mock(Consumer.class), + mock(MessageListenerContainer.class), true); + verify(one).handleOtherException(any(), any(), any(), anyBoolean()); + } + + @Test + void testHandleOneDelegates() { + var def = mock(CommonErrorHandler.class); + var one = mock(CommonErrorHandler.class); + var two = mock(CommonErrorHandler.class); + var three = mock(CommonErrorHandler.class); + var eh = new CommonDelegatingErrorHandler(def); + eh.setErrorHandlers(Map.of(IllegalStateException.class, one, IllegalArgumentException.class, two)); + eh.addDelegate(RuntimeException.class, three); + + eh.handleOne(wrap(new IOException()), mock(ConsumerRecord.class), mock(Consumer.class), + mock(MessageListenerContainer.class)); + verify(def).handleOne(any(), any(), any(), any()); + eh.handleOne(wrap(new KafkaException("test")), mock(ConsumerRecord.class), mock(Consumer.class), + mock(MessageListenerContainer.class)); + verify(three).handleOne(any(), any(), any(), any()); + eh.handleOne(wrap(new IllegalArgumentException()), mock(ConsumerRecord.class), mock(Consumer.class), + mock(MessageListenerContainer.class)); + verify(two).handleOne(any(), any(), any(), any()); + eh.handleOne(wrap(new IllegalStateException()), mock(ConsumerRecord.class), mock(Consumer.class), + mock(MessageListenerContainer.class)); + verify(one).handleOne(any(), any(), any(), any()); + } + @Test void testDelegateForThrowableIsAppliedWhenCauseTraversingIsEnabled() { var defaultHandler = mock(CommonErrorHandler.class);