diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarReaderContainerFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarReaderContainerFactory.java index 90bfa68b..37f1605f 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarReaderContainerFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarReaderContainerFactory.java @@ -25,20 +25,20 @@ import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.ApplicationEventPublisherAware; import org.springframework.core.log.LogAccessor; import org.springframework.pulsar.core.PulsarReaderFactory; -import org.springframework.pulsar.reader.AbstractPulsarReaderListenerContainer; +import org.springframework.pulsar.reader.AbstractPulsarMessageMessageReaderContainer; +import org.springframework.pulsar.reader.PulsarMessageReaderContainer; import org.springframework.pulsar.reader.PulsarReaderContainerProperties; -import org.springframework.pulsar.reader.PulsarReaderListenerContainer; import org.springframework.pulsar.support.JavaUtils; import org.springframework.pulsar.support.MessageConverter; /** * Base {@link PulsarReaderContainerFactory} implementation. * - * @param the {@link AbstractPulsarReaderListenerContainer} implementation type. + * @param the {@link AbstractPulsarMessageMessageReaderContainer} implementation type. * @param Message payload type. * @author Soby Chacko */ -public abstract class AbstractPulsarReaderContainerFactory, T> +public abstract class AbstractPulsarReaderContainerFactory, T> implements PulsarReaderContainerFactory, ApplicationEventPublisherAware, ApplicationContextAware { protected final LogAccessor logger = new LogAccessor(this.getClass()); @@ -99,7 +99,7 @@ public abstract class AbstractPulsarReaderContainerFactory endpoint) { + public C createReaderContainer(PulsarReaderEndpoint endpoint) { C instance = createContainerInstance(endpoint); JavaUtils.INSTANCE.acceptIfNotNull(endpoint.getId(), instance::setBeanName); if (endpoint instanceof AbstractPulsarReaderEndpoint) { @@ -112,13 +112,13 @@ public abstract class AbstractPulsarReaderContainerFactory endpoint); + protected abstract C createContainerInstance(PulsarReaderEndpoint endpoint); private void configureEndpoint(AbstractPulsarReaderEndpoint aplEndpoint) { } - protected void initializeContainer(C instance, PulsarReaderEndpoint endpoint) { + protected void initializeContainer(C instance, PulsarReaderEndpoint endpoint) { PulsarReaderContainerProperties instanceProperties = instance.getContainerProperties(); if (instanceProperties.getSchema() == null) { diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarReaderEndpoint.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarReaderEndpoint.java index 4cea0ef9..9bfc7e15 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarReaderEndpoint.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarReaderEndpoint.java @@ -35,7 +35,7 @@ import org.springframework.context.expression.BeanFactoryResolver; import org.springframework.expression.BeanResolver; import org.springframework.lang.Nullable; import org.springframework.pulsar.listener.adapter.PulsarMessagingMessageListenerAdapter; -import org.springframework.pulsar.reader.PulsarReaderListenerContainer; +import org.springframework.pulsar.reader.PulsarMessageReaderContainer; import org.springframework.pulsar.support.MessageConverter; import org.springframework.util.Assert; @@ -46,7 +46,7 @@ import org.springframework.util.Assert; * @author Soby Chacko */ public abstract class AbstractPulsarReaderEndpoint - implements PulsarReaderEndpoint, BeanFactoryAware, InitializingBean { + implements PulsarReaderEndpoint, BeanFactoryAware, InitializingBean { private String subscriptionName; @@ -146,14 +146,14 @@ public abstract class AbstractPulsarReaderEndpoint } @Override - public void setupListenerContainer(PulsarReaderListenerContainer listenerContainer, + public void setupListenerContainer(PulsarMessageReaderContainer listenerContainer, @Nullable MessageConverter messageConverter) { setupMessageListener(listenerContainer, messageConverter); } @SuppressWarnings("unchecked") - private void setupMessageListener(PulsarReaderListenerContainer container, + private void setupMessageListener(PulsarMessageReaderContainer container, @Nullable MessageConverter messageConverter) { PulsarMessagingMessageListenerAdapter adapter = createReaderListener(container, messageConverter); @@ -163,7 +163,7 @@ public abstract class AbstractPulsarReaderEndpoint } protected abstract PulsarMessagingMessageListenerAdapter createReaderListener( - PulsarReaderListenerContainer container, @Nullable MessageConverter messageConverter); + PulsarMessageReaderContainer container, @Nullable MessageConverter messageConverter); public SchemaType getSchemaType() { return this.schemaType; diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/DefaultPulsarReaderContainerFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/DefaultPulsarReaderContainerFactory.java index 4871aebe..99eb94c2 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/DefaultPulsarReaderContainerFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/DefaultPulsarReaderContainerFactory.java @@ -17,9 +17,9 @@ package org.springframework.pulsar.config; import org.springframework.pulsar.core.PulsarReaderFactory; -import org.springframework.pulsar.reader.DefaultPulsarReaderListenerContainer; +import org.springframework.pulsar.reader.DefaultPulsarMessageReaderContainer; +import org.springframework.pulsar.reader.PulsarMessageReaderContainer; import org.springframework.pulsar.reader.PulsarReaderContainerProperties; -import org.springframework.pulsar.reader.PulsarReaderListenerContainer; import org.springframework.util.CollectionUtils; import org.springframework.util.StringUtils; @@ -30,7 +30,7 @@ import org.springframework.util.StringUtils; * @author Soby Chacko */ public class DefaultPulsarReaderContainerFactory - extends AbstractPulsarReaderContainerFactory, T> { + extends AbstractPulsarReaderContainerFactory, T> { public DefaultPulsarReaderContainerFactory(PulsarReaderFactory readerFactory, PulsarReaderContainerProperties containerProperties) { @@ -38,8 +38,8 @@ public class DefaultPulsarReaderContainerFactory } @Override - protected DefaultPulsarReaderListenerContainer createContainerInstance( - PulsarReaderEndpoint endpoint) { + protected DefaultPulsarMessageReaderContainer createContainerInstance( + PulsarReaderEndpoint endpoint) { PulsarReaderContainerProperties properties = new PulsarReaderContainerProperties(); properties.setSchemaResolver(this.getContainerProperties().getSchemaResolver()); @@ -55,17 +55,17 @@ public class DefaultPulsarReaderContainerFactory properties.setSchemaType(endpoint.getSchemaType()); properties.setStartMessageId(endpoint.getStartMessageId()); - return new DefaultPulsarReaderListenerContainer<>(this.getReaderFactory(), properties); + return new DefaultPulsarMessageReaderContainer<>(this.getReaderFactory(), properties); } @Override - protected void initializeContainer(DefaultPulsarReaderListenerContainer instance, - PulsarReaderEndpoint endpoint) { + protected void initializeContainer(DefaultPulsarMessageReaderContainer instance, + PulsarReaderEndpoint endpoint) { super.initializeContainer(instance, endpoint); } @Override - public DefaultPulsarReaderListenerContainer createReaderContainer(String... topics) { + public DefaultPulsarMessageReaderContainer createReaderContainer(String... topics) { // TODO return null; } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/GenericReaderEndpointRegistry.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/GenericReaderEndpointRegistry.java index 8ed805e1..19a1c4e6 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/GenericReaderEndpointRegistry.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/GenericReaderEndpointRegistry.java @@ -36,8 +36,8 @@ import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.SmartLifecycle; import org.springframework.context.event.ContextRefreshedEvent; import org.springframework.lang.Nullable; +import org.springframework.pulsar.reader.PulsarMessageReaderContainer; import org.springframework.pulsar.reader.PulsarReaderContainerRegistry; -import org.springframework.pulsar.reader.PulsarReaderListenerContainer; import org.springframework.util.Assert; /** @@ -57,7 +57,7 @@ import org.springframework.util.Assert; * @param endpoint type * @author Soby Chacko */ -public class GenericReaderEndpointRegistry> +public class GenericReaderEndpointRegistry> implements PulsarReaderContainerRegistry, DisposableBean, SmartLifecycle, ApplicationContextAware, ApplicationListener { diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarReaderEndpoint.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarReaderEndpoint.java index 9a25ab32..f5ae0059 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarReaderEndpoint.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarReaderEndpoint.java @@ -39,9 +39,9 @@ import org.springframework.pulsar.listener.Acknowledgement; import org.springframework.pulsar.listener.adapter.HandlerAdapter; import org.springframework.pulsar.listener.adapter.PulsarMessagingMessageListenerAdapter; import org.springframework.pulsar.listener.adapter.PulsarRecordMessagingReaderListenerAdapter; -import org.springframework.pulsar.reader.DefaultPulsarReaderListenerContainer; +import org.springframework.pulsar.reader.DefaultPulsarMessageReaderContainer; +import org.springframework.pulsar.reader.PulsarMessageReaderContainer; import org.springframework.pulsar.reader.PulsarReaderContainerProperties; -import org.springframework.pulsar.reader.PulsarReaderListenerContainer; import org.springframework.pulsar.support.MessageConverter; import org.springframework.pulsar.support.converter.PulsarRecordMessageConverter; import org.springframework.util.Assert; @@ -84,7 +84,7 @@ public class MethodPulsarReaderEndpoint extends AbstractPulsarReaderEndpoint< } @Override - protected PulsarMessagingMessageListenerAdapter createReaderListener(PulsarReaderListenerContainer container, + protected PulsarMessagingMessageListenerAdapter createReaderListener(PulsarMessageReaderContainer container, @Nullable MessageConverter messageConverter) { PulsarMessagingMessageListenerAdapter readerListener = createMessageListenerInstance(messageConverter); HandlerAdapter handlerMethod = configureListenerAdapter(readerListener); @@ -112,7 +112,7 @@ public class MethodPulsarReaderEndpoint extends AbstractPulsarReaderEndpoint< messageParameter = parameter.get(); } - DefaultPulsarReaderListenerContainer containerInstance = (DefaultPulsarReaderListenerContainer) container; + DefaultPulsarMessageReaderContainer containerInstance = (DefaultPulsarMessageReaderContainer) container; PulsarReaderContainerProperties pulsarContainerProperties = containerInstance.getContainerProperties(); SchemaResolver schemaResolver = pulsarContainerProperties.getSchemaResolver(); SchemaType schemaType = pulsarContainerProperties.getSchemaType(); diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarReaderContainerFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarReaderContainerFactory.java index c34a2af7..dd5c0618 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarReaderContainerFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarReaderContainerFactory.java @@ -16,14 +16,14 @@ package org.springframework.pulsar.config; -import org.springframework.pulsar.reader.PulsarReaderListenerContainer; +import org.springframework.pulsar.reader.PulsarMessageReaderContainer; /** - * Container factory for {@link PulsarReaderListenerContainer}. + * Container factory for {@link PulsarMessageReaderContainer}. * * @author Soby Chacko */ public interface PulsarReaderContainerFactory extends - ReaderContainerFactory> { + ReaderContainerFactory> { } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarReaderEndpoint.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarReaderEndpoint.java index fd6018f3..d89fe142 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarReaderEndpoint.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarReaderEndpoint.java @@ -22,7 +22,7 @@ import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.common.schema.SchemaType; import org.springframework.lang.Nullable; -import org.springframework.pulsar.reader.PulsarReaderListenerContainer; +import org.springframework.pulsar.reader.PulsarMessageReaderContainer; import org.springframework.pulsar.support.MessageConverter; /** @@ -33,7 +33,7 @@ import org.springframework.pulsar.support.MessageConverter; * @param reader listener container type. * @author Soby Chacko */ -public interface PulsarReaderEndpoint { +public interface PulsarReaderEndpoint { /** * Return the id of this endpoint. diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarReaderEndpointRegistry.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarReaderEndpointRegistry.java index b308734d..37c69e68 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarReaderEndpointRegistry.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarReaderEndpointRegistry.java @@ -16,28 +16,28 @@ package org.springframework.pulsar.config; -import org.springframework.pulsar.reader.PulsarReaderListenerContainer; +import org.springframework.pulsar.reader.PulsarMessageReaderContainer; /** - * Creates the necessary {@link PulsarReaderListenerContainer} instances for the - * registered {@linkplain PulsarReaderEndpoint endpoints}. Also manages the lifecycle of - * the listener containers, in particular within the lifecycle of the application context. + * Creates the necessary {@link PulsarMessageReaderContainer} instances for the registered + * {@linkplain PulsarReaderEndpoint endpoints}. Also manages the lifecycle of the listener + * containers, in particular within the lifecycle of the application context. * *

- * Contrary to {@link PulsarReaderListenerContainer}s created manually, listener - * containers managed by registry are not beans in the application context and are not - * candidates for autowiring. Use {@link #getReaderContainer(String)} ()} if you need to - * access this registry's listener containers for management purposes. If you need to - * access to a specific message listener container, use - * {@link #getReaderContainer(String)} with the id of the endpoint. + * Contrary to {@link PulsarMessageReaderContainer}s created manually, listener containers + * managed by registry are not beans in the application context and are not candidates for + * autowiring. Use {@link #getReaderContainer(String)} ()} if you need to access this + * registry's listener containers for management purposes. If you need to access to a + * specific message listener container, use {@link #getReaderContainer(String)} with the + * id of the endpoint. * * @author Soby Chacko */ public class PulsarReaderEndpointRegistry extends - GenericReaderEndpointRegistry> { + GenericReaderEndpointRegistry> { public PulsarReaderEndpointRegistry() { - super(PulsarReaderListenerContainer.class); + super(PulsarMessageReaderContainer.class); } } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/ReaderContainerFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/ReaderContainerFactory.java index 23391208..8304dc49 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/ReaderContainerFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/ReaderContainerFactory.java @@ -16,16 +16,16 @@ package org.springframework.pulsar.config; -import org.springframework.pulsar.reader.PulsarReaderListenerContainer; +import org.springframework.pulsar.reader.PulsarMessageReaderContainer; /** - * Base container factory interface for {@link PulsarReaderListenerContainer}. + * Base container factory interface for {@link PulsarMessageReaderContainer}. * * @param Container type * @param Endpoint type * @author Soby Chacko */ -public interface ReaderContainerFactory> { +public interface ReaderContainerFactory> { C createReaderContainer(E endpoint); diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/AbstractPulsarReaderListenerContainer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/AbstractPulsarMessageMessageReaderContainer.java similarity index 93% rename from spring-pulsar/src/main/java/org/springframework/pulsar/reader/AbstractPulsarReaderListenerContainer.java rename to spring-pulsar/src/main/java/org/springframework/pulsar/reader/AbstractPulsarMessageMessageReaderContainer.java index ab163f7d..45828a26 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/AbstractPulsarReaderListenerContainer.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/AbstractPulsarMessageMessageReaderContainer.java @@ -30,12 +30,12 @@ import org.springframework.pulsar.core.PulsarReaderFactory; import org.springframework.util.Assert; /** - * Core implementation for {@link PulsarReaderListenerContainer}. + * Core implementation for {@link PulsarMessageReaderContainer}. * * @param reader data type. * @author Soby Chacko */ -public non-sealed abstract class AbstractPulsarReaderListenerContainer implements PulsarReaderListenerContainer, +public non-sealed abstract class AbstractPulsarMessageMessageReaderContainer implements PulsarMessageReaderContainer, BeanNameAware, ApplicationEventPublisherAware, ApplicationContextAware { protected final LogAccessor logger = new LogAccessor(this.getClass()); @@ -59,7 +59,7 @@ public non-sealed abstract class AbstractPulsarReaderListenerContainer implem private volatile boolean running = false; @SuppressWarnings("unchecked") - protected AbstractPulsarReaderListenerContainer(PulsarReaderFactory pulsarReaderFactory, + protected AbstractPulsarMessageMessageReaderContainer(PulsarReaderFactory pulsarReaderFactory, PulsarReaderContainerProperties pulsarReaderContainerProperties) { this.pulsarReaderFactory = (PulsarReaderFactory) pulsarReaderFactory; this.pulsarReaderContainerProperties = pulsarReaderContainerProperties; diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/DefaultPulsarReaderListenerContainer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/DefaultPulsarMessageReaderContainer.java similarity index 92% rename from spring-pulsar/src/main/java/org/springframework/pulsar/reader/DefaultPulsarReaderListenerContainer.java rename to spring-pulsar/src/main/java/org/springframework/pulsar/reader/DefaultPulsarMessageReaderContainer.java index 10084713..2349cab2 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/DefaultPulsarReaderListenerContainer.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/DefaultPulsarMessageReaderContainer.java @@ -46,7 +46,7 @@ import org.springframework.scheduling.SchedulingAwareRunnable; * @param reader data type. * @author Soby Chacko */ -public class DefaultPulsarReaderListenerContainer extends AbstractPulsarReaderListenerContainer { +public class DefaultPulsarMessageReaderContainer extends AbstractPulsarMessageMessageReaderContainer { private final AtomicReference internalAsyncReader = new AtomicReference<>(); @@ -54,11 +54,11 @@ public class DefaultPulsarReaderListenerContainer extends AbstractPulsarReade private volatile CompletableFuture readerFuture; - private final AbstractPulsarReaderListenerContainer thisOrParentContainer; + private final AbstractPulsarMessageMessageReaderContainer thisOrParentContainer; private final AtomicReference readerThread = new AtomicReference<>(); - public DefaultPulsarReaderListenerContainer(PulsarReaderFactory pulsarReaderFactory, + public DefaultPulsarMessageReaderContainer(PulsarReaderFactory pulsarReaderFactory, PulsarReaderContainerProperties pulsarReaderContainerProperties) { super(pulsarReaderFactory, pulsarReaderContainerProperties); this.thisOrParentContainer = this; @@ -161,7 +161,7 @@ public class DefaultPulsarReaderListenerContainer extends AbstractPulsarReade @Override public void run() { - DefaultPulsarReaderListenerContainer.this.readerThread.set(Thread.currentThread()); + DefaultPulsarMessageReaderContainer.this.readerThread.set(Thread.currentThread()); publishReaderStartingEvent(); publishReaderStartedEvent(); @@ -171,7 +171,7 @@ public class DefaultPulsarReaderListenerContainer extends AbstractPulsarReade this.listener.received(this.reader, message); } catch (PulsarClientException e) { - DefaultPulsarReaderListenerContainer.this.logger.error(e, () -> "Error receiving messages."); + DefaultPulsarMessageReaderContainer.this.logger.error(e, () -> "Error receiving messages."); } } } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/PulsarReaderListenerContainer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/PulsarMessageReaderContainer.java similarity index 89% rename from spring-pulsar/src/main/java/org/springframework/pulsar/reader/PulsarReaderListenerContainer.java rename to spring-pulsar/src/main/java/org/springframework/pulsar/reader/PulsarMessageReaderContainer.java index a8021e75..62bb31b7 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/PulsarReaderListenerContainer.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/PulsarMessageReaderContainer.java @@ -25,8 +25,8 @@ import org.springframework.context.SmartLifecycle; * * @author Soby Chacko */ -public sealed interface PulsarReaderListenerContainer - extends SmartLifecycle, DisposableBean permits AbstractPulsarReaderListenerContainer { +public sealed interface PulsarMessageReaderContainer + extends SmartLifecycle, DisposableBean permits AbstractPulsarMessageMessageReaderContainer { void setupReaderListener(Object messageListener); diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/PulsarReaderContainerRegistry.java b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/PulsarReaderContainerRegistry.java index c0effa43..72227820 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/PulsarReaderContainerRegistry.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/PulsarReaderContainerRegistry.java @@ -29,12 +29,12 @@ import org.springframework.lang.Nullable; public interface PulsarReaderContainerRegistry { @Nullable - PulsarReaderListenerContainer getReaderContainer(String id); + PulsarMessageReaderContainer getReaderContainer(String id); Set getReaderContainerIds(); - Collection getReaderContainers(); + Collection getReaderContainers(); - Collection getAllReaderContainers(); + Collection getAllReaderContainers(); } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/reader/DefaultPulsarReaderListenerContainerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/reader/DefaultPulsarMessageReaderContainerTests.java similarity index 89% rename from spring-pulsar/src/test/java/org/springframework/pulsar/reader/DefaultPulsarReaderListenerContainerTests.java rename to spring-pulsar/src/test/java/org/springframework/pulsar/reader/DefaultPulsarMessageReaderContainerTests.java index 09ecc65e..3c878c40 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/reader/DefaultPulsarReaderListenerContainerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/reader/DefaultPulsarMessageReaderContainerTests.java @@ -41,11 +41,11 @@ import org.springframework.pulsar.core.PulsarTemplate; import org.springframework.pulsar.test.support.PulsarTestContainerSupport; /** - * Basic tests for {@link DefaultPulsarReaderListenerContainer}. + * Basic tests for {@link DefaultPulsarMessageReaderContainer}. * * @author Soby Chacko */ -public class DefaultPulsarReaderListenerContainerTests implements PulsarTestContainerSupport { +public class DefaultPulsarMessageReaderContainerTests implements PulsarTestContainerSupport { private final LogAccessor logger = new LogAccessor(this.getClass()); @@ -78,9 +78,9 @@ public class DefaultPulsarReaderListenerContainerTests implements PulsarTestCont readerContainerProperties.setStartMessageId(MessageId.earliest); readerContainerProperties.setSchema(Schema.STRING); - DefaultPulsarReaderListenerContainer container = null; + DefaultPulsarMessageReaderContainer container = null; try { - container = new DefaultPulsarReaderListenerContainer<>(pulsarReaderFactory, readerContainerProperties); + container = new DefaultPulsarMessageReaderContainer<>(pulsarReaderFactory, readerContainerProperties); container.start(); Map prodConfig = Map.of("topicName", "dprlct-001"); @@ -109,9 +109,9 @@ public class DefaultPulsarReaderListenerContainerTests implements PulsarTestCont containerProps.setStartMessageId(MessageId.earliest); containerProps.setTopics(List.of("dprlct-002")); containerProps.setSchema(Schema.STRING); - DefaultPulsarReaderListenerContainer container = null; + DefaultPulsarMessageReaderContainer container = null; try { - container = new DefaultPulsarReaderListenerContainer<>(pulsarReaderFactory, containerProps); + container = new DefaultPulsarMessageReaderContainer<>(pulsarReaderFactory, containerProps); container.start(); Map prodConfig = Map.of("topicName", "dprlct-002"); @@ -140,9 +140,9 @@ public class DefaultPulsarReaderListenerContainerTests implements PulsarTestCont var readerConfig = Collections.emptyMap(); var readerFactory = new DefaultPulsarReaderFactory(pulsarClient, readerConfig); - DefaultPulsarReaderListenerContainer container = null; + DefaultPulsarMessageReaderContainer container = null; try { - container = new DefaultPulsarReaderListenerContainer<>(readerFactory, containerProps); + container = new DefaultPulsarMessageReaderContainer<>(readerFactory, containerProps); var prodConfig = Map.of("topicName", "dprlct-003"); var producerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, prodConfig); @@ -165,7 +165,7 @@ public class DefaultPulsarReaderListenerContainerTests implements PulsarTestCont } } - private void safeStopContainer(PulsarReaderListenerContainer container) { + private void safeStopContainer(PulsarMessageReaderContainer container) { try { container.stop(); }