Update with new MessageResult::acknowledge(Message) function (#233)
This commit is contained in:
committed by
GitHub
parent
20153dc6b9
commit
1055637ff2
@@ -102,10 +102,8 @@ class DefaultReactivePulsarMessageListenerContainerTests implements PulsarTestCo
|
||||
pulsarConsumerFactory.createConsumer(Schema.STRING).consumeNothing().block(Duration.ofSeconds(10));
|
||||
CountDownLatch latch = new CountDownLatch(5);
|
||||
ReactivePulsarContainerProperties<String> pulsarContainerProperties = new ReactivePulsarContainerProperties<>();
|
||||
pulsarContainerProperties.setMessageHandler((ReactivePulsarStreamingHandler<String>) (msg) -> msg.map(m -> {
|
||||
latch.countDown();
|
||||
return MessageResult.acknowledge(m.getMessageId());
|
||||
}));
|
||||
pulsarContainerProperties.setMessageHandler((ReactivePulsarStreamingHandler<String>) (msg) -> msg
|
||||
.doOnNext((m) -> latch.countDown()).map(MessageResult::acknowledge));
|
||||
pulsarContainerProperties.setSchema(Schema.STRING);
|
||||
DefaultReactivePulsarMessageListenerContainer<String> container = new DefaultReactivePulsarMessageListenerContainer<>(
|
||||
pulsarConsumerFactory, pulsarContainerProperties);
|
||||
@@ -292,13 +290,9 @@ class DefaultReactivePulsarMessageListenerContainerTests implements PulsarTestCo
|
||||
|
||||
CountDownLatch latch = new CountDownLatch(6);
|
||||
ReactivePulsarContainerProperties<String> pulsarContainerProperties = new ReactivePulsarContainerProperties<>();
|
||||
pulsarContainerProperties.setMessageHandler((ReactivePulsarStreamingHandler<String>) (msg) -> msg.map(m -> {
|
||||
latch.countDown();
|
||||
if (m.getValue().endsWith("4")) {
|
||||
return MessageResult.negativeAcknowledge(m.getMessageId());
|
||||
}
|
||||
return MessageResult.acknowledge(m.getMessageId());
|
||||
}));
|
||||
pulsarContainerProperties.setMessageHandler((ReactivePulsarStreamingHandler<String>) (msg) -> msg
|
||||
.doOnNext((m) -> latch.countDown()).map((m) -> m.getValue().endsWith("4")
|
||||
? MessageResult.negativeAcknowledge(m) : MessageResult.acknowledge(m)));
|
||||
pulsarContainerProperties.setSchema(Schema.STRING);
|
||||
pulsarContainerProperties.setSubscriptionType(SubscriptionType.Shared);
|
||||
DefaultReactivePulsarMessageListenerContainer<String> container = new DefaultReactivePulsarMessageListenerContainer<>(
|
||||
@@ -320,7 +314,7 @@ class DefaultReactivePulsarMessageListenerContainerTests implements PulsarTestCo
|
||||
if (message.getValue().endsWith("4")) {
|
||||
dlqLatch.countDown();
|
||||
}
|
||||
return Mono.just(MessageResult.acknowledge(message.getMessageId()));
|
||||
return Mono.just(MessageResult.acknowledge(message));
|
||||
}).block();
|
||||
|
||||
assertThat(dlqLatch.await(10, TimeUnit.SECONDS)).isTrue();
|
||||
|
||||
@@ -259,7 +259,7 @@ public class ReactivePulsarListenerTests implements PulsarTestContainerSupport {
|
||||
|
||||
@ReactivePulsarListener(topics = "streaming-1", stream = true, consumerCustomizer = "consumerCustomizer")
|
||||
Flux<MessageResult<Void>> listen1(Flux<Message<String>> messages) {
|
||||
return messages.doOnNext(m -> latch1.countDown()).map(m -> MessageResult.acknowledge(m.getMessageId()));
|
||||
return messages.doOnNext(m -> latch1.countDown()).map(MessageResult::acknowledge);
|
||||
}
|
||||
|
||||
@ReactivePulsarListener(topics = "streaming-2", stream = true, consumerCustomizer = "consumerCustomizer")
|
||||
|
||||
@@ -98,7 +98,7 @@ public class ReactiveSpringPulsarBootApp {
|
||||
public Flux<MessageResult<Void>> listenStreaming(Flux<Message<Foo>> messages) {
|
||||
return messages
|
||||
.doOnNext((msg) -> this.logger.info("Streaming reactive listener received: {}", msg.getValue()))
|
||||
.map(m -> MessageResult.acknowledge(m.getMessageId()));
|
||||
.map(MessageResult::acknowledge);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -101,7 +101,7 @@ class ReactivePulsarListenerTests implements PulsarTestContainerSupport {
|
||||
@ReactivePulsarListener(subscriptionName = "rplt-subscription2", topics = "rplt-topic2", stream = true,
|
||||
consumerCustomizer = "consumerCustomizer")
|
||||
public Flux<MessageResult<Void>> listen(Flux<Message<String>> messages) {
|
||||
return messages.doOnNext(t -> LATCH2.countDown()).map(m -> MessageResult.acknowledge(m.getMessageId()));
|
||||
return messages.doOnNext(t -> LATCH2.countDown()).map(MessageResult::acknowledge);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user