Immutable props in message listener container (#140)
* Store container props in instance vars to avoid re-retrieving the values from the mutable container props Resolves #133
This commit is contained in:
@@ -170,6 +170,12 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
|
||||
|
||||
private final PulsarConsumerErrorHandler<T> pulsarConsumerErrorHandler;
|
||||
|
||||
private final boolean isBatchListener = this.containerProperties.isBatchListener();
|
||||
|
||||
private final AckMode ackMode = this.containerProperties.getAckMode();
|
||||
|
||||
private final SubscriptionType subscriptionType = this.containerProperties.getSubscriptionType();
|
||||
|
||||
@SuppressWarnings({ "unchecked", "rawtypes" })
|
||||
Listener(MessageListener<?> messageListener) {
|
||||
if (messageListener instanceof PulsarBatchMessageListener) {
|
||||
@@ -282,7 +288,7 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
|
||||
DefaultPulsarMessageListenerContainer.this.logger.error(e, () -> "Error receiving messages.");
|
||||
}
|
||||
Assert.isTrue(messages != null, "Messages cannot be null.");
|
||||
if (this.containerProperties.isBatchListener()) {
|
||||
if (this.isBatchListener) {
|
||||
if (!inRetryMode.get() && !messagesPendingInBatch.get()) {
|
||||
messageList = new ArrayList<>();
|
||||
messages.forEach(messageList::add);
|
||||
@@ -291,13 +297,13 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
|
||||
if (messageList != null && messageList.size() > 0) {
|
||||
if (this.batchMessageListener instanceof PulsarBatchAcknowledgingMessageListener) {
|
||||
this.batchMessageListener.received(this.consumer, messageList,
|
||||
this.containerProperties.getAckMode() == AckMode.MANUAL
|
||||
this.ackMode.equals(AckMode.MANUAL)
|
||||
? new ConsumerBatchAcknowledgment(this.consumer) : null);
|
||||
}
|
||||
else {
|
||||
this.batchMessageListener.received(this.consumer, messageList);
|
||||
}
|
||||
if (this.containerProperties.getAckMode() == AckMode.BATCH) {
|
||||
if (this.ackMode.equals(AckMode.BATCH)) {
|
||||
try {
|
||||
if (isSharedSubscriptionType()) {
|
||||
this.consumer.acknowledge(messages);
|
||||
@@ -335,29 +341,26 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
|
||||
do {
|
||||
try {
|
||||
if (this.listener instanceof PulsarAcknowledgingMessageListener) {
|
||||
this.listener.received(this.consumer, message,
|
||||
this.containerProperties.getAckMode() == AckMode.MANUAL
|
||||
? new ConsumerAcknowledgment(this.consumer, message) : null);
|
||||
this.listener.received(this.consumer, message, this.ackMode.equals(AckMode.MANUAL)
|
||||
? new ConsumerAcknowledgment(this.consumer, message) : null);
|
||||
}
|
||||
else if (this.listener != null) {
|
||||
this.listener.received(this.consumer, message);
|
||||
}
|
||||
if (this.containerProperties.getAckMode() == AckMode.RECORD) {
|
||||
if (this.ackMode.equals(AckMode.RECORD)) {
|
||||
handleAck(message);
|
||||
}
|
||||
if (inRetryMode.get()) {
|
||||
inRetryMode.set(false);
|
||||
}
|
||||
inRetryMode.compareAndSet(true, false);
|
||||
}
|
||||
catch (Exception e) {
|
||||
if (this.pulsarConsumerErrorHandler != null) {
|
||||
invokeRecordListenerErrorHandler(inRetryMode, message, e);
|
||||
}
|
||||
else {
|
||||
if (this.containerProperties.getAckMode() == AckMode.RECORD) {
|
||||
if (this.ackMode.equals(AckMode.RECORD)) {
|
||||
this.consumer.negativeAcknowledge(message);
|
||||
}
|
||||
else if (this.containerProperties.getAckMode() == AckMode.BATCH) {
|
||||
else if (this.ackMode.equals(AckMode.BATCH)) {
|
||||
this.nackableMessages.add(message.getMessageId());
|
||||
}
|
||||
else {
|
||||
@@ -371,7 +374,7 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
|
||||
while (inRetryMode.get());
|
||||
}
|
||||
// All the records are processed at this point. Handle acks.
|
||||
if (this.containerProperties.getAckMode() == AckMode.BATCH) {
|
||||
if (this.ackMode.equals(AckMode.BATCH)) {
|
||||
handleAcks(messages);
|
||||
}
|
||||
}
|
||||
@@ -426,9 +429,7 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
|
||||
inRetryMode.set(true);
|
||||
}
|
||||
else {
|
||||
if (inRetryMode.get()) {
|
||||
inRetryMode.set(false);
|
||||
}
|
||||
inRetryMode.compareAndSet(true, false);
|
||||
// retries exhausted - recover the message
|
||||
this.pulsarConsumerErrorHandler.recoverMessage(this.consumer, pulsarMessage,
|
||||
pulsarBatchListenerFailedException);
|
||||
@@ -453,14 +454,12 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
|
||||
inRetryMode.set(true);
|
||||
}
|
||||
else {
|
||||
if (inRetryMode.get()) {
|
||||
inRetryMode.set(false);
|
||||
}
|
||||
inRetryMode.compareAndSet(true, false);
|
||||
// retries exhausted - recover the message
|
||||
this.pulsarConsumerErrorHandler.recoverMessage(this.consumer, message, e);
|
||||
// retries exhausted - if record ackmode, acknowledge, otherwise normal
|
||||
// batch ack at the end
|
||||
if (this.containerProperties.getAckMode() == AckMode.RECORD) {
|
||||
if (this.ackMode.equals(AckMode.RECORD)) {
|
||||
handleAck(message);
|
||||
}
|
||||
}
|
||||
@@ -468,12 +467,8 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
|
||||
|
||||
private void pendingMessagesHandledSuccessfully(AtomicBoolean inRetryMode,
|
||||
AtomicBoolean messagesPendingInBatch) {
|
||||
if (inRetryMode.get()) {
|
||||
inRetryMode.set(false);
|
||||
}
|
||||
if (messagesPendingInBatch.get()) {
|
||||
messagesPendingInBatch.set(false);
|
||||
}
|
||||
inRetryMode.compareAndSet(true, false);
|
||||
messagesPendingInBatch.compareAndSet(true, false);
|
||||
this.pulsarConsumerErrorHandler.clearMessage();
|
||||
}
|
||||
|
||||
@@ -483,8 +478,8 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
|
||||
}
|
||||
|
||||
private boolean isSharedSubscriptionType() {
|
||||
return this.containerProperties.getSubscriptionType() == SubscriptionType.Shared
|
||||
|| this.containerProperties.getSubscriptionType() == SubscriptionType.Key_Shared;
|
||||
return this.subscriptionType.equals(SubscriptionType.Shared)
|
||||
|| this.subscriptionType.equals(SubscriptionType.Key_Shared);
|
||||
}
|
||||
|
||||
private void handleAcks(Messages<T> messages) {
|
||||
|
||||
Reference in New Issue
Block a user