Remove final keyword in methods (#213)
This commit is contained in:
committed by
GitHub
parent
95a6c4c655
commit
27bba9816e
@@ -244,7 +244,7 @@ public class PulsarListenerAnnotationBeanPostProcessor<V>
|
||||
}
|
||||
|
||||
@Override
|
||||
public Object postProcessAfterInitialization(final Object bean, final String beanName) throws BeansException {
|
||||
public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException {
|
||||
if (!this.nonAnnotatedClasses.contains(bean.getClass())) {
|
||||
Class<?> targetClass = AopUtils.getTargetClass(bean);
|
||||
Map<Method, Set<PulsarListener>> annotatedMethods = MethodIntrospector.selectMethods(targetClass,
|
||||
|
||||
@@ -242,7 +242,7 @@ public class ReactivePulsarListenerAnnotationBeanPostProcessor<V>
|
||||
}
|
||||
|
||||
@Override
|
||||
public Object postProcessAfterInitialization(final Object bean, final String beanName) throws BeansException {
|
||||
public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException {
|
||||
if (!this.nonAnnotatedClasses.contains(bean.getClass())) {
|
||||
Class<?> targetClass = AopUtils.getTargetClass(bean);
|
||||
Map<Method, Set<ReactivePulsarListener>> annotatedMethods = MethodIntrospector.selectMethods(targetClass,
|
||||
|
||||
@@ -120,7 +120,7 @@ public class MethodPulsarListenerEndpoint<V> extends AbstractPulsarListenerEndpo
|
||||
Assert.state(this.messageHandlerMethodFactory != null,
|
||||
"Could not create message listener - MessageHandlerMethodFactory not set");
|
||||
PulsarMessagingMessageListenerAdapter<V> messageListener = createMessageListenerInstance(messageConverter);
|
||||
final HandlerAdapter handlerMethod = configureListenerAdapter(messageListener);
|
||||
HandlerAdapter handlerMethod = configureListenerAdapter(messageListener);
|
||||
messageListener.setHandlerMethod(handlerMethod);
|
||||
|
||||
// Since we have access to the handler method here, check if we can type infer the
|
||||
@@ -128,14 +128,14 @@ public class MethodPulsarListenerEndpoint<V> extends AbstractPulsarListenerEndpo
|
||||
|
||||
// TODO: filter out the payload type by excluding Consumer, Message, Messages etc.
|
||||
|
||||
final MethodParameter[] methodParameters = handlerMethod.getInvokerHandlerMethod().getMethodParameters();
|
||||
MethodParameter[] methodParameters = handlerMethod.getInvokerHandlerMethod().getMethodParameters();
|
||||
MethodParameter messageParameter = null;
|
||||
final Optional<MethodParameter> parameter = Arrays.stream(methodParameters)
|
||||
Optional<MethodParameter> parameter = Arrays.stream(methodParameters)
|
||||
.filter(methodParameter1 -> !methodParameter1.getParameterType().equals(Consumer.class)
|
||||
|| !methodParameter1.getParameterType().equals(Acknowledgement.class)
|
||||
|| !methodParameter1.hasParameterAnnotation(Header.class))
|
||||
.findFirst();
|
||||
final long count = Arrays.stream(methodParameters)
|
||||
long count = Arrays.stream(methodParameters)
|
||||
.filter(methodParameter1 -> !methodParameter1.getParameterType().equals(Consumer.class)
|
||||
&& !methodParameter1.getParameterType().equals(Acknowledgement.class)
|
||||
&& !methodParameter1.hasParameterAnnotation(Header.class))
|
||||
@@ -145,9 +145,9 @@ public class MethodPulsarListenerEndpoint<V> extends AbstractPulsarListenerEndpo
|
||||
messageParameter = parameter.get();
|
||||
}
|
||||
|
||||
final ConcurrentPulsarMessageListenerContainer<?> containerInstance = (ConcurrentPulsarMessageListenerContainer<?>) container;
|
||||
final PulsarContainerProperties pulsarContainerProperties = containerInstance.getContainerProperties();
|
||||
final SchemaType schemaType = pulsarContainerProperties.getSchemaType();
|
||||
ConcurrentPulsarMessageListenerContainer<?> containerInstance = (ConcurrentPulsarMessageListenerContainer<?>) container;
|
||||
PulsarContainerProperties pulsarContainerProperties = containerInstance.getContainerProperties();
|
||||
SchemaType schemaType = pulsarContainerProperties.getSchemaType();
|
||||
if (schemaType != SchemaType.NONE) {
|
||||
switch (schemaType) {
|
||||
case STRING -> pulsarContainerProperties.setSchema(Schema.STRING);
|
||||
@@ -193,7 +193,7 @@ public class MethodPulsarListenerEndpoint<V> extends AbstractPulsarListenerEndpo
|
||||
}
|
||||
}
|
||||
}
|
||||
final SchemaType type = pulsarContainerProperties.getSchema().getSchemaInfo().getType();
|
||||
SchemaType type = pulsarContainerProperties.getSchema().getSchemaInfo().getType();
|
||||
pulsarContainerProperties.setSchemaType(type);
|
||||
|
||||
container.setNegativeAckRedeliveryBackoff(this.negativeAckRedeliveryBackoff);
|
||||
@@ -206,7 +206,7 @@ public class MethodPulsarListenerEndpoint<V> extends AbstractPulsarListenerEndpo
|
||||
|
||||
private Schema<?> getMessageSchema(MethodParameter messageParameter, Function<Class<?>, Schema<?>> schemaFactory) {
|
||||
ResolvableType messageType = resolvableType(messageParameter);
|
||||
final Class<?> messageClass = messageType.getRawClass();
|
||||
Class<?> messageClass = messageType.getRawClass();
|
||||
return schemaFactory.apply(messageClass);
|
||||
}
|
||||
|
||||
@@ -221,7 +221,7 @@ public class MethodPulsarListenerEndpoint<V> extends AbstractPulsarListenerEndpo
|
||||
|
||||
private ResolvableType resolvableType(MethodParameter methodParameter) {
|
||||
ResolvableType resolvableType = ResolvableType.forMethodParameter(methodParameter);
|
||||
final Class<?> rawClass = resolvableType.getRawClass();
|
||||
Class<?> rawClass = resolvableType.getRawClass();
|
||||
if (rawClass != null && isContainerType(rawClass)) {
|
||||
resolvableType = resolvableType.getGeneric(0);
|
||||
}
|
||||
|
||||
@@ -114,7 +114,7 @@ public class MethodReactivePulsarListenerEndpoint<V> extends AbstractReactivePul
|
||||
Assert.state(this.messageHandlerMethodFactory != null,
|
||||
"Could not create message listener - MessageHandlerMethodFactory not set");
|
||||
PulsarMessagingMessageListenerAdapter<V> messageListener = createMessageListenerInstance(messageConverter);
|
||||
final HandlerAdapter handlerMethod = configureListenerAdapter(messageListener);
|
||||
HandlerAdapter handlerMethod = configureListenerAdapter(messageListener);
|
||||
messageListener.setHandlerMethod(handlerMethod);
|
||||
|
||||
// Since we have access to the handler method here, check if we can type infer the
|
||||
@@ -122,14 +122,14 @@ public class MethodReactivePulsarListenerEndpoint<V> extends AbstractReactivePul
|
||||
|
||||
// TODO: filter out the payload type by excluding Consumer, Message, Messages etc.
|
||||
|
||||
final MethodParameter[] methodParameters = handlerMethod.getInvokerHandlerMethod().getMethodParameters();
|
||||
MethodParameter[] methodParameters = handlerMethod.getInvokerHandlerMethod().getMethodParameters();
|
||||
MethodParameter messageParameter = null;
|
||||
final Optional<MethodParameter> parameter = Arrays.stream(methodParameters)
|
||||
Optional<MethodParameter> parameter = Arrays.stream(methodParameters)
|
||||
.filter(methodParameter1 -> !methodParameter1.getParameterType().equals(Consumer.class)
|
||||
|| !methodParameter1.getParameterType().equals(Acknowledgement.class)
|
||||
|| !methodParameter1.hasParameterAnnotation(Header.class))
|
||||
.findFirst();
|
||||
final long count = Arrays.stream(methodParameters)
|
||||
long count = Arrays.stream(methodParameters)
|
||||
.filter(methodParameter1 -> !methodParameter1.getParameterType().equals(Consumer.class)
|
||||
&& !methodParameter1.getParameterType().equals(Acknowledgement.class)
|
||||
&& !methodParameter1.hasParameterAnnotation(Header.class))
|
||||
@@ -139,10 +139,9 @@ public class MethodReactivePulsarListenerEndpoint<V> extends AbstractReactivePul
|
||||
messageParameter = parameter.get();
|
||||
}
|
||||
|
||||
final DefaultReactivePulsarMessageListenerContainer<?> containerInstance = (DefaultReactivePulsarMessageListenerContainer<?>) container;
|
||||
final ReactivePulsarContainerProperties<?> pulsarContainerProperties = containerInstance
|
||||
.getContainerProperties();
|
||||
final SchemaType schemaType = pulsarContainerProperties.getSchemaType();
|
||||
DefaultReactivePulsarMessageListenerContainer<?> containerInstance = (DefaultReactivePulsarMessageListenerContainer<?>) container;
|
||||
ReactivePulsarContainerProperties<?> pulsarContainerProperties = containerInstance.getContainerProperties();
|
||||
SchemaType schemaType = pulsarContainerProperties.getSchemaType();
|
||||
if (schemaType != SchemaType.NONE) {
|
||||
switch (schemaType) {
|
||||
case STRING -> pulsarContainerProperties.setSchema((Schema) Schema.STRING);
|
||||
@@ -188,7 +187,7 @@ public class MethodReactivePulsarListenerEndpoint<V> extends AbstractReactivePul
|
||||
}
|
||||
}
|
||||
}
|
||||
final SchemaType type = pulsarContainerProperties.getSchema().getSchemaInfo().getType();
|
||||
SchemaType type = pulsarContainerProperties.getSchema().getSchemaInfo().getType();
|
||||
pulsarContainerProperties.setSchemaType(type);
|
||||
|
||||
ReactiveMessageConsumerBuilderCustomizer<V> customizer1 = b -> b.deadLetterPolicy(this.deadLetterPolicy);
|
||||
@@ -204,7 +203,7 @@ public class MethodReactivePulsarListenerEndpoint<V> extends AbstractReactivePul
|
||||
|
||||
private Schema<?> getMessageSchema(MethodParameter messageParameter, Function<Class<?>, Schema<?>> schemaFactory) {
|
||||
ResolvableType messageType = resolvableType(messageParameter);
|
||||
final Class<?> messageClass = messageType.getRawClass();
|
||||
Class<?> messageClass = messageType.getRawClass();
|
||||
return schemaFactory.apply(messageClass);
|
||||
}
|
||||
|
||||
@@ -219,7 +218,7 @@ public class MethodReactivePulsarListenerEndpoint<V> extends AbstractReactivePul
|
||||
|
||||
private ResolvableType resolvableType(MethodParameter methodParameter) {
|
||||
ResolvableType resolvableType = ResolvableType.forMethodParameter(methodParameter);
|
||||
final Class<?> rawClass = resolvableType.getRawClass();
|
||||
Class<?> rawClass = resolvableType.getRawClass();
|
||||
if (rawClass != null && isContainerType(rawClass)) {
|
||||
resolvableType = resolvableType.getGeneric(0);
|
||||
}
|
||||
|
||||
@@ -95,7 +95,7 @@ public class CachingPulsarProducerFactory<T> extends DefaultPulsarProducerFactor
|
||||
@Override
|
||||
protected Producer<T> doCreateProducer(Schema<T> schema, @Nullable String topic,
|
||||
@Nullable Collection<String> encryptionKeys, @Nullable List<ProducerBuilderCustomizer<T>> customizers) {
|
||||
final String topicName = ProducerUtils.resolveTopicName(topic, this);
|
||||
String topicName = ProducerUtils.resolveTopicName(topic, this);
|
||||
ProducerCacheKey<T> producerCacheKey = new ProducerCacheKey<>(schema, topicName,
|
||||
encryptionKeys == null ? null : new HashSet<>(encryptionKeys), customizers);
|
||||
return this.producerCache.get(producerCacheKey,
|
||||
|
||||
@@ -158,14 +158,14 @@ public class PulsarTemplate<T> implements PulsarOperations<T>, BeanNameAware {
|
||||
@Nullable Collection<String> encryptionKeys, T message,
|
||||
@Nullable TypedMessageBuilderCustomizer<T> typedMessageBuilderCustomizer,
|
||||
@Nullable ProducerBuilderCustomizer<T> producerCustomizer) throws PulsarClientException {
|
||||
final String topicName = ProducerUtils.resolveTopicName(topic, this.producerFactory);
|
||||
String topicName = ProducerUtils.resolveTopicName(topic, this.producerFactory);
|
||||
this.logger.trace(() -> String.format("Sending msg to '%s' topic", topicName));
|
||||
|
||||
PulsarMessageSenderContext senderContext = PulsarMessageSenderContext.newContext(topicName, this.beanName);
|
||||
Observation observation = newObservation(senderContext);
|
||||
try {
|
||||
observation.start();
|
||||
final Producer<T> producer = prepareProducerForSend(topic, message, encryptionKeys, producerCustomizer);
|
||||
Producer<T> producer = prepareProducerForSend(topic, message, encryptionKeys, producerCustomizer);
|
||||
TypedMessageBuilder<T> messageBuilder = producer.newMessage().value(message);
|
||||
if (typedMessageBuilderCustomizer != null) {
|
||||
typedMessageBuilderCustomizer.customize(messageBuilder);
|
||||
|
||||
@@ -77,9 +77,9 @@ public class DefaultReactivePulsarSenderFactory<T> implements ReactivePulsarSend
|
||||
|
||||
private ReactiveMessageSender<T> doCreateReactiveMessageSender(String topic, Schema<T> schema,
|
||||
List<ReactiveMessageSenderBuilderCustomizer<T>> customizers) {
|
||||
final String resolvedTopic = ReactiveMessageSenderUtils.resolveTopicName(topic, this);
|
||||
String resolvedTopic = ReactiveMessageSenderUtils.resolveTopicName(topic, this);
|
||||
this.logger.trace(() -> String.format("Creating reactive message sender for '%s' topic", resolvedTopic));
|
||||
final ReactiveMessageSenderBuilder<T> sender = this.reactivePulsarClient.messageSender(schema);
|
||||
ReactiveMessageSenderBuilder<T> sender = this.reactivePulsarClient.messageSender(schema);
|
||||
sender.applySpec(this.reactiveMessageSenderSpec);
|
||||
sender.topic(resolvedTopic);
|
||||
if (this.reactiveMessageSenderCache != null) {
|
||||
|
||||
@@ -90,7 +90,7 @@ public class ReactivePulsarTemplate<T> implements ReactivePulsarOperations<T> {
|
||||
private Mono<MessageId> doSend(String topic, T message,
|
||||
MessageSpecBuilderCustomizer<T> messageSpecBuilderCustomizer,
|
||||
ReactiveMessageSenderBuilderCustomizer<T> customizer) {
|
||||
final String topicName = ReactiveMessageSenderUtils.resolveTopicName(topic, this.reactiveMessageSenderFactory);
|
||||
String topicName = ReactiveMessageSenderUtils.resolveTopicName(topic, this.reactiveMessageSenderFactory);
|
||||
this.logger.trace(() -> String.format("Sending reactive message to '%s' topic", topicName));
|
||||
ReactiveMessageSender<T> sender = createMessageSender(topic, message, customizer);
|
||||
return sender.sendOne(getMessageSpec(messageSpecBuilderCustomizer, message)).doOnError(
|
||||
@@ -100,7 +100,7 @@ public class ReactivePulsarTemplate<T> implements ReactivePulsarOperations<T> {
|
||||
}
|
||||
|
||||
private Flux<MessageId> doSendMany(String topic, Flux<T> messages) {
|
||||
final String topicName = ReactiveMessageSenderUtils.resolveTopicName(topic, this.reactiveMessageSenderFactory);
|
||||
String topicName = ReactiveMessageSenderUtils.resolveTopicName(topic, this.reactiveMessageSenderFactory);
|
||||
this.logger.trace(() -> String.format("Sending reactive messages to '%s' topic", topicName));
|
||||
|
||||
if (this.schema != null) {
|
||||
|
||||
@@ -51,7 +51,7 @@ public class DefaultPulsarConsumerErrorHandler<T> implements PulsarConsumerError
|
||||
|
||||
@Override
|
||||
public boolean shouldRetryMessage(Exception exception, Message<T> message) {
|
||||
final Pair pair = this.backOffExecutionThreadLocal.get();
|
||||
Pair pair = this.backOffExecutionThreadLocal.get();
|
||||
long nextBackOff;
|
||||
BackOffExecution backOffExecution;
|
||||
if (pair != null && pair.message.equals(message)) {
|
||||
@@ -85,7 +85,7 @@ public class DefaultPulsarConsumerErrorHandler<T> implements PulsarConsumerError
|
||||
@SuppressWarnings("unchecked")
|
||||
public Message<T> currentMessage() {
|
||||
// there is only one message tracked at any time.
|
||||
final Pair pair = this.backOffExecutionThreadLocal.get();
|
||||
Pair pair = this.backOffExecutionThreadLocal.get();
|
||||
if (pair == null) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -274,28 +274,28 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
|
||||
|
||||
private void populateAllNecessaryPropertiesIfNeedBe(Map<String, Object> currentProperties) {
|
||||
if (currentProperties.containsKey("topicNames")) {
|
||||
final String topicsFromMap = (String) currentProperties.get("topicNames");
|
||||
final String[] topicNames = StringUtils.delimitedListToStringArray(topicsFromMap, ",");
|
||||
final Set<String> propertiesDefinedTopics = Set.of(topicNames);
|
||||
String topicsFromMap = (String) currentProperties.get("topicNames");
|
||||
String[] topicNames = StringUtils.delimitedListToStringArray(topicsFromMap, ",");
|
||||
Set<String> propertiesDefinedTopics = Set.of(topicNames);
|
||||
if (!propertiesDefinedTopics.isEmpty()) {
|
||||
currentProperties.put("topicNames", propertiesDefinedTopics);
|
||||
}
|
||||
}
|
||||
if (!currentProperties.containsKey("subscriptionType")) {
|
||||
final SubscriptionType subscriptionType = this.containerProperties.getSubscriptionType();
|
||||
SubscriptionType subscriptionType = this.containerProperties.getSubscriptionType();
|
||||
if (subscriptionType != null) {
|
||||
currentProperties.put("subscriptionType", subscriptionType);
|
||||
}
|
||||
}
|
||||
if (!currentProperties.containsKey("topicNames")) {
|
||||
final String[] topics = this.containerProperties.getTopics();
|
||||
final Set<String> listenerDefinedTopics = new HashSet<>(Arrays.stream(topics).toList());
|
||||
String[] topics = this.containerProperties.getTopics();
|
||||
Set<String> listenerDefinedTopics = new HashSet<>(Arrays.stream(topics).toList());
|
||||
if (!listenerDefinedTopics.isEmpty()) {
|
||||
currentProperties.put("topicNames", listenerDefinedTopics);
|
||||
}
|
||||
}
|
||||
if (!currentProperties.containsKey("topicsPattern")) {
|
||||
final String topicsPattern = this.containerProperties.getTopicsPattern();
|
||||
String topicsPattern = this.containerProperties.getTopicsPattern();
|
||||
if (topicsPattern != null) {
|
||||
currentProperties.put("topicsPattern", topicsPattern);
|
||||
}
|
||||
@@ -305,15 +305,15 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
|
||||
currentProperties.put("subscriptionName", this.containerProperties.getSubscriptionName());
|
||||
}
|
||||
}
|
||||
final RedeliveryBackoff negativeAckRedeliveryBackoff = DefaultPulsarMessageListenerContainer.this.negativeAckRedeliveryBackoff;
|
||||
RedeliveryBackoff negativeAckRedeliveryBackoff = DefaultPulsarMessageListenerContainer.this.negativeAckRedeliveryBackoff;
|
||||
if (negativeAckRedeliveryBackoff != null) {
|
||||
currentProperties.put("negativeAckRedeliveryBackoff", negativeAckRedeliveryBackoff);
|
||||
}
|
||||
final RedeliveryBackoff ackTimeoutRedeliveryBackoff = DefaultPulsarMessageListenerContainer.this.ackTimeoutRedeliveryBackoff;
|
||||
RedeliveryBackoff ackTimeoutRedeliveryBackoff = DefaultPulsarMessageListenerContainer.this.ackTimeoutRedeliveryBackoff;
|
||||
if (ackTimeoutRedeliveryBackoff != null) {
|
||||
currentProperties.put("ackTimeoutRedeliveryBackoff", ackTimeoutRedeliveryBackoff);
|
||||
}
|
||||
final DeadLetterPolicy deadLetterPolicy = DefaultPulsarMessageListenerContainer.this.deadLetterPolicy;
|
||||
DeadLetterPolicy deadLetterPolicy = DefaultPulsarMessageListenerContainer.this.deadLetterPolicy;
|
||||
if (deadLetterPolicy != null) {
|
||||
currentProperties.put("deadLetterPolicy", deadLetterPolicy);
|
||||
}
|
||||
@@ -380,8 +380,7 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
|
||||
this.consumer.acknowledge(messages);
|
||||
}
|
||||
else {
|
||||
final Stream<Message<T>> stream = StreamSupport.stream(messages.spliterator(),
|
||||
true);
|
||||
Stream<Message<T>> stream = StreamSupport.stream(messages.spliterator(), true);
|
||||
Message<T> last = stream.reduce((a, b) -> b).orElse(null);
|
||||
this.consumer.acknowledgeCumulative(last);
|
||||
}
|
||||
@@ -493,7 +492,7 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
|
||||
|
||||
PulsarBatchListenerFailedException pulsarBatchListenerFailedException = (PulsarBatchListenerFailedException) exception;
|
||||
Message<T> pulsarMessage = getPulsarMessageCausedTheException(pulsarBatchListenerFailedException);
|
||||
final Message<T> theCurrentPulsarMessageTracked = this.pulsarConsumerErrorHandler.currentMessage();
|
||||
Message<T> theCurrentPulsarMessageTracked = this.pulsarConsumerErrorHandler.currentMessage();
|
||||
// Previous message in error handled during retry but another msg in sublist
|
||||
// caused error;
|
||||
// resetting state in order to track it
|
||||
@@ -507,10 +506,10 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
|
||||
// handled on the retry. Otherwise, if we are out of retries then the sublist
|
||||
// does not include
|
||||
// the message in error (it instead gets recovered).
|
||||
final int indexOfFailedMessage = messageList.indexOf(pulsarMessage);
|
||||
int indexOfFailedMessage = messageList.indexOf(pulsarMessage);
|
||||
messageList = messageList.subList(indexOfFailedMessage, messageList.size());
|
||||
final boolean toBeRetried = this.pulsarConsumerErrorHandler
|
||||
.shouldRetryMessage(pulsarBatchListenerFailedException, pulsarMessage);
|
||||
boolean toBeRetried = this.pulsarConsumerErrorHandler.shouldRetryMessage(pulsarBatchListenerFailedException,
|
||||
pulsarMessage);
|
||||
if (toBeRetried) {
|
||||
inRetryMode.set(true);
|
||||
}
|
||||
@@ -535,7 +534,7 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
|
||||
}
|
||||
|
||||
private void invokeRecordListenerErrorHandler(AtomicBoolean inRetryMode, Message<T> message, Exception e) {
|
||||
final boolean toBeRetried = this.pulsarConsumerErrorHandler.shouldRetryMessage(e, message);
|
||||
boolean toBeRetried = this.pulsarConsumerErrorHandler.shouldRetryMessage(e, message);
|
||||
if (toBeRetried) {
|
||||
inRetryMode.set(true);
|
||||
}
|
||||
@@ -576,7 +575,7 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
|
||||
this.consumer.acknowledge(messages);
|
||||
}
|
||||
else {
|
||||
final Stream<Message<T>> stream = StreamSupport.stream(messages.spliterator(), true);
|
||||
Stream<Message<T>> stream = StreamSupport.stream(messages.spliterator(), true);
|
||||
Message<T> last = stream.reduce((a, b) -> b).orElse(null);
|
||||
this.consumer.acknowledgeCumulative(last);
|
||||
}
|
||||
|
||||
@@ -79,7 +79,7 @@ public class PulsarBatchMessagingMessageListenerAdapter<V> extends PulsarMessagi
|
||||
else if (isPulsarMessageList() && isHeaderFound()) { // List<PulsarMessage>,
|
||||
// @Header
|
||||
List<Message<?>> messages = toSpringMessages(consumer, msg);
|
||||
final Map<String, List<Object>> aggregatedHeaders = withAggregatedHeaders(messages);
|
||||
Map<String, List<Object>> aggregatedHeaders = withAggregatedHeaders(messages);
|
||||
List<Object> list1 = new ArrayList<>(msg);
|
||||
message = MessageBuilder.withPayload(list1).copyHeaders(aggregatedHeaders).build();
|
||||
}
|
||||
@@ -89,7 +89,7 @@ public class PulsarBatchMessagingMessageListenerAdapter<V> extends PulsarMessagi
|
||||
}
|
||||
else if (isMessageList() && isHeaderFound()) { // List<SpringMessage>, @Header
|
||||
List<Message<?>> messages = toSpringMessages(consumer, msg);
|
||||
final Map<String, List<Object>> aggregatedHeaders = withAggregatedHeaders(messages);
|
||||
Map<String, List<Object>> aggregatedHeaders = withAggregatedHeaders(messages);
|
||||
message = MessageBuilder.withPayload(messages).copyHeaders(aggregatedHeaders).build();
|
||||
}
|
||||
else if (this.isSimpleExtraction()) { // List<Object>
|
||||
@@ -99,7 +99,7 @@ public class PulsarBatchMessagingMessageListenerAdapter<V> extends PulsarMessagi
|
||||
}
|
||||
else if (isHeaderFound()) { // List<Object>, @Header
|
||||
List<Message<?>> messages = toSpringMessages(consumer, msg);
|
||||
final Map<String, List<Object>> aggregatedHeaders = withAggregatedHeaders(messages);
|
||||
Map<String, List<Object>> aggregatedHeaders = withAggregatedHeaders(messages);
|
||||
List<V> list = new ArrayList<>(msg.size());
|
||||
msg.stream().iterator().forEachRemaining(vMessage -> list.add(vMessage.getValue()));
|
||||
message = MessageBuilder.withPayload(list).copyHeaders(aggregatedHeaders).build();
|
||||
@@ -123,7 +123,7 @@ public class PulsarBatchMessagingMessageListenerAdapter<V> extends PulsarMessagi
|
||||
}
|
||||
|
||||
private Map<String, List<Object>> withAggregatedHeaders(List<Message<?>> messages) {
|
||||
final Map<String, List<Object>> aggregatedHeaders = new HashMap<>();
|
||||
Map<String, List<Object>> aggregatedHeaders = new HashMap<>();
|
||||
for (Message<?> message : messages) {
|
||||
message.getHeaders().forEach((s, o) -> {
|
||||
List<Object> objects = aggregatedHeaders.computeIfAbsent(s, k -> new ArrayList<>());
|
||||
@@ -139,8 +139,7 @@ public class PulsarBatchMessagingMessageListenerAdapter<V> extends PulsarMessagi
|
||||
return messages;
|
||||
}
|
||||
|
||||
protected void invoke(Object records, Consumer<V> consumer, final Message<?> message,
|
||||
Acknowledgement acknowledgement) {
|
||||
protected void invoke(Object records, Consumer<V> consumer, Message<?> message, Acknowledgement acknowledgement) {
|
||||
try {
|
||||
invokeHandler(records, message, consumer, acknowledgement);
|
||||
}
|
||||
|
||||
@@ -49,7 +49,7 @@ public class PulsarMessagingMessageConverter<V> implements PulsarRecordMessageCo
|
||||
@Override
|
||||
public Message<?> toMessage(org.apache.pulsar.client.api.Message<V> record, Consumer<V> consumer, Type type) {
|
||||
|
||||
final Map<String, Object> messageHeaders = new HashMap<>();
|
||||
Map<String, Object> messageHeaders = new HashMap<>();
|
||||
this.pulsarMessageHeaderMapper.toHeaders(record, messageHeaders);
|
||||
Message<?> message = MessageBuilder.createMessage(extractAndConvertValue(record),
|
||||
new MessageHeaders(messageHeaders));
|
||||
|
||||
@@ -64,13 +64,13 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
@Test
|
||||
void testRecordAck() throws Exception {
|
||||
Map<String, Object> config = new HashMap<>();
|
||||
final Set<String> strings = new HashSet<>();
|
||||
Set<String> strings = new HashSet<>();
|
||||
strings.add("cons-ack-tests-011");
|
||||
config.put("topicNames", strings);
|
||||
config.put("subscriptionName", "cons-ack-tests-sb-011");
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
|
||||
new DefaultPulsarConsumerFactory<>(pulsarClient, config));
|
||||
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
@@ -80,7 +80,7 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
pulsarContainerProperties.setAckMode(AckMode.RECORD);
|
||||
DefaultPulsarMessageListenerContainer<String> container = new DefaultPulsarMessageListenerContainer<>(
|
||||
pulsarConsumerFactory, pulsarContainerProperties);
|
||||
final Consumer<String> containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container);
|
||||
Consumer<String> containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container);
|
||||
|
||||
CountDownLatch latch = new CountDownLatch(10);
|
||||
|
||||
@@ -91,9 +91,9 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
|
||||
Map<String, Object> prodConfig = new HashMap<>();
|
||||
prodConfig.put("topicName", "cons-ack-tests-011");
|
||||
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
|
||||
pulsarClient, prodConfig);
|
||||
final PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
prodConfig);
|
||||
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
pulsarTemplate.sendAsync("hello john doe");
|
||||
}
|
||||
@@ -105,13 +105,13 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
@Test
|
||||
void testBatchAck() throws Exception {
|
||||
Map<String, Object> config = new HashMap<>();
|
||||
final Set<String> strings = new HashSet<>();
|
||||
Set<String> strings = new HashSet<>();
|
||||
strings.add("cons-ack-tests-012");
|
||||
config.put("topicNames", strings);
|
||||
config.put("subscriptionName", "cons-ack-tests-sb-012");
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
|
||||
new DefaultPulsarConsumerFactory<>(pulsarClient, config));
|
||||
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
@@ -121,13 +121,13 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
pulsarContainerProperties.setSchema(Schema.STRING);
|
||||
DefaultPulsarMessageListenerContainer<String> container = new DefaultPulsarMessageListenerContainer<>(
|
||||
pulsarConsumerFactory, pulsarContainerProperties);
|
||||
final Consumer<String> containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container);
|
||||
Consumer<String> containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container);
|
||||
|
||||
Map<String, Object> prodConfig = new HashMap<>();
|
||||
prodConfig.put("topicName", "cons-ack-tests-012");
|
||||
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
|
||||
pulsarClient, prodConfig);
|
||||
final PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
prodConfig);
|
||||
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
pulsarTemplate.sendAsync("hello john doe");
|
||||
}
|
||||
@@ -145,9 +145,9 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
Map<String, Object> config = new HashMap<>();
|
||||
config.put("topicNames", Collections.singleton("cons-ack-tests-013"));
|
||||
config.put("subscriptionName", "cons-ack-tests-sb-013");
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
|
||||
new DefaultPulsarConsumerFactory<>(pulsarClient, config));
|
||||
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
@@ -162,7 +162,7 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
pulsarContainerProperties.setSchema(Schema.STRING);
|
||||
DefaultPulsarMessageListenerContainer<String> container = new DefaultPulsarMessageListenerContainer<>(
|
||||
pulsarConsumerFactory, pulsarContainerProperties);
|
||||
final Consumer<String> containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container);
|
||||
Consumer<String> containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container);
|
||||
|
||||
AtomicInteger ackCallCount = new AtomicInteger(0);
|
||||
doAnswer(invocation -> {
|
||||
@@ -172,9 +172,9 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
|
||||
Map<String, Object> prodConfig = new HashMap<>();
|
||||
prodConfig.put("topicName", "cons-ack-tests-013");
|
||||
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
|
||||
pulsarClient, prodConfig);
|
||||
final PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
prodConfig);
|
||||
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
pulsarTemplate.sendAsync("hello john doe");
|
||||
}
|
||||
@@ -186,7 +186,7 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
await().atMost(Duration.ofSeconds(10))
|
||||
.untilAsserted(() -> verify(containerConsumer, times(5)).negativeAcknowledge(any(Message.class)));
|
||||
|
||||
final int ackCalls = ackCallCount.get();
|
||||
int ackCalls = ackCallCount.get();
|
||||
if (ackCalls < 5) {
|
||||
await().atMost(Duration.ofSeconds(10))
|
||||
.untilAsserted(() -> verify(containerConsumer, atMost(4)).acknowledge(any(MessageId.class)));
|
||||
@@ -209,17 +209,17 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
@SuppressWarnings("unchecked")
|
||||
void testManualAckForRecordListener() throws Exception {
|
||||
Map<String, Object> config = new HashMap<>();
|
||||
final Set<String> strings = new HashSet<>();
|
||||
Set<String> strings = new HashSet<>();
|
||||
strings.add("cons-ack-tests-014");
|
||||
config.put("topicNames", strings);
|
||||
config.put("subscriptionName", "cons-ack-tests-sb-014");
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
|
||||
new DefaultPulsarConsumerFactory<>(pulsarClient, config));
|
||||
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
final List<Acknowledgement> acksObjects = new ArrayList<>();
|
||||
List<Acknowledgement> acksObjects = new ArrayList<>();
|
||||
PulsarAcknowledgingMessageListener<?> pulsarAcknowledgingMessageListener = (consumer, msg, acknowledgement) -> {
|
||||
acksObjects.add(acknowledgement);
|
||||
acknowledgement.acknowledge();
|
||||
@@ -231,7 +231,7 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
pulsarContainerProperties.setAckMode(AckMode.MANUAL);
|
||||
DefaultPulsarMessageListenerContainer<String> container = new DefaultPulsarMessageListenerContainer<>(
|
||||
pulsarConsumerFactory, pulsarContainerProperties);
|
||||
final Consumer<String> containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container);
|
||||
Consumer<String> containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container);
|
||||
|
||||
CountDownLatch latch = new CountDownLatch(10);
|
||||
|
||||
@@ -242,9 +242,9 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
|
||||
Map<String, Object> prodConfig = new HashMap<>();
|
||||
prodConfig.put("topicName", "cons-ack-tests-014");
|
||||
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
|
||||
pulsarClient, prodConfig);
|
||||
final PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
prodConfig);
|
||||
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
pulsarTemplate.sendAsync("hello john doe");
|
||||
}
|
||||
@@ -263,13 +263,13 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
@SuppressWarnings("unchecked")
|
||||
void testBatchAckForBatchListener() throws Exception {
|
||||
Map<String, Object> config = new HashMap<>();
|
||||
final Set<String> strings = new HashSet<>();
|
||||
Set<String> strings = new HashSet<>();
|
||||
strings.add("cons-ack-tests-015");
|
||||
config.put("topicNames", strings);
|
||||
config.put("subscriptionName", "cons-ack-tests-sb-015");
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
|
||||
new DefaultPulsarConsumerFactory<>(pulsarClient, config));
|
||||
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
@@ -277,7 +277,7 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
pulsarContainerProperties.setBatchTimeoutMillis(60_000);
|
||||
pulsarContainerProperties.setBatchListener(true);
|
||||
CountDownLatch latch = new CountDownLatch(1);
|
||||
final PulsarBatchMessageListener<?> pulsarBatchMessageListener = mock(PulsarBatchMessageListener.class);
|
||||
PulsarBatchMessageListener<?> pulsarBatchMessageListener = mock(PulsarBatchMessageListener.class);
|
||||
|
||||
doAnswer(invocation -> {
|
||||
latch.countDown();
|
||||
@@ -288,13 +288,13 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
pulsarContainerProperties.setSchema(Schema.STRING);
|
||||
DefaultPulsarMessageListenerContainer<String> container = new DefaultPulsarMessageListenerContainer<>(
|
||||
pulsarConsumerFactory, pulsarContainerProperties);
|
||||
final Consumer<String> containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container);
|
||||
Consumer<String> containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container);
|
||||
|
||||
Map<String, Object> prodConfig = new HashMap<>();
|
||||
prodConfig.put("topicName", "cons-ack-tests-015");
|
||||
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
|
||||
pulsarClient, prodConfig);
|
||||
final PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
prodConfig);
|
||||
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
pulsarTemplate.sendAsync("hello john doe");
|
||||
}
|
||||
@@ -311,20 +311,20 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
@SuppressWarnings("unchecked")
|
||||
void testBatchNackForEntireBatchWhenUsingBatchListener() throws Exception {
|
||||
Map<String, Object> config = new HashMap<>();
|
||||
final Set<String> strings = new HashSet<>();
|
||||
Set<String> strings = new HashSet<>();
|
||||
strings.add("cons-ack-tests-016");
|
||||
config.put("topicNames", strings);
|
||||
config.put("subscriptionName", "cons-ack-tests-sb-016");
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
|
||||
new DefaultPulsarConsumerFactory<>(pulsarClient, config));
|
||||
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
pulsarContainerProperties.setMaxNumMessages(10);
|
||||
pulsarContainerProperties.setBatchTimeoutMillis(60_000);
|
||||
pulsarContainerProperties.setBatchListener(true);
|
||||
final PulsarBatchMessageListener<?> pulsarBatchMessageListener = mock(PulsarBatchMessageListener.class);
|
||||
PulsarBatchMessageListener<?> pulsarBatchMessageListener = mock(PulsarBatchMessageListener.class);
|
||||
CountDownLatch latch = new CountDownLatch(1);
|
||||
|
||||
doAnswer(invocation -> {
|
||||
@@ -336,13 +336,13 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
pulsarContainerProperties.setSchema(Schema.STRING);
|
||||
DefaultPulsarMessageListenerContainer<String> container = new DefaultPulsarMessageListenerContainer<>(
|
||||
pulsarConsumerFactory, pulsarContainerProperties);
|
||||
final Consumer<String> containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container);
|
||||
Consumer<String> containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container);
|
||||
|
||||
Map<String, Object> prodConfig = new HashMap<>();
|
||||
prodConfig.put("topicName", "cons-ack-tests-016");
|
||||
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
|
||||
pulsarClient, prodConfig);
|
||||
final PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
prodConfig);
|
||||
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
pulsarTemplate.sendAsync("hello john doe");
|
||||
}
|
||||
@@ -361,13 +361,13 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
Map<String, Object> config = new HashMap<>();
|
||||
config.put("topicNames", Set.of("duplicate-message-test"));
|
||||
config.put("subscriptionName", "duplicate-sub-1");
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient,
|
||||
config);
|
||||
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
final AtomicInteger counter1 = new AtomicInteger(0);
|
||||
AtomicInteger counter1 = new AtomicInteger(0);
|
||||
pulsarContainerProperties.setMessageListener((PulsarRecordMessageListener<?>) (consumer, msg) -> {
|
||||
counter1.getAndIncrement();
|
||||
});
|
||||
@@ -377,9 +377,9 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
container1.start();
|
||||
|
||||
Map<String, Object> prodConfig = Collections.singletonMap("topicName", "duplicate-message-test");
|
||||
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
|
||||
pulsarClient, prodConfig);
|
||||
final PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
prodConfig);
|
||||
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
pulsarTemplate.send("hello john doe");
|
||||
|
||||
while (counter1.get() == 0) {
|
||||
@@ -391,7 +391,7 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
// unacked message.
|
||||
container1.stop();
|
||||
|
||||
final AtomicInteger counter2 = new AtomicInteger(0);
|
||||
AtomicInteger counter2 = new AtomicInteger(0);
|
||||
pulsarContainerProperties.setMessageListener((PulsarRecordMessageListener<?>) (consumer, msg) -> {
|
||||
counter2.getAndIncrement();
|
||||
});
|
||||
|
||||
@@ -54,14 +54,14 @@ class FailoverConsumerTests implements PulsarTestContainerSupport {
|
||||
admin.topics().createPartitionedTopic(topicName, numPartitions);
|
||||
|
||||
Map<String, Object> config = new HashMap<>();
|
||||
final HashSet<String> topics = new HashSet<>();
|
||||
HashSet<String> topics = new HashSet<>();
|
||||
topics.add("my-part-topic-1");
|
||||
config.put("topicNames", topics);
|
||||
config.put("subscriptionName", "my-part-subscription-1");
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient,
|
||||
config);
|
||||
CountDownLatch latch = new CountDownLatch(3);
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
pulsarContainerProperties
|
||||
@@ -80,9 +80,9 @@ class FailoverConsumerTests implements PulsarTestContainerSupport {
|
||||
Map<String, Object> prodConfig = new HashMap<>();
|
||||
prodConfig.put("topicName", "my-part-topic-1");
|
||||
prodConfig.put("messageRoutingMode", MessageRoutingMode.CustomPartition);
|
||||
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
|
||||
pulsarClient, prodConfig);
|
||||
final PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
prodConfig);
|
||||
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
|
||||
pulsarTemplate.newMessage("hello john doe")
|
||||
.withProducerCustomizer(builder -> builder.messageRouter(new FooRouter())).sendAsync();
|
||||
@@ -91,7 +91,7 @@ class FailoverConsumerTests implements PulsarTestContainerSupport {
|
||||
pulsarTemplate.newMessage("hello buzz doe")
|
||||
.withProducerCustomizer(builder -> builder.messageRouter(new BuzzRouter())).sendAsync();
|
||||
|
||||
final boolean await = latch.await(10, TimeUnit.SECONDS);
|
||||
boolean await = latch.await(10, TimeUnit.SECONDS);
|
||||
assertThat(await).isTrue();
|
||||
}
|
||||
|
||||
|
||||
@@ -61,10 +61,10 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
|
||||
config.put("topicNames", Collections.singleton("default-error-handler-tests-1"));
|
||||
config.put("subscriptionName", "default-error-handler-tests-sub-1");
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient,
|
||||
config);
|
||||
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
PulsarRecordMessageListener<?> messageListener = mock(PulsarRecordMessageListener.class);
|
||||
@@ -112,16 +112,16 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
|
||||
config.put("topicNames", Collections.singleton("default-error-handler-tests-2"));
|
||||
config.put("subscriptionName", "default-error-handler-tests-sub-2");
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient,
|
||||
config);
|
||||
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
PulsarRecordMessageListener<?> messageListener = mock(PulsarRecordMessageListener.class);
|
||||
AtomicInteger count = new AtomicInteger(0);
|
||||
doAnswer(invocation -> {
|
||||
final int currentCount = count.incrementAndGet();
|
||||
int currentCount = count.incrementAndGet();
|
||||
if (currentCount <= 3) {
|
||||
throw new RuntimeException();
|
||||
}
|
||||
@@ -161,16 +161,16 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
|
||||
config.put("topicNames", Collections.singleton("default-error-handler-tests-3"));
|
||||
config.put("subscriptionName", "default-error-handler-tests-sub-3");
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient,
|
||||
config);
|
||||
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
PulsarRecordMessageListener<?> messageListener = mock(PulsarRecordMessageListener.class);
|
||||
doAnswer(invocation -> {
|
||||
final Message<Integer> message = invocation.getArgument(1);
|
||||
final Integer value = message.getValue();
|
||||
Message<Integer> message = invocation.getArgument(1);
|
||||
Integer value = message.getValue();
|
||||
if (value % 2 == 0) {
|
||||
throw new RuntimeException();
|
||||
}
|
||||
@@ -220,26 +220,26 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
|
||||
config.put("topicNames", Collections.singleton("default-error-handler-tests-4"));
|
||||
config.put("subscriptionName", "default-error-handler-tests-sub-4");
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient,
|
||||
config);
|
||||
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
pulsarContainerProperties.setMaxNumMessages(10);
|
||||
pulsarContainerProperties.setBatchTimeoutMillis(60_000);
|
||||
pulsarContainerProperties.setBatchListener(true);
|
||||
final PulsarBatchAcknowledgingMessageListener<?> pulsarBatchMessageListener = mock(
|
||||
PulsarBatchAcknowledgingMessageListener<?> pulsarBatchMessageListener = mock(
|
||||
PulsarBatchAcknowledgingMessageListener.class);
|
||||
|
||||
doAnswer(invocation -> {
|
||||
final List<Message<Integer>> message = invocation.getArgument(1);
|
||||
final Message<Integer> integerMessage = message.get(0);
|
||||
final Integer value = integerMessage.getValue();
|
||||
List<Message<Integer>> message = invocation.getArgument(1);
|
||||
Message<Integer> integerMessage = message.get(0);
|
||||
Integer value = integerMessage.getValue();
|
||||
if (value == 0) {
|
||||
throw new PulsarBatchListenerFailedException("failed", integerMessage);
|
||||
}
|
||||
final Acknowledgement acknowledgment = invocation.getArgument(2);
|
||||
Acknowledgement acknowledgment = invocation.getArgument(2);
|
||||
List<MessageId> messageIds = new ArrayList<>();
|
||||
for (Message<Integer> integerMessage1 : message) {
|
||||
messageIds.add(integerMessage1.getMessageId());
|
||||
@@ -262,9 +262,9 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
|
||||
|
||||
Map<String, Object> prodConfig = new HashMap<>();
|
||||
prodConfig.put("topicName", "default-error-handler-tests-4");
|
||||
final DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
|
||||
pulsarClient, prodConfig);
|
||||
final PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
prodConfig);
|
||||
PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
pulsarTemplate.sendAsync(i);
|
||||
}
|
||||
@@ -292,27 +292,27 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
|
||||
config.put("topicNames", Collections.singleton("default-error-handler-tests-5"));
|
||||
config.put("subscriptionName", "default-error-handler-tests-sub-5");
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient,
|
||||
config);
|
||||
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
pulsarContainerProperties.setMaxNumMessages(10);
|
||||
pulsarContainerProperties.setBatchTimeoutMillis(60_000);
|
||||
pulsarContainerProperties.setBatchListener(true);
|
||||
final PulsarBatchAcknowledgingMessageListener<?> pulsarBatchMessageListener = mock(
|
||||
PulsarBatchAcknowledgingMessageListener<?> pulsarBatchMessageListener = mock(
|
||||
PulsarBatchAcknowledgingMessageListener.class);
|
||||
|
||||
doAnswer(invocation -> {
|
||||
final List<Message<Integer>> messages = invocation.getArgument(1);
|
||||
List<Message<Integer>> messages = invocation.getArgument(1);
|
||||
|
||||
for (Message<Integer> message : messages) {
|
||||
if (message.getValue() == 5) {
|
||||
throw new PulsarBatchListenerFailedException("failed", message);
|
||||
}
|
||||
else {
|
||||
final Acknowledgement acknowledgment = invocation.getArgument(2);
|
||||
Acknowledgement acknowledgment = invocation.getArgument(2);
|
||||
acknowledgment.acknowledge(message.getMessageId());
|
||||
}
|
||||
}
|
||||
@@ -333,9 +333,9 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
|
||||
|
||||
Map<String, Object> prodConfig = new HashMap<>();
|
||||
prodConfig.put("topicName", "default-error-handler-tests-5");
|
||||
final DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
|
||||
pulsarClient, prodConfig);
|
||||
final PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
prodConfig);
|
||||
PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
pulsarTemplate.sendAsync(i);
|
||||
}
|
||||
@@ -362,27 +362,27 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
|
||||
config.put("topicNames", Collections.singleton("default-error-handler-tests-6"));
|
||||
config.put("subscriptionName", "default-error-handler-tests-sub-6");
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient,
|
||||
config);
|
||||
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
pulsarContainerProperties.setMaxNumMessages(10);
|
||||
pulsarContainerProperties.setBatchTimeoutMillis(60_000);
|
||||
pulsarContainerProperties.setBatchListener(true);
|
||||
final PulsarBatchAcknowledgingMessageListener<?> pulsarBatchMessageListener = mock(
|
||||
PulsarBatchAcknowledgingMessageListener<?> pulsarBatchMessageListener = mock(
|
||||
PulsarBatchAcknowledgingMessageListener.class);
|
||||
|
||||
doAnswer(invocation -> {
|
||||
final List<Message<Integer>> messages = invocation.getArgument(1);
|
||||
List<Message<Integer>> messages = invocation.getArgument(1);
|
||||
|
||||
for (Message<Integer> message : messages) {
|
||||
if (message.getValue() == 2 || message.getValue() == 5) {
|
||||
throw new PulsarBatchListenerFailedException("failed", message);
|
||||
}
|
||||
else {
|
||||
final Acknowledgement acknowledgment = invocation.getArgument(2);
|
||||
Acknowledgement acknowledgment = invocation.getArgument(2);
|
||||
acknowledgment.acknowledge(message.getMessageId());
|
||||
}
|
||||
}
|
||||
@@ -403,9 +403,9 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
|
||||
|
||||
Map<String, Object> prodConfig = new HashMap<>();
|
||||
prodConfig.put("topicName", "default-error-handler-tests-6");
|
||||
final DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
|
||||
pulsarClient, prodConfig);
|
||||
final PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
prodConfig);
|
||||
PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
pulsarTemplate.sendAsync(i);
|
||||
}
|
||||
@@ -432,25 +432,25 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
|
||||
config.put("topicNames", Collections.singleton("default-error-handler-tests-7"));
|
||||
config.put("subscriptionName", "default-error-handler-tests-sub-7");
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient,
|
||||
config);
|
||||
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
pulsarContainerProperties.setMaxNumMessages(10);
|
||||
pulsarContainerProperties.setBatchTimeoutMillis(60_000);
|
||||
pulsarContainerProperties.setBatchListener(true);
|
||||
final PulsarBatchAcknowledgingMessageListener<?> pulsarBatchMessageListener = mock(
|
||||
PulsarBatchAcknowledgingMessageListener<?> pulsarBatchMessageListener = mock(
|
||||
PulsarBatchAcknowledgingMessageListener.class);
|
||||
|
||||
AtomicInteger count = new AtomicInteger(0);
|
||||
doAnswer(invocation -> {
|
||||
final List<Message<Integer>> messages = invocation.getArgument(1);
|
||||
final Acknowledgement acknowledgment = invocation.getArgument(2);
|
||||
List<Message<Integer>> messages = invocation.getArgument(1);
|
||||
Acknowledgement acknowledgment = invocation.getArgument(2);
|
||||
for (Message<Integer> message : messages) {
|
||||
if (message.getValue() == 5) {
|
||||
final int currentCount = count.getAndIncrement();
|
||||
int currentCount = count.getAndIncrement();
|
||||
if (currentCount < 3) {
|
||||
throw new PulsarBatchListenerFailedException("failed", message);
|
||||
}
|
||||
@@ -479,9 +479,9 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
|
||||
|
||||
Map<String, Object> prodConfig = new HashMap<>();
|
||||
prodConfig.put("topicName", "default-error-handler-tests-7");
|
||||
final DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
|
||||
pulsarClient, prodConfig);
|
||||
final PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
prodConfig);
|
||||
PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
pulsarTemplate.sendAsync(i);
|
||||
}
|
||||
@@ -501,25 +501,25 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
|
||||
config.put("topicNames", Collections.singleton("default-error-handler-tests-8"));
|
||||
config.put("subscriptionName", "default-error-handler-tests-sub-8");
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient,
|
||||
config);
|
||||
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
pulsarContainerProperties.setMaxNumMessages(10);
|
||||
pulsarContainerProperties.setBatchTimeoutMillis(60_000);
|
||||
pulsarContainerProperties.setBatchListener(true);
|
||||
final PulsarBatchAcknowledgingMessageListener<?> pulsarBatchMessageListener = mock(
|
||||
PulsarBatchAcknowledgingMessageListener<?> pulsarBatchMessageListener = mock(
|
||||
PulsarBatchAcknowledgingMessageListener.class);
|
||||
|
||||
AtomicInteger count = new AtomicInteger(0);
|
||||
doAnswer(invocation -> {
|
||||
final List<Message<Integer>> messages = invocation.getArgument(1);
|
||||
final Acknowledgement acknowledgment = invocation.getArgument(2);
|
||||
List<Message<Integer>> messages = invocation.getArgument(1);
|
||||
Acknowledgement acknowledgment = invocation.getArgument(2);
|
||||
for (Message<Integer> message : messages) {
|
||||
if (message.getValue() == 5) {
|
||||
final int currentCount = count.getAndIncrement();
|
||||
int currentCount = count.getAndIncrement();
|
||||
if (currentCount < 3) {
|
||||
throw new PulsarBatchListenerFailedException("failed", message);
|
||||
}
|
||||
@@ -551,9 +551,9 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
|
||||
|
||||
Map<String, Object> prodConfig = new HashMap<>();
|
||||
prodConfig.put("topicName", "default-error-handler-tests-8");
|
||||
final DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
|
||||
pulsarClient, prodConfig);
|
||||
final PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
prodConfig);
|
||||
PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
pulsarTemplate.sendAsync(i);
|
||||
}
|
||||
|
||||
@@ -60,14 +60,14 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
|
||||
@Test
|
||||
void basicDefaultConsumer() throws Exception {
|
||||
Map<String, Object> config = new HashMap<>();
|
||||
final HashSet<String> strings = new HashSet<>();
|
||||
HashSet<String> strings = new HashSet<>();
|
||||
strings.add("dpmlct-012");
|
||||
config.put("topicNames", strings);
|
||||
config.put("subscriptionName", "dpmlct-sb-012");
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient,
|
||||
config);
|
||||
CountDownLatch latch = new CountDownLatch(1);
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
pulsarContainerProperties
|
||||
@@ -78,9 +78,9 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
|
||||
container.start();
|
||||
Map<String, Object> prodConfig = new HashMap<>();
|
||||
prodConfig.put("topicName", "dpmlct-012");
|
||||
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
|
||||
pulsarClient, prodConfig);
|
||||
final PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
prodConfig);
|
||||
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
pulsarTemplate.sendAsync("hello john doe");
|
||||
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
|
||||
container.stop();
|
||||
@@ -90,15 +90,15 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
|
||||
@Test
|
||||
void subscriptionInitialPositionEarliest() throws Exception {
|
||||
Map<String, Object> config = new HashMap<>();
|
||||
final HashSet<String> strings = new HashSet<>();
|
||||
HashSet<String> strings = new HashSet<>();
|
||||
strings.add("dpmlct-013");
|
||||
config.put("topicNames", strings);
|
||||
config.put("subscriptionName", "dpmlct-sb-013");
|
||||
config.put("subscriptionInitialPosition", SubscriptionInitialPosition.Earliest);
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient,
|
||||
config);
|
||||
CountDownLatch latch = new CountDownLatch(5);
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
pulsarContainerProperties
|
||||
@@ -109,9 +109,9 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
|
||||
|
||||
Map<String, Object> prodConfig = new HashMap<>();
|
||||
prodConfig.put("topicName", "dpmlct-013");
|
||||
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
|
||||
pulsarClient, prodConfig);
|
||||
final PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
prodConfig);
|
||||
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
for (int i = 0; i < 5; i++) {
|
||||
pulsarTemplate.send("hello john doe" + i);
|
||||
}
|
||||
@@ -125,14 +125,14 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
|
||||
@Test
|
||||
void subscriptionInitialPositionDefaultLatest() throws Exception {
|
||||
Map<String, Object> config = new HashMap<>();
|
||||
final HashSet<String> strings = new HashSet<>();
|
||||
HashSet<String> strings = new HashSet<>();
|
||||
strings.add("dpmlct-014");
|
||||
config.put("topicNames", strings);
|
||||
config.put("subscriptionName", "dpmlct-sb-014");
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient,
|
||||
config);
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
List<String> messages = new ArrayList<>();
|
||||
pulsarContainerProperties.setMessageListener(
|
||||
@@ -143,9 +143,9 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
|
||||
|
||||
Map<String, Object> prodConfig = new HashMap<>();
|
||||
prodConfig.put("topicName", "dpmlct-014");
|
||||
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
|
||||
pulsarClient, prodConfig);
|
||||
final PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
prodConfig);
|
||||
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
for (int i = 0; i < 5; i++) {
|
||||
pulsarTemplate.send("hello john doe" + i);
|
||||
}
|
||||
@@ -169,9 +169,9 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
|
||||
.maxDelayMs(5 * 1000).build();
|
||||
config.put("negativeAckRedeliveryBackoff", redeliveryBackoff);
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
|
||||
new DefaultPulsarConsumerFactory<>(pulsarClient, config));
|
||||
CountDownLatch latch = new CountDownLatch(10);
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
@@ -186,12 +186,12 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
|
||||
DefaultPulsarMessageListenerContainer<String> container = new DefaultPulsarMessageListenerContainer<>(
|
||||
pulsarConsumerFactory, pulsarContainerProperties);
|
||||
|
||||
final Consumer<String> containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container);
|
||||
Consumer<String> containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container);
|
||||
|
||||
Map<String, Object> prodConfig = Collections.singletonMap("topicName", "dpmlct-015");
|
||||
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
|
||||
pulsarClient, prodConfig);
|
||||
final PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
prodConfig);
|
||||
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
for (int i = 0; i < 5; i++) {
|
||||
pulsarTemplate.send("hello john doe" + i);
|
||||
}
|
||||
@@ -219,10 +219,10 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
|
||||
.deadLetterTopic("dpmlct-016-dlq-topic").build();
|
||||
config.put("deadLetterPolicy", deadLetterPolicy);
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient,
|
||||
config);
|
||||
|
||||
CountDownLatch dlqLatch = new CountDownLatch(1);
|
||||
CountDownLatch latch = new CountDownLatch(6);
|
||||
@@ -251,9 +251,9 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
|
||||
container.start();
|
||||
|
||||
Map<String, Object> prodConfig = Collections.singletonMap("topicName", "dpmlct-016");
|
||||
final DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
|
||||
pulsarClient, prodConfig);
|
||||
final PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
prodConfig);
|
||||
PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
for (int i = 1; i < 6; i++) {
|
||||
pulsarTemplate.send(i);
|
||||
}
|
||||
@@ -277,10 +277,10 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
|
||||
.build();
|
||||
config.put("deadLetterPolicy", deadLetterPolicy);
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build();
|
||||
DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient,
|
||||
config);
|
||||
|
||||
CountDownLatch dlqLatch = new CountDownLatch(1);
|
||||
CountDownLatch latch = new CountDownLatch(6);
|
||||
@@ -312,9 +312,9 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
|
||||
container.start();
|
||||
|
||||
Map<String, Object> prodConfig = Collections.singletonMap("topicName", "dpmlct-016");
|
||||
final DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
|
||||
pulsarClient, prodConfig);
|
||||
final PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
prodConfig);
|
||||
PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
for (int i = 1; i < 6; i++) {
|
||||
pulsarTemplate.send(i);
|
||||
}
|
||||
|
||||
@@ -129,7 +129,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport {
|
||||
@Bean
|
||||
PulsarListenerContainerFactory pulsarListenerContainerFactory(
|
||||
PulsarConsumerFactory<Object> pulsarConsumerFactory) {
|
||||
final ConcurrentPulsarListenerContainerFactory<?> pulsarListenerContainerFactory = new ConcurrentPulsarListenerContainerFactory<>(
|
||||
ConcurrentPulsarListenerContainerFactory<?> pulsarListenerContainerFactory = new ConcurrentPulsarListenerContainerFactory<>(
|
||||
pulsarConsumerFactory, new PulsarContainerProperties(), null);
|
||||
return pulsarListenerContainerFactory;
|
||||
}
|
||||
@@ -168,9 +168,9 @@ public class PulsarListenerTests implements PulsarTestContainerSupport {
|
||||
void testPulsarListenerProvidedConsumerProperties(@Autowired PulsarListenerEndpointRegistry registry)
|
||||
throws Exception {
|
||||
|
||||
final PulsarContainerProperties pulsarContainerProperties = registry.getListenerContainer("foo")
|
||||
PulsarContainerProperties pulsarContainerProperties = registry.getListenerContainer("foo")
|
||||
.getContainerProperties();
|
||||
final Properties pulsarConsumerProperties = pulsarContainerProperties.getPulsarConsumerProperties();
|
||||
Properties pulsarConsumerProperties = pulsarContainerProperties.getPulsarConsumerProperties();
|
||||
assertThat(pulsarConsumerProperties.size()).isEqualTo(2);
|
||||
assertThat(pulsarConsumerProperties.get("topicNames")).isEqualTo("foo-1");
|
||||
assertThat(pulsarConsumerProperties.get("subscriptionName")).isEqualTo("subscription-1");
|
||||
@@ -185,7 +185,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport {
|
||||
Map.of("batchingEnabled", false));
|
||||
PulsarTemplate<String> customTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
|
||||
final ConcurrentPulsarMessageListenerContainer<?> bar = (ConcurrentPulsarMessageListenerContainer<?>) registry
|
||||
ConcurrentPulsarMessageListenerContainer<?> bar = (ConcurrentPulsarMessageListenerContainer<?>) registry
|
||||
.getListenerContainer("bar");
|
||||
|
||||
assertThat(bar.getConcurrency()).isEqualTo(3);
|
||||
@@ -204,7 +204,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport {
|
||||
Map.of("batchingEnabled", false));
|
||||
PulsarTemplate<String> customTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
|
||||
final ConcurrentPulsarMessageListenerContainer<?> bar = (ConcurrentPulsarMessageListenerContainer<?>) registry
|
||||
ConcurrentPulsarMessageListenerContainer<?> bar = (ConcurrentPulsarMessageListenerContainer<?>) registry
|
||||
.getListenerContainer("bar");
|
||||
|
||||
assertThat(bar.getConcurrency()).isEqualTo(3);
|
||||
@@ -218,7 +218,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport {
|
||||
|
||||
@Test
|
||||
void ackModeAppliedToContainerFromListener(@Autowired PulsarListenerEndpointRegistry registry) {
|
||||
final PulsarContainerProperties pulsarContainerProperties = registry.getListenerContainer("ackMode-test-id")
|
||||
PulsarContainerProperties pulsarContainerProperties = registry.getListenerContainer("ackMode-test-id")
|
||||
.getContainerProperties();
|
||||
assertThat(pulsarContainerProperties.getAckMode()).isEqualTo(AckMode.RECORD);
|
||||
}
|
||||
@@ -643,7 +643,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport {
|
||||
|
||||
@Test
|
||||
void simpleListenerWithHeaders() throws Exception {
|
||||
final MessageId messageId = pulsarTemplate.newMessage("hello-simple-listener")
|
||||
MessageId messageId = pulsarTemplate.newMessage("hello-simple-listener")
|
||||
.withMessageCustomizer(
|
||||
messageBuilder -> messageBuilder.property("foo", "simpleListenerWithHeaders"))
|
||||
.withTopic("simpleListenerWithHeaders").send();
|
||||
@@ -657,7 +657,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport {
|
||||
|
||||
@Test
|
||||
void pulsarMessageListenerWithHeaders() throws Exception {
|
||||
final MessageId messageId = pulsarTemplate.newMessage("hello-pulsar-message-listener")
|
||||
MessageId messageId = pulsarTemplate.newMessage("hello-pulsar-message-listener")
|
||||
.withMessageCustomizer(
|
||||
messageBuilder -> messageBuilder.property("foo", "pulsarMessageListenerWithHeaders"))
|
||||
.withTopic("pulsarMessageListenerWithHeaders").send();
|
||||
@@ -671,7 +671,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport {
|
||||
|
||||
@Test
|
||||
void springMessagingMessageListenerWithHeaders() throws Exception {
|
||||
final MessageId messageId = pulsarTemplate.newMessage("hello-spring-messaging-message-listener")
|
||||
MessageId messageId = pulsarTemplate.newMessage("hello-spring-messaging-message-listener")
|
||||
.withMessageCustomizer(messageBuilder -> messageBuilder.property("foo",
|
||||
"springMessagingMessageListenerWithHeaders"))
|
||||
.withTopic("springMessagingMessageListenerWithHeaders").send();
|
||||
@@ -685,7 +685,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport {
|
||||
|
||||
@Test
|
||||
void simpleBatchListenerWithHeaders() throws Exception {
|
||||
final MessageId messageId = pulsarTemplate.newMessage("hello-simple-batch-listener")
|
||||
MessageId messageId = pulsarTemplate.newMessage("hello-simple-batch-listener")
|
||||
.withMessageCustomizer(
|
||||
messageBuilder -> messageBuilder.property("foo", "simpleBatchListenerWithHeaders"))
|
||||
.withTopic("simpleBatchListenerWithHeaders").send();
|
||||
@@ -698,7 +698,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport {
|
||||
|
||||
@Test
|
||||
void pulsarMessageBatchListenerWithHeaders() throws Exception {
|
||||
final MessageId messageId = pulsarTemplate.newMessage("hello-pulsar-message-batch-listener")
|
||||
MessageId messageId = pulsarTemplate.newMessage("hello-pulsar-message-batch-listener")
|
||||
.withMessageCustomizer(
|
||||
messageBuilder -> messageBuilder.property("foo", "pulsarMessageBatchListenerWithHeaders"))
|
||||
.withTopic("pulsarMessageBatchListenerWithHeaders").send();
|
||||
@@ -712,7 +712,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport {
|
||||
|
||||
@Test
|
||||
void springMessagingMessageBatchListenerWithHeaders() throws Exception {
|
||||
final MessageId messageId = pulsarTemplate.newMessage("hello-spring-messaging-message-batch-listener")
|
||||
MessageId messageId = pulsarTemplate.newMessage("hello-spring-messaging-message-batch-listener")
|
||||
.withMessageCustomizer(messageBuilder -> messageBuilder.property("foo",
|
||||
"springMessagingMessageBatchListenerWithHeaders"))
|
||||
.withTopic("springMessagingMessageBatchListenerWithHeaders").send();
|
||||
@@ -726,7 +726,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport {
|
||||
|
||||
@Test
|
||||
void pulsarMessagesBatchListenerWithHeaders() throws Exception {
|
||||
final MessageId messageId = pulsarTemplate.newMessage("hello-pulsar-messages-batch-listener")
|
||||
MessageId messageId = pulsarTemplate.newMessage("hello-pulsar-messages-batch-listener")
|
||||
.withMessageCustomizer(
|
||||
messageBuilder -> messageBuilder.property("foo", "pulsarMessagesBatchListenerWithHeaders"))
|
||||
.withTopic("pulsarMessagesBatchListenerWithHeaders").send();
|
||||
|
||||
@@ -509,7 +509,7 @@ public class ReactivePulsarListenerTests implements PulsarTestContainerSupport {
|
||||
|
||||
@Test
|
||||
void simpleListenerWithHeaders() throws Exception {
|
||||
final MessageId messageId = pulsarTemplate.newMessage("hello-simple-listener")
|
||||
MessageId messageId = pulsarTemplate.newMessage("hello-simple-listener")
|
||||
.withMessageCustomizer(
|
||||
messageBuilder -> messageBuilder.property("foo", "simpleListenerWithHeaders"))
|
||||
.withTopic("simpleListenerWithHeaders").send();
|
||||
@@ -523,7 +523,7 @@ public class ReactivePulsarListenerTests implements PulsarTestContainerSupport {
|
||||
|
||||
@Test
|
||||
void pulsarMessageListenerWithHeaders() throws Exception {
|
||||
final MessageId messageId = pulsarTemplate.newMessage("hello-pulsar-message-listener")
|
||||
MessageId messageId = pulsarTemplate.newMessage("hello-pulsar-message-listener")
|
||||
.withMessageCustomizer(
|
||||
messageBuilder -> messageBuilder.property("foo", "pulsarMessageListenerWithHeaders"))
|
||||
.withTopic("pulsarMessageListenerWithHeaders").send();
|
||||
@@ -537,7 +537,7 @@ public class ReactivePulsarListenerTests implements PulsarTestContainerSupport {
|
||||
|
||||
@Test
|
||||
void springMessagingMessageListenerWithHeaders() throws Exception {
|
||||
final MessageId messageId = pulsarTemplate.newMessage("hello-spring-messaging-message-listener")
|
||||
MessageId messageId = pulsarTemplate.newMessage("hello-spring-messaging-message-listener")
|
||||
.withMessageCustomizer(messageBuilder -> messageBuilder.property("foo",
|
||||
"springMessagingMessageListenerWithHeaders"))
|
||||
.withTopic("springMessagingMessageListenerWithHeaders").send();
|
||||
|
||||
Reference in New Issue
Block a user