Remove Listener name from the Reader core classes
This commit is contained in:
@@ -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 <C> the {@link AbstractPulsarReaderListenerContainer} implementation type.
|
||||
* @param <C> the {@link AbstractPulsarMessageMessageReaderContainer} implementation type.
|
||||
* @param <T> Message payload type.
|
||||
* @author Soby Chacko
|
||||
*/
|
||||
public abstract class AbstractPulsarReaderContainerFactory<C extends AbstractPulsarReaderListenerContainer<T>, T>
|
||||
public abstract class AbstractPulsarReaderContainerFactory<C extends AbstractPulsarMessageMessageReaderContainer<T>, T>
|
||||
implements PulsarReaderContainerFactory, ApplicationEventPublisherAware, ApplicationContextAware {
|
||||
|
||||
protected final LogAccessor logger = new LogAccessor(this.getClass());
|
||||
@@ -99,7 +99,7 @@ public abstract class AbstractPulsarReaderContainerFactory<C extends AbstractPul
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@Override
|
||||
public C createReaderContainer(PulsarReaderEndpoint<PulsarReaderListenerContainer> endpoint) {
|
||||
public C createReaderContainer(PulsarReaderEndpoint<PulsarMessageReaderContainer> endpoint) {
|
||||
C instance = createContainerInstance(endpoint);
|
||||
JavaUtils.INSTANCE.acceptIfNotNull(endpoint.getId(), instance::setBeanName);
|
||||
if (endpoint instanceof AbstractPulsarReaderEndpoint) {
|
||||
@@ -112,13 +112,13 @@ public abstract class AbstractPulsarReaderContainerFactory<C extends AbstractPul
|
||||
return instance;
|
||||
}
|
||||
|
||||
protected abstract C createContainerInstance(PulsarReaderEndpoint<PulsarReaderListenerContainer> endpoint);
|
||||
protected abstract C createContainerInstance(PulsarReaderEndpoint<PulsarMessageReaderContainer> endpoint);
|
||||
|
||||
private void configureEndpoint(AbstractPulsarReaderEndpoint<C> aplEndpoint) {
|
||||
|
||||
}
|
||||
|
||||
protected void initializeContainer(C instance, PulsarReaderEndpoint<PulsarReaderListenerContainer> endpoint) {
|
||||
protected void initializeContainer(C instance, PulsarReaderEndpoint<PulsarMessageReaderContainer> endpoint) {
|
||||
PulsarReaderContainerProperties instanceProperties = instance.getContainerProperties();
|
||||
|
||||
if (instanceProperties.getSchema() == null) {
|
||||
|
||||
@@ -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<K>
|
||||
implements PulsarReaderEndpoint<PulsarReaderListenerContainer>, BeanFactoryAware, InitializingBean {
|
||||
implements PulsarReaderEndpoint<PulsarMessageReaderContainer>, BeanFactoryAware, InitializingBean {
|
||||
|
||||
private String subscriptionName;
|
||||
|
||||
@@ -146,14 +146,14 @@ public abstract class AbstractPulsarReaderEndpoint<K>
|
||||
}
|
||||
|
||||
@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<K> adapter = createReaderListener(container, messageConverter);
|
||||
@@ -163,7 +163,7 @@ public abstract class AbstractPulsarReaderEndpoint<K>
|
||||
}
|
||||
|
||||
protected abstract PulsarMessagingMessageListenerAdapter<K> createReaderListener(
|
||||
PulsarReaderListenerContainer container, @Nullable MessageConverter messageConverter);
|
||||
PulsarMessageReaderContainer container, @Nullable MessageConverter messageConverter);
|
||||
|
||||
public SchemaType getSchemaType() {
|
||||
return this.schemaType;
|
||||
|
||||
@@ -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<T>
|
||||
extends AbstractPulsarReaderContainerFactory<DefaultPulsarReaderListenerContainer<T>, T> {
|
||||
extends AbstractPulsarReaderContainerFactory<DefaultPulsarMessageReaderContainer<T>, T> {
|
||||
|
||||
public DefaultPulsarReaderContainerFactory(PulsarReaderFactory<? super T> readerFactory,
|
||||
PulsarReaderContainerProperties containerProperties) {
|
||||
@@ -38,8 +38,8 @@ public class DefaultPulsarReaderContainerFactory<T>
|
||||
}
|
||||
|
||||
@Override
|
||||
protected DefaultPulsarReaderListenerContainer<T> createContainerInstance(
|
||||
PulsarReaderEndpoint<PulsarReaderListenerContainer> endpoint) {
|
||||
protected DefaultPulsarMessageReaderContainer<T> createContainerInstance(
|
||||
PulsarReaderEndpoint<PulsarMessageReaderContainer> endpoint) {
|
||||
|
||||
PulsarReaderContainerProperties properties = new PulsarReaderContainerProperties();
|
||||
properties.setSchemaResolver(this.getContainerProperties().getSchemaResolver());
|
||||
@@ -55,17 +55,17 @@ public class DefaultPulsarReaderContainerFactory<T>
|
||||
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<T> instance,
|
||||
PulsarReaderEndpoint<PulsarReaderListenerContainer> endpoint) {
|
||||
protected void initializeContainer(DefaultPulsarMessageReaderContainer<T> instance,
|
||||
PulsarReaderEndpoint<PulsarMessageReaderContainer> endpoint) {
|
||||
super.initializeContainer(instance, endpoint);
|
||||
}
|
||||
|
||||
@Override
|
||||
public DefaultPulsarReaderListenerContainer<T> createReaderContainer(String... topics) {
|
||||
public DefaultPulsarMessageReaderContainer<T> createReaderContainer(String... topics) {
|
||||
// TODO
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -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 <E> endpoint type
|
||||
* @author Soby Chacko
|
||||
*/
|
||||
public class GenericReaderEndpointRegistry<C extends PulsarReaderListenerContainer, E extends PulsarReaderEndpoint<C>>
|
||||
public class GenericReaderEndpointRegistry<C extends PulsarMessageReaderContainer, E extends PulsarReaderEndpoint<C>>
|
||||
implements PulsarReaderContainerRegistry, DisposableBean, SmartLifecycle, ApplicationContextAware,
|
||||
ApplicationListener<ContextRefreshedEvent> {
|
||||
|
||||
|
||||
@@ -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<V> extends AbstractPulsarReaderEndpoint<
|
||||
}
|
||||
|
||||
@Override
|
||||
protected PulsarMessagingMessageListenerAdapter<V> createReaderListener(PulsarReaderListenerContainer container,
|
||||
protected PulsarMessagingMessageListenerAdapter<V> createReaderListener(PulsarMessageReaderContainer container,
|
||||
@Nullable MessageConverter messageConverter) {
|
||||
PulsarMessagingMessageListenerAdapter<V> readerListener = createMessageListenerInstance(messageConverter);
|
||||
HandlerAdapter handlerMethod = configureListenerAdapter(readerListener);
|
||||
@@ -112,7 +112,7 @@ public class MethodPulsarReaderEndpoint<V> 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();
|
||||
|
||||
@@ -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<PulsarReaderListenerContainer, PulsarReaderEndpoint<PulsarReaderListenerContainer>> {
|
||||
ReaderContainerFactory<PulsarMessageReaderContainer, PulsarReaderEndpoint<PulsarMessageReaderContainer>> {
|
||||
|
||||
}
|
||||
|
||||
@@ -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 <C> reader listener container type.
|
||||
* @author Soby Chacko
|
||||
*/
|
||||
public interface PulsarReaderEndpoint<C extends PulsarReaderListenerContainer> {
|
||||
public interface PulsarReaderEndpoint<C extends PulsarMessageReaderContainer> {
|
||||
|
||||
/**
|
||||
* Return the id of this endpoint.
|
||||
|
||||
@@ -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.
|
||||
*
|
||||
* <p>
|
||||
* 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<PulsarReaderListenerContainer, PulsarReaderEndpoint<PulsarReaderListenerContainer>> {
|
||||
GenericReaderEndpointRegistry<PulsarMessageReaderContainer, PulsarReaderEndpoint<PulsarMessageReaderContainer>> {
|
||||
|
||||
public PulsarReaderEndpointRegistry() {
|
||||
super(PulsarReaderListenerContainer.class);
|
||||
super(PulsarMessageReaderContainer.class);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -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 <C> Container type
|
||||
* @param <E> Endpoint type
|
||||
* @author Soby Chacko
|
||||
*/
|
||||
public interface ReaderContainerFactory<C extends PulsarReaderListenerContainer, E extends PulsarReaderEndpoint<C>> {
|
||||
public interface ReaderContainerFactory<C extends PulsarMessageReaderContainer, E extends PulsarReaderEndpoint<C>> {
|
||||
|
||||
C createReaderContainer(E endpoint);
|
||||
|
||||
|
||||
@@ -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 <T> reader data type.
|
||||
* @author Soby Chacko
|
||||
*/
|
||||
public non-sealed abstract class AbstractPulsarReaderListenerContainer<T> implements PulsarReaderListenerContainer,
|
||||
public non-sealed abstract class AbstractPulsarMessageMessageReaderContainer<T> implements PulsarMessageReaderContainer,
|
||||
BeanNameAware, ApplicationEventPublisherAware, ApplicationContextAware {
|
||||
|
||||
protected final LogAccessor logger = new LogAccessor(this.getClass());
|
||||
@@ -59,7 +59,7 @@ public non-sealed abstract class AbstractPulsarReaderListenerContainer<T> implem
|
||||
private volatile boolean running = false;
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
protected AbstractPulsarReaderListenerContainer(PulsarReaderFactory<? super T> pulsarReaderFactory,
|
||||
protected AbstractPulsarMessageMessageReaderContainer(PulsarReaderFactory<? super T> pulsarReaderFactory,
|
||||
PulsarReaderContainerProperties pulsarReaderContainerProperties) {
|
||||
this.pulsarReaderFactory = (PulsarReaderFactory<T>) pulsarReaderFactory;
|
||||
this.pulsarReaderContainerProperties = pulsarReaderContainerProperties;
|
||||
@@ -46,7 +46,7 @@ import org.springframework.scheduling.SchedulingAwareRunnable;
|
||||
* @param <T> reader data type.
|
||||
* @author Soby Chacko
|
||||
*/
|
||||
public class DefaultPulsarReaderListenerContainer<T> extends AbstractPulsarReaderListenerContainer<T> {
|
||||
public class DefaultPulsarMessageReaderContainer<T> extends AbstractPulsarMessageMessageReaderContainer<T> {
|
||||
|
||||
private final AtomicReference<InternalAsyncReader> internalAsyncReader = new AtomicReference<>();
|
||||
|
||||
@@ -54,11 +54,11 @@ public class DefaultPulsarReaderListenerContainer<T> extends AbstractPulsarReade
|
||||
|
||||
private volatile CompletableFuture<?> readerFuture;
|
||||
|
||||
private final AbstractPulsarReaderListenerContainer<?> thisOrParentContainer;
|
||||
private final AbstractPulsarMessageMessageReaderContainer<?> thisOrParentContainer;
|
||||
|
||||
private final AtomicReference<Thread> readerThread = new AtomicReference<>();
|
||||
|
||||
public DefaultPulsarReaderListenerContainer(PulsarReaderFactory<? super T> pulsarReaderFactory,
|
||||
public DefaultPulsarMessageReaderContainer(PulsarReaderFactory<? super T> pulsarReaderFactory,
|
||||
PulsarReaderContainerProperties pulsarReaderContainerProperties) {
|
||||
super(pulsarReaderFactory, pulsarReaderContainerProperties);
|
||||
this.thisOrParentContainer = this;
|
||||
@@ -161,7 +161,7 @@ public class DefaultPulsarReaderListenerContainer<T> 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<T> 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.");
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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);
|
||||
|
||||
@@ -29,12 +29,12 @@ import org.springframework.lang.Nullable;
|
||||
public interface PulsarReaderContainerRegistry {
|
||||
|
||||
@Nullable
|
||||
PulsarReaderListenerContainer getReaderContainer(String id);
|
||||
PulsarMessageReaderContainer getReaderContainer(String id);
|
||||
|
||||
Set<String> getReaderContainerIds();
|
||||
|
||||
Collection<? extends PulsarReaderListenerContainer> getReaderContainers();
|
||||
Collection<? extends PulsarMessageReaderContainer> getReaderContainers();
|
||||
|
||||
Collection<? extends PulsarReaderListenerContainer> getAllReaderContainers();
|
||||
Collection<? extends PulsarMessageReaderContainer> getAllReaderContainers();
|
||||
|
||||
}
|
||||
|
||||
@@ -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<String> container = null;
|
||||
DefaultPulsarMessageReaderContainer<String> container = null;
|
||||
try {
|
||||
container = new DefaultPulsarReaderListenerContainer<>(pulsarReaderFactory, readerContainerProperties);
|
||||
container = new DefaultPulsarMessageReaderContainer<>(pulsarReaderFactory, readerContainerProperties);
|
||||
container.start();
|
||||
|
||||
Map<String, Object> 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<String> container = null;
|
||||
DefaultPulsarMessageReaderContainer<String> container = null;
|
||||
try {
|
||||
container = new DefaultPulsarReaderListenerContainer<>(pulsarReaderFactory, containerProps);
|
||||
container = new DefaultPulsarMessageReaderContainer<>(pulsarReaderFactory, containerProps);
|
||||
container.start();
|
||||
|
||||
Map<String, Object> prodConfig = Map.of("topicName", "dprlct-002");
|
||||
@@ -140,9 +140,9 @@ public class DefaultPulsarReaderListenerContainerTests implements PulsarTestCont
|
||||
|
||||
var readerConfig = Collections.<String, Object>emptyMap();
|
||||
var readerFactory = new DefaultPulsarReaderFactory<String>(pulsarClient, readerConfig);
|
||||
DefaultPulsarReaderListenerContainer<String> container = null;
|
||||
DefaultPulsarMessageReaderContainer<String> container = null;
|
||||
try {
|
||||
container = new DefaultPulsarReaderListenerContainer<>(readerFactory, containerProps);
|
||||
container = new DefaultPulsarMessageReaderContainer<>(readerFactory, containerProps);
|
||||
|
||||
var prodConfig = Map.<String, Object>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();
|
||||
}
|
||||
Reference in New Issue
Block a user