diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java index 0100865f..c722af92 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java @@ -92,6 +92,7 @@ import org.springframework.util.concurrent.ListenableFutureCallback; * @author Artem Bilan * @author Loic Talhouarne * @author Vladimir Tsanev + * @author Chen Binbin * @author Yang Qiju * @author Tom van den Berge */ @@ -719,11 +720,23 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener break; } catch (Exception e) { - if (this.containerProperties.getGenericErrorHandler() != null) { - this.containerProperties.getGenericErrorHandler().handle(e, null); + try { + GenericErrorHandler containerErrorHandler = this.containerProperties.getGenericErrorHandler(); + if (containerErrorHandler != null) { + if (containerErrorHandler instanceof ConsumerAwareErrorHandler + || containerErrorHandler instanceof ConsumerAwareBatchErrorHandler) { + containerErrorHandler.handle(e, null, this.consumer); + } + else { + containerErrorHandler.handle(e, null); + } + } + else { + this.logger.error("Container exception", e); + } } - else { - this.logger.error("Container exception", e); + catch (Exception ex) { + this.logger.error("Container exception", ex); } } } diff --git a/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java b/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java index 95910841..7595860a 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java @@ -134,11 +134,13 @@ public class KafkaMessageListenerContainerTests { private static String topic18 = "testTopic18"; + private static String topic19 = "testTopic19"; + @ClassRule public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, topic1, topic2, topic3, topic4, topic5, topic6, topic7, topic8, topic9, topic10, topic11, topic12, topic13, topic14, topic15, topic16, topic17, - topic18); + topic18, topic19); @Rule public TestName testName = new TestName(); @@ -1718,6 +1720,69 @@ public class KafkaMessageListenerContainerTests { container.stop(); } + @Test + public void testExceptionWhenCommitAfterRebalance() throws Exception { + final CountDownLatch rebalanceLatch = new CountDownLatch(2); + final CountDownLatch consumeLatch = new CountDownLatch(7); + + Map props = KafkaTestUtils.consumerProps("test19", "false", embeddedKafka); + props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest"); + props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 15000); + DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(props); + ContainerProperties containerProps = new ContainerProperties(topic19); + containerProps.setMessageListener((MessageListener) messages -> { + logger.info("listener: " + messages); + consumeLatch.countDown(); + try { + Thread.sleep(3000); + } + catch (InterruptedException e) { + e.printStackTrace(); + } + }); + containerProps.setSyncCommits(true); + containerProps.setAckMode(AckMode.BATCH); + containerProps.setPollTimeout(100); + containerProps.setAckOnError(false); + containerProps.setErrorHandler(new SeekToCurrentErrorHandler()); + + Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); + ProducerFactory pf = new DefaultKafkaProducerFactory<>(senderProps); + KafkaTemplate template = new KafkaTemplate<>(pf); + template.setDefaultTopic(topic19); + + containerProps.setConsumerRebalanceListener(new ConsumerRebalanceListener() { + + @Override + public void onPartitionsRevoked(Collection partitions) { + } + + @Override + public void onPartitionsAssigned(Collection partitions) { + logger.info("rebalance occurred."); + rebalanceLatch.countDown(); + } + }); + + KafkaMessageListenerContainer container = + new KafkaMessageListenerContainer<>(cf, containerProps); + container.setBeanName("testContainerException"); + container.start(); + ContainerTestUtils.waitForAssignment(container, embeddedKafka.getPartitionsPerTopic()); + container.pause(); + + for (int i = 0; i < 6; i++) { + template.sendDefault(0, 0, "a"); + } + template.flush(); + + container.resume(); + // should be rebalanced and consume again + assertThat(rebalanceLatch.await(60, TimeUnit.SECONDS)).isTrue(); + assertThat(consumeLatch.await(60, TimeUnit.SECONDS)).isTrue(); + container.stop(); + } + private Consumer spyOnConsumer(KafkaMessageListenerContainer container) { Consumer consumer = spy( KafkaTestUtils.getPropertyValue(container, "listenerConsumer.consumer", Consumer.class));