diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java index 7bdc5c9ef1154..f8d7098d04291 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java @@ -2292,7 +2292,8 @@ private void advanceCursorsIfNecessary(List ledgersToDelete) { // move the mark delete position to the highestPositionToDelete only if it is smaller than the add confirmed // to prevent the edge case where the cursor is caught up to the latest and highestPositionToDelete may be larger than the last add confirmed if (highestPositionToDelete.compareTo((PositionImpl) cursor.getMarkDeletedPosition()) > 0 - && highestPositionToDelete.compareTo((PositionImpl) cursor.getManagedLedger().getLastConfirmedEntry()) <= 0 ) { + && highestPositionToDelete.compareTo((PositionImpl) cursor.getManagedLedger().getLastConfirmedEntry()) <= 0 + && !(!cursor.isDurable() && cursor instanceof NonDurableCursorImpl && ((NonDurableCursorImpl) cursor).isReadCompacted())) { cursor.asyncMarkDelete(highestPositionToDelete, new MarkDeleteCallback() { @Override public void markDeleteComplete(Object ctx) { diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/NonDurableCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/NonDurableCursorImpl.java index 28ed13d808c86..263e6faa605f0 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/NonDurableCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/NonDurableCursorImpl.java @@ -33,7 +33,7 @@ public class NonDurableCursorImpl extends ManagedCursorImpl { - private final boolean readCompacted; + private volatile boolean readCompacted; NonDurableCursorImpl(BookKeeper bookkeeper, ManagedLedgerConfig config, ManagedLedgerImpl ledger, String cursorName, PositionImpl startCursorPosition, PulsarApi.CommandSubscribe.InitialPosition initialPosition, @@ -120,6 +120,10 @@ public void asyncDeleteCursor(final String consumerName, final DeleteCursorCallb callback.deleteCursorComplete(ctx); } + public void setReadCompacted(boolean readCompacted) { + this.readCompacted = readCompacted; + } + public boolean isReadCompacted() { return readCompacted; } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherSingleActiveConsumer.java index 5d486528f7f8c..0a432aa065ba4 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherSingleActiveConsumer.java @@ -27,6 +27,9 @@ import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; import java.util.concurrent.atomic.AtomicReferenceFieldUpdater; + +import org.apache.bookkeeper.mledger.ManagedCursor; +import org.apache.bookkeeper.mledger.impl.NonDurableCursorImpl; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.service.BrokerServiceException.ConsumerBusyException; import org.apache.pulsar.broker.service.BrokerServiceException.ServerMetadataException; @@ -47,7 +50,7 @@ public abstract class AbstractDispatcherSingleActiveConsumer extends AbstractBas protected boolean isKeyHashRangeFiltered = false; protected CompletableFuture closeFuture = null; protected final int partitionIndex; - + protected final ManagedCursor cursor; // This dispatcher supports both the Exclusive and Failover subscription types protected final SubType subscriptionType; @@ -59,12 +62,13 @@ public abstract class AbstractDispatcherSingleActiveConsumer extends AbstractBas protected boolean isFirstRead = true; public AbstractDispatcherSingleActiveConsumer(SubType subscriptionType, int partitionIndex, - String topicName, Subscription subscription, ServiceConfiguration serviceConfig) { + String topicName, Subscription subscription, ServiceConfiguration serviceConfig, ManagedCursor cursor) { super(subscription, serviceConfig); this.topicName = topicName; this.consumers = new CopyOnWriteArrayList<>(); this.partitionIndex = partitionIndex; this.subscriptionType = subscriptionType; + this.cursor = cursor; ACTIVE_CONSUMER_UPDATER.set(this, null); } @@ -180,7 +184,9 @@ public synchronized void addConsumer(Consumer consumer) throws BrokerServiceExce consumer.notifyActiveConsumerChange(currentActiveConsumer); } } - + if (cursor != null && !cursor.isDurable() && cursor instanceof NonDurableCursorImpl) { + ((NonDurableCursorImpl) cursor).setReadCompacted(ACTIVE_CONSUMER_UPDATER.get(this).readCompacted()); + } } public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceException { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherSingleActiveConsumer.java index 69e9a95b3ae4f..a5a3515541951 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherSingleActiveConsumer.java @@ -51,7 +51,7 @@ public final class NonPersistentDispatcherSingleActiveConsumer extends AbstractD public NonPersistentDispatcherSingleActiveConsumer(SubType subscriptionType, int partitionIndex, NonPersistentTopic topic, Subscription subscription) { super(subscriptionType, partitionIndex, topic.getName(), subscription, - topic.getBrokerService().pulsar().getConfiguration()); + topic.getBrokerService().pulsar().getConfiguration(), null); this.topic = topic; this.subscription = subscription; this.msgDrop = new Rate(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java index 00b549d31c60d..47a9cf22f14b0 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java @@ -83,7 +83,7 @@ public final class PersistentDispatcherSingleActiveConsumer extends AbstractDisp public PersistentDispatcherSingleActiveConsumer(ManagedCursor cursor, SubType subscriptionType, int partitionIndex, PersistentTopic topic, Subscription subscription) { super(subscriptionType, partitionIndex, topic.getName(), subscription, - topic.getBrokerService().pulsar().getConfiguration()); + topic.getBrokerService().pulsar().getConfiguration(), cursor); this.topic = topic; this.name = topic.getName() + " / " + (cursor.getName() != null ? Codec.decode(cursor.getName()) : ""/* NonDurableCursor doesn't have name */);