From 7881c8f99d2174678adfdcf08d14ef73266bc512 Mon Sep 17 00:00:00 2001 From: Chris Bono Date: Sun, 19 Jan 2025 18:13:51 -0600 Subject: [PATCH] Fix batch listener error handling (#1007) Updates the batch listener error handler code to handle the case where there is only a single outstanding message in the batch list when retries are expired. See #998 Signed-off-by: darshimo Co-authored-by: darshimo --- ...DefaultPulsarMessageListenerContainer.java | 3 +- ...efaultPulsarConsumerErrorHandlerTests.java | 77 +++++++++++++++++++ 2 files changed, 78 insertions(+), 2 deletions(-) diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java index 3c599a82..72e49938 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java @@ -593,12 +593,11 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess pulsarBatchListenerFailedException); handleAck(pulsarMessage); if (messageList.size() == 1) { + messageList.remove(0); messagesPendingInBatch.set(false); } else { messageList = messageList.subList(1, messageList.size()); - } - if (!messageList.isEmpty()) { messagesPendingInBatch.set(true); } this.pulsarConsumerErrorHandler.clearMessage(); diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarConsumerErrorHandlerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarConsumerErrorHandlerTests.java index 285ef285..8b9adb57 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarConsumerErrorHandlerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarConsumerErrorHandlerTests.java @@ -38,6 +38,7 @@ import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.Schema; import org.junit.jupiter.api.Test; +import org.springframework.core.log.LogAccessor; import org.springframework.pulsar.core.DefaultPulsarConsumerFactory; import org.springframework.pulsar.core.DefaultPulsarProducerFactory; import org.springframework.pulsar.core.PulsarOperations; @@ -51,6 +52,8 @@ import org.springframework.util.backoff.FixedBackOff; */ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContainerSupport { + private final LogAccessor logger = new LogAccessor(this.getClass()); + @Test @SuppressWarnings("unchecked") void happyPathErrorHandlingForRecordMessageListener() throws Exception { @@ -564,4 +567,78 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain pulsarClient.close(); } + @Test + @SuppressWarnings("unchecked") + void whenBatchRecordListenerOneMessageBatchFailsThenSentToDltProperly() throws Exception { + var topicName = "default-error-handler-tests-9"; + var pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); + var pulsarConsumerFactory = new DefaultPulsarConsumerFactory(pulsarClient, + List.of((consumerBuilder) -> { + consumerBuilder.topic(topicName); + consumerBuilder.subscriptionName("%s-sub".formatted(topicName)); + })); + // Prepare container for batch consume + var pulsarContainerProperties = new PulsarContainerProperties(); + pulsarContainerProperties.setSchema(Schema.INT32); + pulsarContainerProperties.setAckMode(AckMode.MANUAL); + pulsarContainerProperties.setBatchListener(true); + pulsarContainerProperties.setMaxNumMessages(1); + pulsarContainerProperties.setBatchTimeoutMillis(60_000); + PulsarBatchAcknowledgingMessageListener pulsarBatchMessageListener = mock(); + doAnswer(invocation -> { + List> message = invocation.getArgument(1); + Message integerMessage = message.get(0); + Integer value = integerMessage.getValue(); + if (value == 0) { + throw new PulsarBatchListenerFailedException("failed", integerMessage); + } + Acknowledgement acknowledgment = invocation.getArgument(2); + List messageIds = new ArrayList<>(); + for (Message integerMessage1 : message) { + messageIds.add(integerMessage1.getMessageId()); + } + acknowledgment.acknowledge(messageIds); + return new Object(); + }).when(pulsarBatchMessageListener).received(any(Consumer.class), any(List.class), any(Acknowledgement.class)); + pulsarContainerProperties.setMessageListener(pulsarBatchMessageListener); + var container = new DefaultPulsarMessageListenerContainer<>(pulsarConsumerFactory, pulsarContainerProperties); + + // Set error handler to recover after 2 retries + PulsarTemplate mockPulsarTemplate = mock(RETURNS_DEEP_STUBS); + PulsarOperations.SendMessageBuilder sendMessageBuilderMock = mock(); + when(mockPulsarTemplate.newMessage(any(Integer.class)) + .withTopic(any(String.class)) + .withMessageCustomizer(any(TypedMessageBuilderCustomizer.class))).thenReturn(sendMessageBuilderMock); + container.setPulsarConsumerErrorHandler(new DefaultPulsarConsumerErrorHandler<>( + new PulsarDeadLetterPublishingRecoverer<>(mockPulsarTemplate), new FixedBackOff(100, 2))); + try { + container.start(); + // Send single message in batch + var pulsarProducerFactory = new DefaultPulsarProducerFactory(pulsarClient, topicName); + var pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + pulsarTemplate.sendAsync(0); + // Initial call should fail + // Next 2 calls should fail (retries 2) + // No more calls after that - msg should go to DLT + await().atMost(Duration.ofSeconds(30)) + .untilAsserted(() -> verify(pulsarBatchMessageListener, times(3)).received(any(Consumer.class), + any(List.class), any(Acknowledgement.class))); + await().atMost(Duration.ofSeconds(30)) + .untilAsserted(() -> verify(sendMessageBuilderMock, times(1)).sendAsync()); + } + finally { + safeStopContainer(container); + } + pulsarClient.close(); + } + + private void safeStopContainer(PulsarMessageListenerContainer container) { + try { + container.stop(); + } + catch (Exception ex) { + logger.warn(ex, "Failed to stop container %s: %s".formatted(container, ex.getMessage())); + } + } + }