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 @@ -385,9 +385,18 @@ public Future<Void> sendMessages(final List<? extends Entry> entries,
} else {
stickyKeyHash = stickyKeyHashes.get(i);
}
boolean sendingAllowed =
pendingAcks.addPendingAckIfAllowed(entry.getLedgerId(), entry.getEntryId(), batchSize,
stickyKeyHash);
boolean sendingAllowed;
long[] ackSet = batchIndexesAcks == null ? null : batchIndexesAcks.getAckSet(i);
int remainingUnacked;
if (ackSet != null) {
remainingUnacked = BitSet.valueOf(ackSet).cardinality();
unackedMessages -= (batchSize - remainingUnacked);
} else {
remainingUnacked = batchSize;
}
sendingAllowed =
pendingAcks.addPendingAckIfAllowed(entry.getLedgerId(), entry.getEntryId(),
remainingUnacked, stickyKeyHash);
if (!sendingAllowed) {
// sending isn't allowed when pending acks doesn't accept adding the entry
// this happens when Key_Shared draining hashes contains the stickyKeyHash
Expand All @@ -401,10 +410,6 @@ public Future<Void> sendMessages(final List<? extends Entry> entries,
.attr("batchSize", batchSize)
.log("Skipping sending of entry since adding to pending acks failed");
} else {
long[] ackSet = batchIndexesAcks == null ? null : batchIndexesAcks.getAckSet(i);
if (ackSet != null) {
unackedMessages -= (batchSize - BitSet.valueOf(ackSet).cardinality());
}
log.debug()
.attr("ledgerId", entry.getLedgerId())
.attr("entryId", entry.getEntryId())
Expand Down Expand Up @@ -580,7 +585,7 @@ private CompletableFuture<Long> individualAckNormal(CommandAck ack, Map<String,
ObjectIntPair<Consumer> ackOwnerConsumerAndBatchSize =
getAckOwnerConsumerAndBatchSize(msgId.getLedgerId(), msgId.getEntryId());
Consumer ackOwnerConsumer = ackOwnerConsumerAndBatchSize.left();
long ackedCount;
long ackedCount = 0;
int batchSize = ackOwnerConsumerAndBatchSize.rightInt();
if (msgId.getAckSetsCount() > 0) {
long[] ackSets = new long[msgId.getAckSetsCount()];
Expand All @@ -596,11 +601,19 @@ private CompletableFuture<Long> individualAckNormal(CommandAck ack, Map<String,
.syncBatchPositionBitSetForPendingAck(position);
}
}
addAndGetUnAckedMsgs(ackOwnerConsumer, -(int) ackedCount);
if (ackedCount > 0) {
boolean updated = ackOwnerConsumer.updateRemainingUnacked(
position.getLedgerId(), position.getEntryId(), (int) ackedCount);
if (updated) {
addAndGetUnAckedMsgs(ackOwnerConsumer, -(int) ackedCount);
}
}
} else {
position = PositionFactory.create(msgId.getLedgerId(), msgId.getEntryId());
ackedCount = getAckedCountForMsgIdNoAckSets(batchSize, position, ackOwnerConsumer);
if (checkCanRemovePendingAcksAndHandle(ackOwnerConsumer, position, msgId)) {
IntIntPair removed = ackOwnerConsumer.removePendingAckAndGet(
position.getLedgerId(), position.getEntryId());
if (removed != null) {
ackedCount = removed.leftInt();
addAndGetUnAckedMsgs(ackOwnerConsumer, -(int) ackedCount);
updateBlockedConsumerOnUnackedMsgs(ackOwnerConsumer);
}
Expand Down Expand Up @@ -679,12 +692,22 @@ private CompletableFuture<Long> individualAckWithTransaction(CommandAck ack) {
}
AckSetStateUtil.getAckSetState(position).setAckSet(ackSets);
ackedCount = getAckedCountForTransactionAck(batchSize, ackSets);
if (ackedCount > 0) {
boolean updated = ackOwnerConsumer.updateRemainingUnacked(
position.getLedgerId(), position.getEntryId(), (int) ackedCount);
if (updated) {
addAndGetUnAckedMsgs(ackOwnerConsumer, -(int) ackedCount);
}
}
} else {
IntIntPair removed = ackOwnerConsumer.removePendingAckAndGet(
position.getLedgerId(), position.getEntryId());
if (removed != null) {
addAndGetUnAckedMsgs(ackOwnerConsumer, -removed.leftInt());
updateBlockedConsumerOnUnackedMsgs(ackOwnerConsumer);
}
}

addAndGetUnAckedMsgs(ackOwnerConsumer, -(int) ackedCount);

checkCanRemovePendingAcksAndHandle(ackOwnerConsumer, position, msgId);

checkAckValidationError(ack, position);

totalAckCount.add(ackedCount);
Expand All @@ -708,16 +731,6 @@ private CompletableFuture<Long> individualAckWithTransaction(CommandAck ack) {
return completableFuture.thenApply(__ -> totalAckCount.sum());
}

private long getAckedCountForMsgIdNoAckSets(int batchSize, Position position, Consumer consumer) {
if (isAcknowledgmentAtBatchIndexLevelEnabled && Subscription.isIndividualAckMode(subType)) {
long[] cursorAckSet = getCursorAckSet(position);
if (cursorAckSet != null) {
return getAckedCountForBatchIndexLevelEnabled(position, batchSize, EMPTY_ACK_SET, consumer);
}
}
return batchSize;
}

private long getAckedCountForBatchIndexLevelEnabled(Position position, int batchSize, long[] ackSets,
Consumer consumer) {
long ackedCount = 0;
Expand Down Expand Up @@ -747,19 +760,6 @@ private long getAckedCountForTransactionAck(int batchSize, long[] ackSets) {
return ackedCount;
}

private long getUnAckedCountForBatchIndexLevelEnabled(Position position, int batchSize) {
long unAckedCount = batchSize;
if (isAcknowledgmentAtBatchIndexLevelEnabled) {
long[] cursorAckSet = getCursorAckSet(position);
if (cursorAckSet != null) {
BitSetRecyclable cursorBitSet = BitSetRecyclable.create().resetWords(cursorAckSet);
unAckedCount = cursorBitSet.cardinality();
cursorBitSet.recycle();
}
}
return unAckedCount;
}

private void checkAckValidationError(CommandAck ack, Position position) {
if (ack.hasValidationError()) {
log.warn()
Expand All @@ -769,14 +769,6 @@ private void checkAckValidationError(CommandAck ack, Position position) {
}
}

private boolean checkCanRemovePendingAcksAndHandle(Consumer ackOwnedConsumer,
Position position, MessageIdData msgId) {
if (Subscription.isIndividualAckMode(subType) && msgId.getAckSetsCount() == 0) {
return removePendingAcks(ackOwnedConsumer, position);
}
return false;
}

/**
* Retrieves the acknowledgment owner consumer and batch size for the specified ledgerId and entryId.
*
Expand Down Expand Up @@ -1121,6 +1113,37 @@ public PendingAcksMap getPendingAcks() {
return pendingAcks;
}

/**
* Atomically decrement the remaining unacked count for the specified position
* by the given acknowledged delta.
*
* <p>No-op if {@code pendingAcks} is not initialized.
*
* @return {@code true} if the update succeeds or pendingAcks is null;
* {@code false} otherwise
*/
public boolean updateRemainingUnacked(long ledgerId, long entryId, int ackedDelta) {
if (pendingAcks != null) {
return pendingAcks.updateRemainingUnacked(ledgerId, entryId, ackedDelta);
}
return true;
}

/**
* Atomically remove the pending ack entry and return its stored values.
*
* <p>No-op if {@code pendingAcks} is not initialized.
*
* @return the removed {@link IntIntPair#leftInt() remainingUnacked} and
* {@link IntIntPair#rightInt() stickyKeyHash}, or {@code null} if not found
*/
public IntIntPair removePendingAckAndGet(long ledgerId, long entryId) {
if (pendingAcks != null) {
return pendingAcks.removeAndGet(ledgerId, entryId);
}
return null;
}

/**
* Remove all pending acks up to the given mark-delete position and decrement the consumer's unacked message
* counter by the remaining unacked count for each removed entry.
Expand All @@ -1139,9 +1162,8 @@ public void removePendingAcksUpToPositionAndDecrementUnacked(long markDeleteLedg

MutableInt mutableTotalUnacked = new MutableInt(0);
pendingAcks.removeAllUpTo(markDeleteLedgerId, markDeleteEntryId,
(ledgerId, entryId, batchSize, stickyKeyHash) -> {
mutableTotalUnacked.add((int) getUnAckedCountForBatchIndexLevelEnabled(
PositionFactory.create(ledgerId, entryId), batchSize));
(ledgerId, entryId, remainingUnacked, stickyKeyHash) -> {
mutableTotalUnacked.add(remainingUnacked);
});
int totalUnacked = mutableTotalUnacked.intValue();
if (totalUnacked > 0) {
Expand All @@ -1160,11 +1182,8 @@ public void redeliverUnacknowledgedMessages(long consumerEpoch) {
if (pendingAcks != null) {
List<Position> pendingPositions = new ArrayList<>((int) pendingAcks.size());
MutableInt totalRedeliveryMessages = new MutableInt(0);
pendingAcks.forEachAndClear((ledgerId, entryId, batchSize, stickyKeyHash) -> {
int unAckedCount =
(int) getUnAckedCountForBatchIndexLevelEnabled(PositionFactory.create(ledgerId, entryId),
batchSize);
totalRedeliveryMessages.add(unAckedCount);
pendingAcks.forEachAndClear((ledgerId, entryId, remainingUnacked, stickyKeyHash) -> {
totalRedeliveryMessages.add(remainingUnacked);
pendingPositions.add(PositionFactory.create(ledgerId, entryId));
});

Expand Down Expand Up @@ -1193,8 +1212,7 @@ public void redeliverUnacknowledgedMessages(List<MessageIdData> messageIds) {
Position position = PositionFactory.create(msg.getLedgerId(), msg.getEntryId());
IntIntPair pendingAck = pendingAcks.removeAndGet(position.getLedgerId(), position.getEntryId());
if (pendingAck != null) {
int unAckedCount = (int) getUnAckedCountForBatchIndexLevelEnabled(position, pendingAck.leftInt());
totalRedeliveryMessages += unAckedCount;
totalRedeliveryMessages += pendingAck.leftInt();
pendingPositions.add(position);
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -131,6 +131,15 @@ default void cursorIsReset() {
//No-op
}

/**
* This hook is invoked after cursor mark-delete operations triggered by
* message removal flows such as expiry, skip, or clear backlog, but not for
* regular ack-driven mark-delete operations due to their higher frequency.
*
* <p>Since the cursor ack set may no longer be available after mark-delete,
* the cleanup logic relies on the remaining unacked count stored in
* {@code PendingAcksMap} entries.
*/
default void markDeletePositionMoveForward() {
// No-op
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -87,12 +87,13 @@ public interface PendingAcksConsumer {
/**
* Accept a pending acknowledgment.
*
* @param ledgerId the ledger ID
* @param entryId the entry ID
* @param batchSize the batch size
* @param stickyKeyHash the sticky key hash
* @param ledgerId the ledger ID
* @param entryId the entry ID
* @param remainingUnacked the number of remaining unacked messages in this entry
* (accounts for batch index level acknowledgments)
* @param stickyKeyHash the sticky key hash
*/
void accept(long ledgerId, long entryId, int batchSize, int stickyKeyHash);
void accept(long ledgerId, long entryId, int remainingUnacked, int stickyKeyHash);
}

private final Consumer consumer;
Expand Down Expand Up @@ -122,11 +123,12 @@ public interface PendingAcksConsumer {
*
* @param ledgerId the ledger ID
* @param entryId the entry ID
* @param batchSize the batch size
* @param remainingUnacked the number of remaining unacked messages in this entry
* (for batch entries with some indexes already acked, this may be less than batchSize)
* @param stickyKeyHash the sticky key hash
* @return true if the pending ack was added, and it's allowed to send a message, false otherwise
*/
public boolean addPendingAckIfAllowed(long ledgerId, long entryId, int batchSize, int stickyKeyHash) {
public boolean addPendingAckIfAllowed(long ledgerId, long entryId, int remainingUnacked, int stickyKeyHash) {
try {
writeLock.lock();
// prevent adding sticky hash to pending acks if the PendingAcksMap has already been closed
Expand All @@ -143,7 +145,7 @@ public boolean addPendingAckIfAllowed(long ledgerId, long entryId, int batchSize
}
TreeMap<Long, IntIntPair> ledgerPendingAcks =
pendingAcks.computeIfAbsent(ledgerId, k -> new TreeMap<>());
ledgerPendingAcks.put(entryId, IntIntPair.of(batchSize, stickyKeyHash));
ledgerPendingAcks.put(entryId, IntIntPair.of(remainingUnacked, stickyKeyHash));
return true;
} finally {
writeLock.unlock();
Expand Down Expand Up @@ -311,6 +313,34 @@ public boolean remove(long ledgerId, long entryId, int batchSize, int stickyKeyH
}
}

/**
* Atomically update the remaining unacked count for a pending ack entry by subtracting the given delta.
* Called from the ack handler after computing the number of batch indexes acknowledged in a partial ack.
*
* @param ledgerId the ledger ID
* @param entryId the entry ID
* @param ackedDelta the number of batch indexes that were just acknowledged
* @return true if the entry was found and updated, false otherwise
*/
public boolean updateRemainingUnacked(long ledgerId, long entryId, int ackedDelta) {
try {
writeLock.lock();
TreeMap<Long, IntIntPair> ledgerMap = pendingAcks.get(ledgerId);
if (ledgerMap == null) {
return false;
}
IntIntPair current = ledgerMap.get(entryId);
if (current == null) {
return false;
}
int newRemaining = current.leftInt() - ackedDelta;
ledgerMap.put(entryId, IntIntPair.of(newRemaining, current.rightInt()));
return true;
} finally {
writeLock.unlock();
}
}

/**
* Remove the pending ack for the given ledger ID and entry ID.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -145,7 +145,6 @@ public class PersistentDispatcherMultipleConsumers extends AbstractPersistentDis
protected enum ReadType {
Normal, Replay
}
private Position lastMarkDeletePositionBeforeReadMoreEntries;
private volatile long readMoreEntriesCallCount;

public PersistentDispatcherMultipleConsumers(PersistentTopic topic, ManagedCursor cursor,
Expand Down Expand Up @@ -362,17 +361,6 @@ public synchronized void readMoreEntries() {
// increment the counter for readMoreEntries calls, to track the number of times readMoreEntries is called
readMoreEntriesCallCount++;

// remove possible expired messages from redelivery tracker and pending acks
Position markDeletePosition = cursor.getMarkDeletedPosition();
if (lastMarkDeletePositionBeforeReadMoreEntries != markDeletePosition) {
Comment thread
lhotari marked this conversation as resolved.
redeliveryMessages.removeAllUpTo(markDeletePosition.getLedgerId(), markDeletePosition.getEntryId());
for (Consumer consumer : consumerList) {
consumer.removePendingAcksUpToPositionAndDecrementUnacked(
markDeletePosition.getLedgerId(), markDeletePosition.getEntryId());
}
lastMarkDeletePositionBeforeReadMoreEntries = markDeletePosition;
}

// totalAvailablePermits may be updated by other threads
int firstAvailableConsumerPermits = getFirstAvailableConsumerPermits();
int currentTotalAvailablePermits = Math.max(totalAvailablePermits, firstAvailableConsumerPermits);
Expand Down Expand Up @@ -600,6 +588,18 @@ public CopyOnWriteArrayList<Consumer> getConsumers() {
return consumerList;
}

@Override
public void markDeletePositionMoveForward() {
Position markDeletePosition = cursor.getMarkDeletedPosition();
if (markDeletePosition != null) {
redeliveryMessages.removeAllUpTo(markDeletePosition.getLedgerId(), markDeletePosition.getEntryId());
for (Consumer consumer : consumerList) {
consumer.removePendingAcksUpToPositionAndDecrementUnacked(
markDeletePosition.getLedgerId(), markDeletePosition.getEntryId());
}
}
}

@Override
public synchronized boolean canUnsubscribe(Consumer consumer) {
return consumerList.size() == 1 && consumerSet.contains(consumer);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,9 +38,9 @@
import org.apache.bookkeeper.mledger.PositionFactory;
import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl;
import org.apache.bookkeeper.mledger.proto.ManagedLedgerInfo;
import org.apache.pulsar.broker.service.Dispatcher;
import org.apache.pulsar.broker.service.MessageExpirer;
import org.apache.pulsar.client.impl.MessageImpl;
import org.apache.pulsar.common.api.proto.CommandSubscribe.SubType;
import org.apache.pulsar.common.stats.Rate;
import org.jspecify.annotations.Nullable;
/**
Expand Down Expand Up @@ -211,9 +211,11 @@ public void markDeleteComplete(Object ctx) {
long numMessagesExpired = (long) ctx - cursor.getNumberOfEntriesInBacklog(false);
msgExpired.recordMultipleEvents(numMessagesExpired, 0 /* no value stats */);
totalMsgExpired.add(numMessagesExpired);
// If the subscription is a Key_Shared subscription, we should to trigger message dispatch.
if (subscription != null && subscription.getType() == SubType.Key_Shared) {
subscription.getDispatcher().markDeletePositionMoveForward();
if (subscription != null) {
Dispatcher dispatcher = subscription.getDispatcher();
if (dispatcher != null) {
dispatcher.markDeletePositionMoveForward();
}
}
expirationCheckInProgress = FALSE;
log.debug()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -642,6 +642,7 @@ protected int getStickyKeyHash(Entry entry) {

@Override
public void markDeletePositionMoveForward() {
super.markDeletePositionMoveForward();
// reschedule a read with a backoff after moving the mark-delete position forward since there might have
// been consumers that were blocked by hash and couldn't make progress
reScheduleReadWithKeySharedUnblockingInterval();
Expand Down
Loading
Loading