Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -2292,7 +2292,8 @@ private void advanceCursorsIfNecessary(List<LedgerInfo> 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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -47,7 +50,7 @@ public abstract class AbstractDispatcherSingleActiveConsumer extends AbstractBas
protected boolean isKeyHashRangeFiltered = false;
protected CompletableFuture<Void> closeFuture = null;
protected final int partitionIndex;

protected final ManagedCursor cursor;
// This dispatcher supports both the Exclusive and Failover subscription types
protected final SubType subscriptionType;

Expand All @@ -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);
}

Expand Down Expand Up @@ -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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 */);
Expand Down