@@ -32,6 +32,7 @@ import java.util.stream.Collectors;
|
||||
import java.util.stream.Stream;
|
||||
import java.util.stream.StreamSupport;
|
||||
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
import org.apache.pulsar.client.api.BatchReceivePolicy;
|
||||
import org.apache.pulsar.client.api.Consumer;
|
||||
import org.apache.pulsar.client.api.DeadLetterPolicy;
|
||||
@@ -45,6 +46,7 @@ import org.apache.pulsar.client.api.Schema;
|
||||
import org.apache.pulsar.client.api.SubscriptionType;
|
||||
|
||||
import org.springframework.context.ApplicationEventPublisher;
|
||||
import org.springframework.core.log.LogAccessor;
|
||||
import org.springframework.core.task.AsyncTaskExecutor;
|
||||
import org.springframework.core.task.SimpleAsyncTaskExecutor;
|
||||
import org.springframework.pulsar.core.PulsarConsumerFactory;
|
||||
@@ -514,44 +516,34 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
|
||||
}
|
||||
|
||||
private void handleAck(Message<T> message) {
|
||||
try {
|
||||
this.consumer.acknowledge(message);
|
||||
}
|
||||
catch (PulsarClientException pce) {
|
||||
this.consumer.negativeAcknowledge(message);
|
||||
}
|
||||
AbstractAcknowledgement.handleAckByMessageId(this.consumer, message.getMessageId());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
private static final class ConsumerAcknowledgment implements Acknowledgement {
|
||||
private static abstract class AbstractAcknowledgement implements Acknowledgement {
|
||||
|
||||
private final Consumer<?> consumer;
|
||||
private static final LogAccessor logger = new LogAccessor(LogFactory.getLog(AbstractAcknowledgement.class));
|
||||
|
||||
private final Message<?> message;
|
||||
protected final Consumer<?> consumer;
|
||||
|
||||
ConsumerAcknowledgment(Consumer<?> consumer, Message<?> message) {
|
||||
AbstractAcknowledgement(Consumer<?> consumer) {
|
||||
this.consumer = consumer;
|
||||
this.message = message;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void acknowledge() {
|
||||
try {
|
||||
this.consumer.acknowledge(this.message);
|
||||
}
|
||||
catch (PulsarClientException e) {
|
||||
this.consumer.negativeAcknowledge(this.message);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void acknowledge(MessageId messageId) {
|
||||
handleAckByMessageId(this.consumer, messageId);
|
||||
}
|
||||
|
||||
private static void handleAckByMessageId(Consumer<?> consumer, MessageId messageId) {
|
||||
try {
|
||||
this.consumer.acknowledge(messageId);
|
||||
consumer.acknowledge(messageId);
|
||||
}
|
||||
catch (PulsarClientException e) {
|
||||
this.consumer.negativeAcknowledge(messageId);
|
||||
catch (PulsarClientException pce) {
|
||||
AbstractAcknowledgement.logger.warn(pce,
|
||||
() -> String.format("Acknowledgment failed for message: [%s]", messageId));
|
||||
consumer.negativeAcknowledge(messageId);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -562,34 +554,43 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
|
||||
}
|
||||
catch (PulsarClientException e) {
|
||||
for (MessageId messageId : messageIds) {
|
||||
try {
|
||||
this.consumer.acknowledge(messageId);
|
||||
}
|
||||
catch (PulsarClientException ex) {
|
||||
this.consumer.negativeAcknowledge(messageId);
|
||||
}
|
||||
handleAckByMessageId(this.consumer, messageId);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void nack(MessageId messageId) {
|
||||
this.consumer.negativeAcknowledge(messageId);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
private static final class ConsumerAcknowledgment extends AbstractAcknowledgement {
|
||||
|
||||
private final Message<?> message;
|
||||
|
||||
ConsumerAcknowledgment(Consumer<?> consumer, Message<?> message) {
|
||||
super(consumer);
|
||||
this.message = message;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void acknowledge() {
|
||||
acknowledge(this.message.getMessageId());
|
||||
}
|
||||
|
||||
@Override
|
||||
public void nack() {
|
||||
this.consumer.negativeAcknowledge(this.message);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void nack(MessageId messageId) {
|
||||
this.consumer.negativeAcknowledge(messageId);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
private static final class ConsumerBatchAcknowledgment implements Acknowledgement {
|
||||
|
||||
private final Consumer<?> consumer;
|
||||
private static final class ConsumerBatchAcknowledgment extends AbstractAcknowledgement {
|
||||
|
||||
ConsumerBatchAcknowledgment(Consumer<?> consumer) {
|
||||
this.consumer = consumer;
|
||||
super(consumer);
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -597,43 +598,11 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
|
||||
throw new UnsupportedOperationException();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void acknowledge(MessageId messageId) {
|
||||
try {
|
||||
this.consumer.acknowledge(messageId);
|
||||
}
|
||||
catch (PulsarClientException e) {
|
||||
this.consumer.negativeAcknowledge(messageId);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void acknowledge(List<MessageId> messageIds) {
|
||||
try {
|
||||
this.consumer.acknowledge(messageIds);
|
||||
}
|
||||
catch (PulsarClientException e) {
|
||||
for (MessageId messageId : messageIds) {
|
||||
try {
|
||||
this.consumer.acknowledge(messageId);
|
||||
}
|
||||
catch (PulsarClientException ex) {
|
||||
this.consumer.negativeAcknowledge(messageId);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void nack() {
|
||||
throw new UnsupportedOperationException();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void nack(MessageId messageId) {
|
||||
this.consumer.negativeAcknowledge(messageId);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -41,6 +41,7 @@ import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
import org.apache.pulsar.client.api.Consumer;
|
||||
import org.apache.pulsar.client.api.Message;
|
||||
import org.apache.pulsar.client.api.MessageId;
|
||||
import org.apache.pulsar.client.api.Messages;
|
||||
import org.apache.pulsar.client.api.PulsarClient;
|
||||
import org.apache.pulsar.client.api.Schema;
|
||||
@@ -85,7 +86,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
|
||||
doAnswer(invocation -> {
|
||||
latch.countDown();
|
||||
return invocation.callRealMethod();
|
||||
}).when(containerConsumer).acknowledge(any(Message.class));
|
||||
}).when(containerConsumer).acknowledge(any(MessageId.class));
|
||||
|
||||
Map<String, Object> prodConfig = new HashMap<>();
|
||||
prodConfig.put("topicName", "cons-ack-tests-011");
|
||||
@@ -166,7 +167,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
|
||||
doAnswer(invocation -> {
|
||||
ackCallCount.incrementAndGet();
|
||||
return invocation.callRealMethod();
|
||||
}).when(containerConsumer).acknowledge(any(Message.class));
|
||||
}).when(containerConsumer).acknowledge(any(MessageId.class));
|
||||
|
||||
Map<String, Object> prodConfig = new HashMap<>();
|
||||
prodConfig.put("topicName", "cons-ack-tests-013");
|
||||
@@ -187,7 +188,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
|
||||
final int ackCalls = ackCallCount.get();
|
||||
if (ackCalls < 5) {
|
||||
await().atMost(Duration.ofSeconds(10))
|
||||
.untilAsserted(() -> verify(containerConsumer, atMost(4)).acknowledge(any(Message.class)));
|
||||
.untilAsserted(() -> verify(containerConsumer, atMost(4)).acknowledge(any(MessageId.class)));
|
||||
await().atMost(Duration.ofSeconds(10)).untilAsserted(
|
||||
() -> verify(containerConsumer, atLeastOnce()).acknowledgeCumulative(any(Message.class)));
|
||||
if (ackCalls == 0) {
|
||||
@@ -236,7 +237,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
|
||||
doAnswer(invocation -> {
|
||||
latch.countDown();
|
||||
return invocation.callRealMethod();
|
||||
}).when(containerConsumer).acknowledge(any(Message.class));
|
||||
}).when(containerConsumer).acknowledge(any(MessageId.class));
|
||||
|
||||
Map<String, Object> prodConfig = new HashMap<>();
|
||||
prodConfig.put("topicName", "cons-ack-tests-014");
|
||||
@@ -251,7 +252,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
|
||||
// invocation.
|
||||
assertThat(acksObjects.size()).isEqualTo(10);
|
||||
await().atMost(Duration.ofSeconds(10))
|
||||
.untilAsserted(() -> verify(containerConsumer, times(10)).acknowledge(any(Message.class)));
|
||||
.untilAsserted(() -> verify(containerConsumer, times(10)).acknowledge(any(MessageId.class)));
|
||||
|
||||
container.stop();
|
||||
pulsarClient.close();
|
||||
|
||||
Reference in New Issue
Block a user