diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java index ec51fecccfe00..60a7969c5f216 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java @@ -145,9 +145,10 @@ public int filterEntriesForConsumer(Optional entryWrapper, int e MessageMetadata msgMetadata = entryWrapper.isPresent() && entryWrapper.get()[entryWrapperIndex] != null ? entryWrapper.get()[entryWrapperIndex].getMetadata() : null; - msgMetadata = msgMetadata == null - ? Commands.peekMessageMetadata(metadataAndPayload, subscription.toString(), -1) - : msgMetadata; + if (msgMetadata == null) { + msgMetadata = new MessageMetadata(); + msgMetadata.copyFrom(Commands.peekMessageMetadata(metadataAndPayload, subscription.toString(), -1)); + } if (CollectionUtils.isNotEmpty(entryFilters)) { fillContext(filterContext, msgMetadata, subscription); if (EntryFilter.FilterResult.REJECT == getFilterResult(filterContext, entry, entryFilters)) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryWrapper.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryWrapper.java index a1b4e58a74792..1e526085f3111 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryWrapper.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryWrapper.java @@ -24,7 +24,7 @@ public class EntryWrapper { private Entry entry = null; - private MessageMetadata metadata = new MessageMetadata(); + private final MessageMetadata metadata = new MessageMetadata(); private boolean hasMetadata = false; public static EntryWrapper get(Entry entry, MessageMetadata metadata) { @@ -34,7 +34,6 @@ public static EntryWrapper get(Entry entry, MessageMetadata metadata) { entryWrapper.hasMetadata = true; entryWrapper.metadata.copyFrom(metadata); } - entryWrapper.metadata.copyFrom(metadata); return entryWrapper; } @@ -64,4 +63,4 @@ public void recycle() { metadata.clear(); handle.recycle(this); } -} \ No newline at end of file +} diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/FilterContext.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/FilterContext.java index e520e1011f9ff..d74e57bcfde2a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/FilterContext.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/FilterContext.java @@ -25,12 +25,23 @@ @Data public class FilterContext { private Subscription subscription; - private MessageMetadata msgMetadata; + private final MessageMetadata msgMetadata = new MessageMetadata(); + + public FilterContext() { + } public void reset() { subscription = null; - msgMetadata = null; + msgMetadata.clear(); } public static final FilterContext FILTER_CONTEXT_DISABLED = new FilterContext(); + + public MessageMetadata getMsgMetadata() { + return this.msgMetadata; + } + + public void setMsgMetadata(MessageMetadata msgMetadata) { + this.msgMetadata.clear().copyFrom(msgMetadata); + } } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagePayloadContextImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagePayloadContextImpl.java index c4eb2643b5265..2dec1850a0d09 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagePayloadContextImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagePayloadContextImpl.java @@ -45,7 +45,7 @@ protected MessagePayloadContextImpl newObject(Handle private final Recycler.Handle recyclerHandle; private BrokerEntryMetadata brokerEntryMetadata; - private MessageMetadata messageMetadata; + private final MessageMetadata messageMetadata = new MessageMetadata(); private SingleMessageMetadata singleMessageMetadata; private MessageIdImpl messageId; private ConsumerImpl consumer; @@ -68,7 +68,7 @@ public static MessagePayloadContextImpl get(final BrokerEntryMetadata brokerEntr final MessagePayloadContextImpl context = RECYCLER.get(); context.consumerEpoch = consumerEpoch; context.brokerEntryMetadata = brokerEntryMetadata; - context.messageMetadata = messageMetadata; + context.messageMetadata.copyFrom(messageMetadata); context.singleMessageMetadata = new SingleMessageMetadata(); context.messageId = messageId; context.consumer = consumer; @@ -82,7 +82,7 @@ public static MessagePayloadContextImpl get(final BrokerEntryMetadata brokerEntr public void recycle() { brokerEntryMetadata = null; - messageMetadata = null; + messageMetadata.clear(); singleMessageMetadata = null; messageId = null; consumer = null;