From 16e0c9ee8c2bafa6a9b16f9dc4fb977eb2b71545 Mon Sep 17 00:00:00 2001 From: Chris Bono Date: Fri, 17 Jan 2025 19:22:48 -0600 Subject: [PATCH] Polish "Fix batch listener error handling" Updates the error handler logic to only use `List.remove` in the special case when 1 item in batch. Also adds a test for the special case when 1 item in batch. See #998 --- ...DefaultPulsarMessageListenerContainer.java | 10 ++- ...efaultPulsarConsumerErrorHandlerTests.java | 77 +++++++++++++++++++ 2 files changed, 85 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 0d071958..baba3e84 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 @@ -781,8 +781,14 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess this.pulsarConsumerErrorHandler.recoverMessage(this.consumer, pulsarMessage, pulsarBatchListenerFailedException); handleAck(pulsarMessage, txn); - messageList.remove(0); - messagesPendingInBatch.set(!messageList.isEmpty()); + if (messageList.size() == 1) { + messageList.remove(0); + messagesPendingInBatch.set(false); + } + else { + messageList = messageList.subList(1, messageList.size()); + messagesPendingInBatch.set(true); + } this.pulsarConsumerErrorHandler.clearMessage(); } return messageList; 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())); + } + } + }