From 0447b19aa5a8398c3a9f5adc8c6fb661bae7b58d Mon Sep 17 00:00:00 2001 From: Christophe Bornet Date: Fri, 21 Oct 2022 01:59:31 +0200 Subject: [PATCH] Add cache to DefaultReactivePulsarSenderFactory (#169) --- .../DefaultReactivePulsarSenderFactory.java | 16 ++- ...aultReactiveMessageSenderFactoryTests.java | 130 +++++++----------- .../reactive/ReactivePulsarTemplateTests.java | 4 +- 3 files changed, 61 insertions(+), 89 deletions(-) diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/DefaultReactivePulsarSenderFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/DefaultReactivePulsarSenderFactory.java index 6a2d8ba3..48e4e719 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/DefaultReactivePulsarSenderFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/DefaultReactivePulsarSenderFactory.java @@ -24,6 +24,7 @@ import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.reactive.client.adapter.AdaptedReactivePulsarClientFactory; import org.apache.pulsar.reactive.client.api.ReactiveMessageSender; import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderBuilder; +import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderCache; import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderSpec; import org.apache.pulsar.reactive.client.api.ReactivePulsarClient; @@ -44,15 +45,21 @@ public class DefaultReactivePulsarSenderFactory implements ReactivePulsarSend private final ReactiveMessageSenderSpec reactiveMessageSenderSpec; + private final ReactiveMessageSenderCache reactiveMessageSenderCache; + public DefaultReactivePulsarSenderFactory(PulsarClient pulsarClient, - ReactiveMessageSenderSpec reactiveMessageSenderSpec) { - this(AdaptedReactivePulsarClientFactory.create(pulsarClient), reactiveMessageSenderSpec); + ReactiveMessageSenderSpec reactiveMessageSenderSpec, + ReactiveMessageSenderCache reactiveMessageSenderCache) { + this(AdaptedReactivePulsarClientFactory.create(pulsarClient), reactiveMessageSenderSpec, + reactiveMessageSenderCache); } public DefaultReactivePulsarSenderFactory(ReactivePulsarClient reactivePulsarClient, - ReactiveMessageSenderSpec reactiveMessageSenderSpec) { + ReactiveMessageSenderSpec reactiveMessageSenderSpec, + ReactiveMessageSenderCache reactiveMessageSenderCache) { this.reactivePulsarClient = reactivePulsarClient; this.reactiveMessageSenderSpec = reactiveMessageSenderSpec; + this.reactiveMessageSenderCache = reactiveMessageSenderCache; } @Override @@ -81,6 +88,9 @@ public class DefaultReactivePulsarSenderFactory implements ReactivePulsarSend sender.applySpec(this.reactiveMessageSenderSpec); } sender.topic(resolvedTopic); + if (this.reactiveMessageSenderCache != null) { + sender.cache(this.reactiveMessageSenderCache); + } if (messageRouter != null) { sender.messageRouter(messageRouter); } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/DefaultReactiveMessageSenderFactoryTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/DefaultReactiveMessageSenderFactoryTests.java index 78623e7d..3cd1352a 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/DefaultReactiveMessageSenderFactoryTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/DefaultReactiveMessageSenderFactoryTests.java @@ -16,143 +16,105 @@ package org.springframework.pulsar.core.reactive; +import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatIllegalArgumentException; -import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.never; -import static org.mockito.Mockito.verify; -import static org.mockito.Mockito.when; -import java.time.Duration; import java.util.Arrays; import java.util.Collections; -import java.util.concurrent.CompletableFuture; +import java.util.List; -import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.MessageRouter; -import org.apache.pulsar.client.api.Producer; -import org.apache.pulsar.client.api.ProducerBuilder; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.Schema; -import org.apache.pulsar.client.api.TypedMessageBuilder; -import org.apache.pulsar.reactive.client.api.MessageSpec; +import org.apache.pulsar.reactive.client.adapter.AdaptedReactivePulsarClientFactory; import org.apache.pulsar.reactive.client.api.MutableReactiveMessageSenderSpec; import org.apache.pulsar.reactive.client.api.ReactiveMessageSender; -import org.junit.jupiter.api.BeforeEach; +import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderCache; +import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderSpec; +import org.assertj.core.api.InstanceOfAssertFactories; +import org.assertj.core.api.ObjectAssert; import org.junit.jupiter.api.Test; -import reactor.core.publisher.Mono; - /** * Common tests for {@link DefaultReactivePulsarSenderFactory} * * @author Christophe Bornet */ -@SuppressWarnings("unchecked") class DefaultReactiveMessageSenderFactoryTests { protected final Schema schema = Schema.STRING; - private ProducerBuilder producerBuilder; - - private PulsarClient pulsarClient; - - @BeforeEach - void createPulsarClient() { - pulsarClient = mock(PulsarClient.class); - producerBuilder = mock(ProducerBuilder.class); - Producer producer = mock(Producer.class); - TypedMessageBuilder mockMessage = mock(TypedMessageBuilder.class); - - when(mockMessage.sendAsync()).thenReturn(CompletableFuture.completedFuture(MessageId.latest)); - when(producer.newMessage()).thenReturn(mockMessage); - when(producer.closeAsync()).thenReturn(CompletableFuture.completedFuture(null)); - when(producerBuilder.createAsync()).thenReturn(CompletableFuture.completedFuture(producer)); - when(pulsarClient.newProducer(schema)).thenReturn(producerBuilder); + @Test + void createSenderWithSpecificTopic() { + testCreateSender(null, null, "topic1", null, null, "topic1", null); } @Test - void createProducerWithSpecificTopic() { - ReactivePulsarSenderFactory senderFactory = new DefaultReactivePulsarSenderFactory<>(pulsarClient, - null); - ReactiveMessageSender sender = senderFactory.createReactiveMessageSender("topic1", schema); - sender.sendMessage(Mono.just(MessageSpec.of("test"))).block(Duration.ofSeconds(5)); - assertSenderHasTopicAndRouter("topic1", null); - } - - @Test - void createProducerWithSpecificTopicAndMessageRouter() { - ReactivePulsarSenderFactory senderFactory = new DefaultReactivePulsarSenderFactory<>(pulsarClient, - null); + void createSenderWithSpecificTopicAndMessageRouter() { MessageRouter router = mock(MessageRouter.class); - ReactiveMessageSender sender = senderFactory.createReactiveMessageSender("topic1", schema, router); - sender.sendMessage(Mono.just(MessageSpec.of("test"))).block(Duration.ofSeconds(5)); - assertSenderHasTopicAndRouter("topic1", router); + + testCreateSender(null, null, "topic1", router, null, "topic1", router); } @Test - void createProducerWithDefaultTopic() { + void createSenderWithDefaultTopic() { MutableReactiveMessageSenderSpec senderSpec = new MutableReactiveMessageSenderSpec(); senderSpec.setTopicName("topic0"); - ReactivePulsarSenderFactory senderFactory = new DefaultReactivePulsarSenderFactory<>(pulsarClient, - senderSpec); - ReactiveMessageSender sender = senderFactory.createReactiveMessageSender(null, schema); - sender.sendMessage(Mono.just(MessageSpec.of("test"))).block(Duration.ofSeconds(5)); - assertSenderHasTopicAndRouter("topic0", null); + + testCreateSender(senderSpec, null, null, null, null, "topic0", null); } @Test - void createProducerWithDefaultTopicAndMessageRouter() { + void createSenderWithDefaultTopicAndMessageRouter() { MutableReactiveMessageSenderSpec senderSpec = new MutableReactiveMessageSenderSpec(); senderSpec.setTopicName("topic0"); - ReactivePulsarSenderFactory senderFactory = new DefaultReactivePulsarSenderFactory<>(pulsarClient, - senderSpec); MessageRouter router = mock(MessageRouter.class); - ReactiveMessageSender sender = senderFactory.createReactiveMessageSender(null, schema, router); - sender.sendMessage(Mono.just(MessageSpec.of("test"))).block(Duration.ofSeconds(5)); - assertSenderHasTopicAndRouter("topic0", router); + + testCreateSender(senderSpec, null, null, router, null, "topic0", router); + } @Test - void createProducerWithSingleProducerCustomizer() { - ReactivePulsarSenderFactory senderFactory = new DefaultReactivePulsarSenderFactory<>(pulsarClient, - null); - ReactiveMessageSenderBuilderCustomizer customizer = builder -> builder.topic("topic1"); - ReactiveMessageSender sender = senderFactory.createReactiveMessageSender("topic0", schema, null, - Collections.singletonList(customizer)); - sender.sendMessage(Mono.just(MessageSpec.of("test"))).block(Duration.ofSeconds(5)); - assertSenderHasTopicAndRouter("topic1", null); + void createSenderWithSingleSenderCustomizer() { + testCreateSender(null, null, "topic1", null, Collections.singletonList(builder -> builder.topic("topic1")), + "topic1", null); } @Test - void createProducerWithMultipleProducerCustomizer() { - ReactivePulsarSenderFactory senderFactory = new DefaultReactivePulsarSenderFactory<>(pulsarClient, - null); + void createSenderWithMultipleSenderCustomizer() { ReactiveMessageSenderBuilderCustomizer customizer1 = builder -> builder.topic("topic1"); MessageRouter router = mock(MessageRouter.class); ReactiveMessageSenderBuilderCustomizer customizer2 = builder -> builder.messageRouter(router); - ReactiveMessageSender sender = senderFactory.createReactiveMessageSender("topic0", schema, null, - Arrays.asList(customizer1, customizer2)); - sender.sendMessage(Mono.just(MessageSpec.of("test"))).block(Duration.ofSeconds(5)); - assertSenderHasTopicAndRouter("topic1", router); + + testCreateSender(null, null, "topic0", null, Arrays.asList(customizer1, customizer2), "topic1", router); } @Test - void createProducerWithNoTopic() { - ReactivePulsarSenderFactory senderFactory = new DefaultReactivePulsarSenderFactory<>(pulsarClient, - null); + void createSenderWithNoTopic() { + ReactivePulsarSenderFactory senderFactory = new DefaultReactivePulsarSenderFactory<>( + (PulsarClient) null, null, null); assertThatIllegalArgumentException().isThrownBy(() -> senderFactory.createReactiveMessageSender(null, schema)) .withMessageContaining("Topic must be specified when no default topic is configured"); } - protected void assertSenderHasTopicAndRouter(String topic, MessageRouter router) { - verify(producerBuilder).topic(topic); - if (router != null) { - verify(producerBuilder).messageRouter(router); - } - else { - verify(producerBuilder, never()).messageRouter(any()); - } + @Test + void createSenderWithCache() { + testCreateSender(null, AdaptedReactivePulsarClientFactory.createCache(), "topic1", null, null, "topic1", null); + } + + private void testCreateSender(ReactiveMessageSenderSpec spec, ReactiveMessageSenderCache cache, String topic, + MessageRouter router, List> customizers, + String expectedTopic, MessageRouter expectedRouter) { + ReactivePulsarSenderFactory senderFactory = new DefaultReactivePulsarSenderFactory<>( + (PulsarClient) null, spec, cache); + ReactiveMessageSender sender = senderFactory.createReactiveMessageSender(topic, schema, router, + customizers); + ObjectAssert objectAssert = assertThat(sender).extracting("senderSpec") + .asInstanceOf(InstanceOfAssertFactories.type(ReactiveMessageSenderSpec.class)); + objectAssert.extracting(ReactiveMessageSenderSpec::getTopicName).isEqualTo(expectedTopic); + objectAssert.extracting(ReactiveMessageSenderSpec::getMessageRouter).isSameAs(expectedRouter); + assertThat(sender).extracting("producerCache").isSameAs(cache); } } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/ReactivePulsarTemplateTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/ReactivePulsarTemplateTests.java index 9e691da2..3b477083 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/ReactivePulsarTemplateTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/ReactivePulsarTemplateTests.java @@ -66,7 +66,7 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { MutableReactiveMessageSenderSpec senderSpec = new MutableReactiveMessageSenderSpec(); senderSpec.setTopicName(topic); ReactivePulsarSenderFactory producerFactory = new DefaultReactivePulsarSenderFactory<>(client, - senderSpec); + senderSpec, null); ReactivePulsarSenderTemplate pulsarTemplate = new ReactivePulsarSenderTemplate<>(producerFactory); pulsarTemplate.setSchema(Schema.JSON(Foo.class)); Foo foo = new Foo("Foo-" + UUID.randomUUID(), "Bar-" + UUID.randomUUID()); @@ -113,7 +113,7 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { senderSpec.setTopicName(topic); } ReactivePulsarSenderFactory senderFactory = new DefaultReactivePulsarSenderFactory<>(client, - senderSpec); + senderSpec, null); ReactivePulsarSenderTemplate pulsarTemplate = new ReactivePulsarSenderTemplate<>(senderFactory); Mono sendResponse; if (testArgs.useSimpleApi) {