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 ec850489db8f4..02305e4d922f2 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 @@ -165,19 +165,8 @@ public void resetCloseFuture() { // noop } - public static final String NONE_KEY = "NONE_KEY"; - protected byte[] peekStickyKey(ByteBuf metadataAndPayload) { - metadataAndPayload.markReaderIndex(); - MessageMetadata metadata = Commands.parseMessageMetadata(metadataAndPayload); - metadataAndPayload.resetReaderIndex(); - byte[] key = NONE_KEY.getBytes(); - if (metadata.hasOrderingKey()) { - return metadata.getOrderingKey(); - } else if (metadata.hasPartitionKey()) { - return metadata.getPartitionKey().getBytes(); - } - return key; + return Commands.peekStickyKey(metadataAndPayload, subscription.getTopicName(), subscription.getName()); } protected void addMessageToReplay(long ledgerId, long entryId) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerEntryMetadataE2ETest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerEntryMetadataE2ETest.java new file mode 100644 index 0000000000000..7c1ca280fc199 --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerEntryMetadataE2ETest.java @@ -0,0 +1,92 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.broker.service; + +import lombok.Cleanup; +import org.apache.pulsar.client.api.Consumer; +import org.apache.pulsar.client.api.Message; +import org.apache.pulsar.client.api.Producer; +import org.apache.pulsar.client.api.SubscriptionType; +import org.assertj.core.util.Sets; +import org.testng.Assert; +import org.testng.annotations.AfterClass; +import org.testng.annotations.BeforeClass; +import org.testng.annotations.DataProvider; +import org.testng.annotations.Test; + +/** + * Test for the broker entry metadata. + */ +public class BrokerEntryMetadataE2ETest extends BrokerTestBase { + + @DataProvider(name = "subscriptionTypes") + public static Object[] subscriptionTypes() { + return new Object[] { + SubscriptionType.Exclusive, + SubscriptionType.Failover, + SubscriptionType.Shared, + SubscriptionType.Key_Shared + }; + } + + @BeforeClass + protected void setup() throws Exception { + conf.setBrokerEntryMetadataInterceptors(Sets.newTreeSet( + "org.apache.pulsar.common.intercept.AppendBrokerTimestampMetadataInterceptor", + "org.apache.pulsar.common.intercept.AppendIndexMetadataInterceptor" + )); + baseSetup(); + } + + @AfterClass(alwaysRun = true) + protected void cleanup() throws Exception { + internalCleanup(); + } + + @Test(dataProvider = "subscriptionTypes") + public void testProduceAndConsume(SubscriptionType subType) throws Exception { + final String topic = newTopicName(); + final int messages = 10; + + @Cleanup + Producer producer = pulsarClient.newProducer() + .topic(topic) + .create(); + + @Cleanup + Consumer consumer = pulsarClient.newConsumer() + .topic(topic) + .subscriptionType(subType) + .subscriptionName("my-sub") + .subscribe(); + + for (int i = 0; i < messages; i++) { + producer.send(String.valueOf(i).getBytes()); + } + + int receives = 0; + for (int i = 0; i < messages; i++) { + Message received = consumer.receive(); + ++ receives; + Assert.assertEquals(i, Integer.valueOf(new String(received.getValue())).intValue()); + } + + Assert.assertEquals(messages, receives); + } +} diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java index 50011d41edaba..938ffed77dbe0 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java @@ -44,6 +44,7 @@ import org.apache.pulsar.broker.service.Topic; import org.apache.pulsar.broker.service.persistent.PersistentStickyKeyDispatcherMultipleConsumers; import org.apache.pulsar.broker.service.persistent.PersistentSubscription; +import org.apache.pulsar.common.protocol.Commands; import org.apache.pulsar.common.schema.KeyValue; import org.apache.pulsar.common.util.Murmur3_32Hash; import org.slf4j.Logger; @@ -306,7 +307,7 @@ public void testNonKeySendAndReceiveWithHashRangeExclusiveStickyKeyConsumerSelec .value(i) .send(); } - int slot = Murmur3_32Hash.getInstance().makeHash(PersistentStickyKeyDispatcherMultipleConsumers.NONE_KEY.getBytes()) + int slot = Murmur3_32Hash.getInstance().makeHash("NONE_KEY".getBytes()) % KeySharedPolicy.DEFAULT_HASH_RANGE_SIZE; List, Integer>> checkList = new ArrayList<>(); if (slot <= 20000) { diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java index a29ceb1b285f0..55497fc6e9bd0 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java @@ -28,12 +28,12 @@ import io.netty.util.concurrent.FastThreadLocal; import java.io.IOException; +import java.nio.charset.StandardCharsets; import java.util.Collections; import java.util.List; import java.util.Map; import java.util.Optional; import java.util.Set; -import java.util.stream.Collectors; import lombok.experimental.UtilityClass; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; @@ -1654,6 +1654,7 @@ public static MessageMetadata peekMessageMetadata(ByteBuf metadataAndPayload, St try { // save the reader index and restore after parsing int readerIdx = metadataAndPayload.readerIndex(); + skipBrokerEntryMetadataIfExist(metadataAndPayload); MessageMetadata metadata = Commands.parseMessageMetadata(metadataAndPayload); metadataAndPayload.readerIndex(readerIdx); @@ -1664,6 +1665,24 @@ public static MessageMetadata peekMessageMetadata(ByteBuf metadataAndPayload, St } } + private static final byte[] NONE_KEY = "NONE_KEY".getBytes(StandardCharsets.UTF_8); + public static byte[] peekStickyKey(ByteBuf metadataAndPayload, String topic, String subscription) { + try { + int readerIdx = metadataAndPayload.readerIndex(); + skipBrokerEntryMetadataIfExist(metadataAndPayload); + MessageMetadata metadata = Commands.parseMessageMetadata(metadataAndPayload); + metadataAndPayload.readerIndex(readerIdx); + if (metadata.hasOrderingKey()) { + return metadata.getOrderingKey(); + } else if (metadata.hasPartitionKey()) { + return metadata.getPartitionKey().getBytes(StandardCharsets.UTF_8); + } + } catch (Throwable t) { + log.error("[{}] [{}] Failed to peek sticky key from the message metadata", topic, subscription, t); + } + return Commands.NONE_KEY; + } + public static int getCurrentProtocolVersion() { // Return the last ProtocolVersion enum value return ProtocolVersion.values()[ProtocolVersion.values().length - 1].getValue(); diff --git a/pulsar-common/src/main/resources/findbugsExclude.xml b/pulsar-common/src/main/resources/findbugsExclude.xml index 1706bca8422c0..df161c4b621a7 100644 --- a/pulsar-common/src/main/resources/findbugsExclude.xml +++ b/pulsar-common/src/main/resources/findbugsExclude.xml @@ -43,6 +43,9 @@ + + +