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 @@ -76,29 +76,6 @@ protected AbstractBaseDispatcher(Subscription subscription, ServiceConfiguration
}
}

/**
* Update Entries with the metadata of each entry.
*
* @param entries
* @return
*/
protected int updateEntryWrapperWithMetadata(EntryWrapper[] entryWrappers, List<Entry> entries) {
int totalMessages = 0;
for (int i = 0, entriesSize = entries.size(); i < entriesSize; i++) {
Entry entry = entries.get(i);
if (entry == null) {
continue;
}

ByteBuf metadataAndPayload = entry.getDataBuffer();
MessageMetadata msgMetadata = Commands.peekMessageMetadata(metadataAndPayload, subscription.toString(), -1);
EntryWrapper entryWrapper = EntryWrapper.get(entry, msgMetadata);
entryWrappers[i] = entryWrapper;
int batchSize = msgMetadata.getNumMessagesInBatch();
totalMessages += batchSize;
}
return totalMessages;
}

/**
* Filter messages that are being sent to a consumers.
Expand Down Expand Up @@ -127,7 +104,17 @@ public int filterEntriesForConsumer(List<Entry> entries, EntryBatchSizes batchSi
isReplayRead, consumer);
}

public int filterEntriesForConsumer(Optional<EntryWrapper[]> entryWrapper, int entryWrapperOffset,

/**
* Filter entries with prefetched message metadata range so that there is no need to peek metadata from Entry.
*
* @param optMetadataArray the optional message metadata array
* @param startOffset the index in `optMetadataArray` of the first Entry's message metadata
*
* @see AbstractBaseDispatcher#filterEntriesForConsumer(List, EntryBatchSizes, SendMessageInfo,
* EntryBatchIndexesAcks, ManagedCursor, boolean, Consumer)
*/
public int filterEntriesForConsumer(Optional<MessageMetadata[]> optMetadataArray, int startOffset,
List<Entry> entries, EntryBatchSizes batchSizes, SendMessageInfo sendMessageInfo,
EntryBatchIndexesAcks indexesAcks, ManagedCursor cursor, boolean isReplayRead, Consumer consumer) {
int totalMessages = 0;
Expand All @@ -142,17 +129,12 @@ public int filterEntriesForConsumer(Optional<EntryWrapper[]> entryWrapper, int e
continue;
}
ByteBuf metadataAndPayload = entry.getDataBuffer();
int entryWrapperIndex = i + entryWrapperOffset;
MessageMetadata msgMetadata = entryWrapper.isPresent() && entryWrapper.get()[entryWrapperIndex] != null
? entryWrapper.get()[entryWrapperIndex].getMetadata()
: null;
msgMetadata = msgMetadata == null
? Commands.peekMessageMetadata(metadataAndPayload, subscription.toString(), -1)
: msgMetadata;
EntryFilter.FilterResult filterResult = EntryFilter.FilterResult.ACCEPT;
final int metadataIndex = i + startOffset;
final MessageMetadata msgMetadata = optMetadataArray.map(metadataArray -> metadataArray[metadataIndex])
.orElse(Commands.peekMessageMetadata(metadataAndPayload, subscription.toString(), -1));
if (CollectionUtils.isNotEmpty(entryFilters)) {
fillContext(filterContext, msgMetadata, subscription, consumer);
filterResult = getFilterResult(filterContext, entry, entryFilters);
EntryFilter.FilterResult filterResult = getFilterResult(filterContext, entry, entryFilters);
if (filterResult == EntryFilter.FilterResult.REJECT) {
entriesToFiltered.add(entry.getPosition());
entries.set(i, null);
Expand Down

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -24,12 +24,14 @@
import java.util.Collections;
import java.util.Comparator;
import java.util.List;
import java.util.Objects;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicIntegerFieldUpdater;
import java.util.stream.Stream;
import org.apache.bookkeeper.mledger.AsyncCallbacks.ReadEntriesCallback;
import org.apache.bookkeeper.mledger.Entry;
import org.apache.bookkeeper.mledger.ManagedCursor;
Expand All @@ -47,7 +49,6 @@
import org.apache.pulsar.broker.service.Dispatcher;
import org.apache.pulsar.broker.service.EntryBatchIndexesAcks;
import org.apache.pulsar.broker.service.EntryBatchSizes;
import org.apache.pulsar.broker.service.EntryWrapper;
import org.apache.pulsar.broker.service.InMemoryRedeliveryTracker;
import org.apache.pulsar.broker.service.RedeliveryTracker;
import org.apache.pulsar.broker.service.RedeliveryTrackerDisabled;
Expand All @@ -59,6 +60,7 @@
import org.apache.pulsar.client.impl.Backoff;
import org.apache.pulsar.common.api.proto.CommandSubscribe.SubType;
import org.apache.pulsar.common.api.proto.MessageMetadata;
import org.apache.pulsar.common.protocol.Commands;
import org.apache.pulsar.common.util.Codec;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
Expand Down Expand Up @@ -512,8 +514,13 @@ protected void sendMessagesToConsumers(ReadType readType, List<Entry> entries) {
readMoreEntries();
return;
}
EntryWrapper[] entryWrappers = new EntryWrapper[entries.size()];
int remainingMessages = updateEntryWrapperWithMetadata(entryWrappers, entries);
final MessageMetadata[] metadataArray = entries.stream()
.map(entry -> Commands.peekMessageMetadata(entry.getDataBuffer(), subscription.toString(), -1))
.toArray(MessageMetadata[]::new);
int remainingMessages = Stream.of(metadataArray).filter(Objects::nonNull)
.map(MessageMetadata::getNumMessagesInBatch)
.reduce(0, Integer::sum);

int start = 0;
long totalMessagesSent = 0;
long totalBytesSent = 0;
Expand Down Expand Up @@ -564,7 +571,7 @@ protected void sendMessagesToConsumers(ReadType readType, List<Entry> entries) {

EntryBatchSizes batchSizes = EntryBatchSizes.get(entriesForThisConsumer.size());
EntryBatchIndexesAcks batchIndexesAcks = EntryBatchIndexesAcks.get(entriesForThisConsumer.size());
totalEntries += filterEntriesForConsumer(Optional.ofNullable(entryWrappers), start,
totalEntries += filterEntriesForConsumer(Optional.of(metadataArray), start,
entriesForThisConsumer, batchSizes, sendMessageInfo, batchIndexesAcks, cursor,
readType == ReadType.Replay, c);

Expand All @@ -587,13 +594,6 @@ protected void sendMessagesToConsumers(ReadType readType, List<Entry> entries) {
}
}

// release entry-wrapper
for (EntryWrapper entry : entryWrappers) {
if (entry != null) {
entry.recycle();
}
}

// acquire message-dispatch permits for already delivered messages
long permits = dispatchThrottlingOnBatchMessageEnabled ? totalEntries : totalMessagesSent;
if (serviceConfig.isDispatchThrottlingOnNonBacklogConsumerEnabled() || !cursor.isActive()) {
Expand Down