diff --git a/sdk/spring/spring-cloud-azure-integration-tests/src/test/java/com/azure/spring/cloud/integration/tests/eventhubs/binder/EventHubsBinderBatchModeIT.java b/sdk/spring/spring-cloud-azure-integration-tests/src/test/java/com/azure/spring/cloud/integration/tests/eventhubs/binder/EventHubsBinderBatchModeIT.java index 0c00544f38ac..83d513b152fb 100644 --- a/sdk/spring/spring-cloud-azure-integration-tests/src/test/java/com/azure/spring/cloud/integration/tests/eventhubs/binder/EventHubsBinderBatchModeIT.java +++ b/sdk/spring/spring-cloud-azure-integration-tests/src/test/java/com/azure/spring/cloud/integration/tests/eventhubs/binder/EventHubsBinderBatchModeIT.java @@ -9,6 +9,7 @@ import org.springframework.boot.test.context.SpringBootTest; import org.springframework.boot.test.context.TestConfiguration; import org.springframework.context.annotation.Bean; +import org.springframework.integration.annotation.ServiceActivator; import org.springframework.messaging.Message; import org.springframework.messaging.support.GenericMessage; import org.springframework.test.context.ActiveProfiles; @@ -62,6 +63,11 @@ Consumer>> consume() { } }; } + + @ServiceActivator(inputChannel = "errorChannel") + public void processError(Message sendFailedMsg) { + LOGGER.info("receive error message: '{}'", sendFailedMsg.getPayload()); + } } @Test diff --git a/sdk/spring/spring-cloud-azure-integration-tests/src/test/java/com/azure/spring/cloud/integration/tests/eventhubs/binder/EventHubsBinderManualModeIT.java b/sdk/spring/spring-cloud-azure-integration-tests/src/test/java/com/azure/spring/cloud/integration/tests/eventhubs/binder/EventHubsBinderManualModeIT.java index 0d531b2bf200..be68618b32b0 100644 --- a/sdk/spring/spring-cloud-azure-integration-tests/src/test/java/com/azure/spring/cloud/integration/tests/eventhubs/binder/EventHubsBinderManualModeIT.java +++ b/sdk/spring/spring-cloud-azure-integration-tests/src/test/java/com/azure/spring/cloud/integration/tests/eventhubs/binder/EventHubsBinderManualModeIT.java @@ -12,6 +12,7 @@ import org.springframework.boot.test.context.SpringBootTest; import org.springframework.boot.test.context.TestConfiguration; import org.springframework.context.annotation.Bean; +import org.springframework.integration.annotation.ServiceActivator; import org.springframework.messaging.Message; import org.springframework.messaging.support.GenericMessage; import org.springframework.test.context.ActiveProfiles; @@ -65,6 +66,11 @@ Consumer> consume() { } }; } + + @ServiceActivator(inputChannel = "errorChannel") + public void processError(Message sendFailedMsg) { + LOGGER.info("receive error message: '{}'", sendFailedMsg.getPayload()); + } } @Test diff --git a/sdk/spring/spring-cloud-azure-integration-tests/src/test/java/com/azure/spring/cloud/integration/tests/eventhubs/binder/EventHubsBinderProduceErrorIT.java b/sdk/spring/spring-cloud-azure-integration-tests/src/test/java/com/azure/spring/cloud/integration/tests/eventhubs/binder/EventHubsBinderProduceErrorIT.java new file mode 100644 index 000000000000..6f3a75748af0 --- /dev/null +++ b/sdk/spring/spring-cloud-azure-integration-tests/src/test/java/com/azure/spring/cloud/integration/tests/eventhubs/binder/EventHubsBinderProduceErrorIT.java @@ -0,0 +1,82 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the MIT License. +package com.azure.spring.cloud.integration.tests.eventhubs.binder; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.boot.test.context.TestConfiguration; +import org.springframework.context.annotation.Bean; +import org.springframework.integration.annotation.ServiceActivator; +import org.springframework.messaging.Message; +import org.springframework.messaging.support.GenericMessage; +import org.springframework.test.context.ActiveProfiles; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Sinks; + +import java.util.List; +import java.util.UUID; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.function.Consumer; +import java.util.function.Supplier; + +import static org.assertj.core.api.Assertions.assertThat; + +@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.NONE) +@ActiveProfiles(value = { "eventhubs-binder", "produceerror" }) +class EventHubsBinderProduceErrorIT { + + private static final Logger LOGGER = LoggerFactory.getLogger(EventHubsBinderProduceErrorIT.class); + + private static final String MESSAGE = UUID.randomUUID().toString(); + + private static final CountDownLatch LATCH = new CountDownLatch(1); + + @Autowired + private Sinks.Many> many; + + @TestConfiguration + static class TestConfig { + + @Bean + Sinks.Many> many() { + return Sinks.many().unicast().onBackpressureBuffer(); + } + + @Bean + Supplier>> supply(Sinks.Many> many) { + return () -> many.asFlux() + .doOnNext(m -> LOGGER.info("Manually sending message {}", m.getPayload())) + .doOnError(t -> LOGGER.error("Error encountered", t)); + } + + @Bean + Consumer>> consume() { + return message -> { + List payload = message.getPayload(); + LOGGER.info("EventHubsBinderProduceErrorIT: New message received: '{}'", payload); + Assertions.fail("EventHubsBinderProduceErrorIT: can't be here"); + }; + } + + @ServiceActivator(inputChannel = "errorChannel") + public void processError(Message sendFailedMsg) { + LOGGER.info("receive error message: '{}'", sendFailedMsg); + LATCH.countDown(); + } + } + + @Test + void testSendAndReceiveMessage() throws InterruptedException { + LOGGER.info("EventHubsBinderProduceErrorIT begin."); + EventHubsBinderProduceErrorIT.LATCH.await(15, TimeUnit.SECONDS); + LOGGER.info("Send a message:" + MESSAGE + "."); + many.emitNext(new GenericMessage<>(MESSAGE), Sinks.EmitFailureHandler.FAIL_FAST); + assertThat(EventHubsBinderProduceErrorIT.LATCH.await(300, TimeUnit.SECONDS)).isTrue(); + LOGGER.info("EventHubsBinderProduceErrorIT end."); + } +} diff --git a/sdk/spring/spring-cloud-azure-integration-tests/src/test/java/com/azure/spring/cloud/integration/tests/eventhubs/binder/EventHubsBinderRecordModeIT.java b/sdk/spring/spring-cloud-azure-integration-tests/src/test/java/com/azure/spring/cloud/integration/tests/eventhubs/binder/EventHubsBinderRecordModeIT.java index 386a2391be23..48f379864f71 100644 --- a/sdk/spring/spring-cloud-azure-integration-tests/src/test/java/com/azure/spring/cloud/integration/tests/eventhubs/binder/EventHubsBinderRecordModeIT.java +++ b/sdk/spring/spring-cloud-azure-integration-tests/src/test/java/com/azure/spring/cloud/integration/tests/eventhubs/binder/EventHubsBinderRecordModeIT.java @@ -9,6 +9,7 @@ import org.springframework.boot.test.context.SpringBootTest; import org.springframework.boot.test.context.TestConfiguration; import org.springframework.context.annotation.Bean; +import org.springframework.integration.annotation.ServiceActivator; import org.springframework.messaging.Message; import org.springframework.messaging.support.GenericMessage; import org.springframework.test.context.ActiveProfiles; @@ -58,6 +59,12 @@ Consumer> consume() { } }; } + + @ServiceActivator(inputChannel = "errorChannel") + public void processError(Message sendFailedMsg) { + LOGGER.info("receive error message: '{}'", sendFailedMsg.getPayload()); + } + } @Test diff --git a/sdk/spring/spring-cloud-azure-integration-tests/src/test/java/com/azure/spring/cloud/integration/tests/eventhubs/binder/EventHubsBinderSyncModeIT.java b/sdk/spring/spring-cloud-azure-integration-tests/src/test/java/com/azure/spring/cloud/integration/tests/eventhubs/binder/EventHubsBinderSyncModeIT.java index 2266cfed07d6..7fa9a53668c5 100644 --- a/sdk/spring/spring-cloud-azure-integration-tests/src/test/java/com/azure/spring/cloud/integration/tests/eventhubs/binder/EventHubsBinderSyncModeIT.java +++ b/sdk/spring/spring-cloud-azure-integration-tests/src/test/java/com/azure/spring/cloud/integration/tests/eventhubs/binder/EventHubsBinderSyncModeIT.java @@ -9,6 +9,7 @@ import org.springframework.boot.test.context.SpringBootTest; import org.springframework.boot.test.context.TestConfiguration; import org.springframework.context.annotation.Bean; +import org.springframework.integration.annotation.ServiceActivator; import org.springframework.messaging.Message; import org.springframework.messaging.support.GenericMessage; import org.springframework.test.context.ActiveProfiles; @@ -53,12 +54,17 @@ Supplier>> supply(Sinks.Many> many) { @Bean Consumer> consume() { return message -> { - LOGGER.info("EventHubBinderRecordModeIT: New message received: '{}'", message.getPayload()); + LOGGER.info("EventHubBinderSyncModeIT: New message received: '{}'", message.getPayload()); if (message.getPayload().equals(EventHubsBinderSyncModeIT.MESSAGE) && message.getHeaders().containsKey("x-opt-enqueued-time")) { LATCH.countDown(); } }; } + + @ServiceActivator(inputChannel = "errorChannel") + public void processError(Message sendFailedMsg) { + LOGGER.info("receive error message: '{}'", sendFailedMsg.getPayload()); + } } @Test diff --git a/sdk/spring/spring-cloud-azure-integration-tests/src/test/resources/application-eventhubs-binder.yml b/sdk/spring/spring-cloud-azure-integration-tests/src/test/resources/application-eventhubs-binder.yml index 16f175ef2524..b6ba9703598b 100644 --- a/sdk/spring/spring-cloud-azure-integration-tests/src/test/resources/application-eventhubs-binder.yml +++ b/sdk/spring/spring-cloud-azure-integration-tests/src/test/resources/application-eventhubs-binder.yml @@ -135,3 +135,25 @@ spring: config: activate: on-profile: sync +--- +spring: + cloud: + azure: + eventhubs: + processor: + checkpoint-store: + container-name: ${AZURE_EVENTHUB_NAME_FOR_BINDER_SYNC} + stream: + bindings: + consume-in-0: + destination: notexists + supply-out-0: + destination: notexists + eventhubs: + bindings: + supply-out-0: + producer: + sync: true + config: + activate: + on-profile: produceerror diff --git a/sdk/spring/spring-messaging-azure-eventhubs/src/main/java/com/azure/spring/messaging/eventhubs/core/EventHubsTemplate.java b/sdk/spring/spring-messaging-azure-eventhubs/src/main/java/com/azure/spring/messaging/eventhubs/core/EventHubsTemplate.java index 0b9343efd0e9..49b0dfdea3e2 100644 --- a/sdk/spring/spring-messaging-azure-eventhubs/src/main/java/com/azure/spring/messaging/eventhubs/core/EventHubsTemplate.java +++ b/sdk/spring/spring-messaging-azure-eventhubs/src/main/java/com/azure/spring/messaging/eventhubs/core/EventHubsTemplate.java @@ -100,8 +100,16 @@ public Mono sendAsync(String destination, Message message) { private Mono doSend(String destination, List events, PartitionSupplier partitionSupplier) { EventHubProducerAsyncClient producer = producerFactory.createProducer(destination); CreateBatchOptions options = buildCreateBatchOptions(partitionSupplier); - AtomicReference currentBatch = new AtomicReference<>( - producer.createBatch(options).block()); + + EventDataBatch eventDataBatch = null; + try { + eventDataBatch = producer.createBatch(options).block(); + } catch (Exception e) { + LOGGER.error("EventDataBatch create error.", e); + return Mono.error(e); + } + AtomicReference currentBatch = new AtomicReference<>(eventDataBatch); + Flux.fromIterable(events).flatMap(event -> { final EventDataBatch batch = currentBatch.get(); try {