diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index 1a33fa53bf8e0..0dfa56f4c74d5 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -1144,6 +1144,8 @@ protected void handleProducer(final CommandProducer cmdProducer) { final Optional topicEpoch = cmdProducer.hasTopicEpoch() ? Optional.of(cmdProducer.getTopicEpoch()) : Optional.empty(); final boolean isTxnEnabled = cmdProducer.isTxnEnabled(); + final String initialSubscriptionName = + cmdProducer.hasInitialSubscriptionName() ? cmdProducer.getInitialSubscriptionName() : null; final boolean supportsPartialProducer = supportsPartialProducer(); TopicName topicName = validateTopicName(cmdProducer.getTopic(), requestId, cmdProducer); @@ -1160,8 +1162,15 @@ protected void handleProducer(final CommandProducer cmdProducer) { } CompletableFuture isAuthorizedFuture = isTopicOperationAllowed( - topicName, TopicOperation.PRODUCE + topicName, TopicOperation.PRODUCE ); + + if (!Strings.isNullOrEmpty(initialSubscriptionName)) { + isAuthorizedFuture = + isAuthorizedFuture.thenCombine(isTopicOperationAllowed(topicName, TopicOperation.SUBSCRIBE), + (canProduce, canSubscribe) -> canProduce && canSubscribe); + } + isAuthorizedFuture.thenApply(isAuthorized -> { if (!isAuthorized) { String msg = "Client is not authorized to Produce"; @@ -1245,9 +1254,33 @@ protected void handleProducer(final CommandProducer cmdProducer) { schemaVersionFuture.thenAccept(schemaVersion -> { topic.checkIfTransactionBufferRecoverCompletely(isTxnEnabled).thenAccept(future -> { - buildProducerAndAddTopic(topic, producerId, producerName, requestId, isEncrypted, + CompletableFuture createInitSubFuture; + if (!Strings.isNullOrEmpty(initialSubscriptionName) + && topic.isPersistent() + && !topic.getSubscriptions().containsKey(initialSubscriptionName)) { + createInitSubFuture = + topic.createSubscription(initialSubscriptionName, InitialPosition.Earliest, + false); + } else { + createInitSubFuture = CompletableFuture.completedFuture(null); + } + + createInitSubFuture.whenComplete((sub, ex) -> { + if (ex != null) { + String msg = + "Failed to create the initial subscription: " + ex.getCause().getMessage(); + log.warn("[{}] {} initialSubscriptionName: {}, topic: {}", + remoteAddress, msg, initialSubscriptionName, topicName); + commandSender.sendErrorResponse(requestId, + BrokerServiceException.getClientErrorCode(ex), msg); + producers.remove(producerId, producerFuture); + return; + } + + buildProducerAndAddTopic(topic, producerId, producerName, requestId, isEncrypted, metadata, schemaVersion, epoch, userProvidedProducerName, topicName, producerAccessMode, topicEpoch, supportsPartialProducer, producerFuture); + }); }).exceptionally(exception -> { Throwable cause = exception.getCause(); log.error("producerId {}, requestId {} : TransactionBuffer recover failed", diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/AuthorizationProducerConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/AuthorizationProducerConsumerTest.java index cad210777f6b8..e01ed66cd2146 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/AuthorizationProducerConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/AuthorizationProducerConsumerTest.java @@ -51,6 +51,7 @@ import org.apache.pulsar.broker.resources.PulsarResources; import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.impl.MessageIdImpl; +import org.apache.pulsar.client.impl.ProducerBuilderImpl; import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.policies.data.AuthAction; @@ -538,6 +539,67 @@ public void testAuthData() throws Exception { log.info("-- Exiting {} test --", methodName); } + @Test + public void testPermissionForProducerCreateInitialSubscription() throws Exception { + log.info("-- Starting {} test --", methodName); + + conf.setAuthorizationProvider(PulsarAuthorizationProvider.class.getName()); + setup(); + + String lookupUrl = pulsar.getBrokerServiceUrl(); + + final String invalidRole = "invalid-role"; + final String producerRole = "producer-role"; + final String topic = "persistent://my-property/my-ns/my-topic"; + final String initialSubscriptionName = "init-sub"; + TopicName tn = TopicName.get(topic); + Authentication adminAuthentication = new ClientAuthentication("superUser"); + Authentication authenticationInvalidRole = new ClientAuthentication(invalidRole); + Authentication authenticationProducerRole = new ClientAuthentication(producerRole); + @Cleanup + PulsarAdmin admin = + PulsarAdmin.builder().serviceHttpUrl(brokerUrl.toString()).authentication(adminAuthentication).build(); + + admin.clusters().createCluster("test", ClusterData.builder().serviceUrl(brokerUrl.toString()).build()); + admin.tenants().createTenant("my-property", + new TenantInfoImpl(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); + admin.namespaces().createNamespace("my-property/my-ns", Sets.newHashSet("test")); + admin.topics().grantPermission(topic, invalidRole, Collections.singleton(AuthAction.produce)); + admin.topics().grantPermission(topic, producerRole, Sets.newHashSet(AuthAction.produce, AuthAction.consume)); + + @Cleanup + PulsarClient pulsarClientInvalidRole = PulsarClient.builder().serviceUrl(lookupUrl) + .authentication(authenticationInvalidRole).build(); + + try { + Producer invalidRoleProducer = ((ProducerBuilderImpl) pulsarClientInvalidRole.newProducer()) + .initialSubscriptionName(initialSubscriptionName) + .topic(topic) + .create(); + invalidRoleProducer.close(); + fail("Should not pass"); + } catch (PulsarClientException.AuthorizationException ex) { + // ok + } + + // If the producer doesn't have permission to create the init sub, we should also avoid creating the topic. + Assert.assertFalse(admin.namespaces().getTopics(tn.getNamespace()).contains(tn.getLocalName())); + + @Cleanup + PulsarClient pulsarClientProducerRole = PulsarClient.builder().serviceUrl(lookupUrl) + .authentication(authenticationProducerRole).build(); + + Producer producer = ((ProducerBuilderImpl) pulsarClientProducerRole.newProducer()) + .initialSubscriptionName(initialSubscriptionName) + .topic(topic) + .create(); + producer.close(); + + Assert.assertTrue(admin.topics().getSubscriptions(topic).contains(initialSubscriptionName)); + + log.info("-- Exiting {} test --", methodName); + } + public static class ClientAuthentication implements Authentication { String user; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/DeadLetterTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/DeadLetterTopicTest.java index 645c7674b00c5..03a29e1884eda 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/DeadLetterTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/DeadLetterTopicTest.java @@ -18,10 +18,12 @@ */ package org.apache.pulsar.client.api; +import java.time.Duration; import java.util.ArrayList; import java.util.HashSet; import java.util.List; import java.util.Set; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; @@ -32,6 +34,7 @@ import org.apache.pulsar.client.api.schema.GenericRecord; import lombok.Cleanup; import org.apache.pulsar.client.util.RetryMessageUtil; +import org.awaitility.Awaitility; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.testng.annotations.AfterMethod; @@ -615,4 +618,180 @@ public void testDeadLetterTopicUnderPartitionedTopicWithKeyShareType() throws Ex checkConsumer.close(); } + + @Test + public void testDeadLetterTopicWithInitialSubscription() throws Exception { + final String topic = "persistent://my-property/my-ns/dead-letter-topic"; + + final int maxRedeliveryCount = 1; + + final int sendMessages = 100; + + final String subscriptionName = "my-subscription"; + final String dlqInitialSub = "init-sub"; + + Consumer consumer = pulsarClient.newConsumer(Schema.BYTES) + .topic(topic) + .subscriptionName(subscriptionName) + .subscriptionType(SubscriptionType.Shared) + .ackTimeout(1, TimeUnit.SECONDS) + .deadLetterPolicy(DeadLetterPolicy.builder() + .maxRedeliverCount(maxRedeliveryCount) + .initialSubscriptionName(dlqInitialSub) + .build()) + .receiverQueueSize(100) + .subscriptionInitialPosition(SubscriptionInitialPosition.Earliest) + .subscribe(); + + @Cleanup + PulsarClient newPulsarClient = newPulsarClient(lookupUrl.toString(), 0);// Creates new client connection + + Producer producer = pulsarClient.newProducer(Schema.BYTES) + .topic(topic) + .create(); + + for (int i = 0; i < sendMessages; i++) { + producer.send(String.format("Hello Pulsar [%d]", i).getBytes()); + } + + producer.close(); + + int totalReceived = 0; + do { + Message message = consumer.receive(3, TimeUnit.SECONDS); + if (message == null) { + break; + } + log.info("consumer received message : {} {}", message.getMessageId(), new String(message.getData())); + totalReceived++; + } while (totalReceived < sendMessages * (maxRedeliveryCount + 1)); + + String deadLetterTopic = "persistent://my-property/my-ns/dead-letter-topic-my-subscription-DLQ"; + Awaitility.await().atMost(Duration.ofSeconds(10)).pollInterval(Duration.ofSeconds(1)).untilAsserted(() -> { + assertTrue(admin.namespaces().getTopics("my-property/my-ns").contains(deadLetterTopic)); + assertTrue(admin.topics().getSubscriptions(deadLetterTopic).contains(dlqInitialSub)); + }); + + Consumer deadLetterConsumer = newPulsarClient.newConsumer(Schema.BYTES) + .topic(deadLetterTopic) + .subscriptionName(dlqInitialSub) + .subscribe(); + + int totalInDeadLetter = 0; + do { + Message message = deadLetterConsumer.receive(10, TimeUnit.SECONDS); + assertNotNull(message, "Dead letter consumer can not receive messages."); + log.info("dead letter consumer received message : {} {}", message.getMessageId(), + new String(message.getData())); + deadLetterConsumer.acknowledge(message); + totalInDeadLetter++; + } while (totalInDeadLetter < sendMessages); + + deadLetterConsumer.close(); + consumer.close(); + } + + private CompletableFuture consumerReceiveForDLQ(Consumer consumer, AtomicInteger totalReceived, + int sendMessages, int maxRedeliveryCount) { + return CompletableFuture.runAsync(() -> { + while (true) { + Message message; + try { + message = consumer.receive(3, TimeUnit.SECONDS); + } catch (PulsarClientException e) { + log.info("fail while receiving messages: {}", e.getMessage()); + break; + } + if (message == null) { + break; + } + log.info("consumer received message : {} {}", message.getMessageId(), new String(message.getData())); + totalReceived.incrementAndGet(); + if (totalReceived.get() >= sendMessages * (maxRedeliveryCount + 1)) { + break; + } + } + }); + } + + @Test + public void testDeadLetterTopicWithInitialSubscriptionAndMultiConsumers() throws Exception { + final String topic = "persistent://my-property/my-ns/dead-letter-topic"; + + final int maxRedeliveryCount = 1; + + final int sendMessages = 100; + + final String subscriptionName = "my-subscription"; + final String dlqInitialSub = "init-sub"; + + Consumer consumer = pulsarClient.newConsumer(Schema.BYTES) + .topic(topic) + .subscriptionName(subscriptionName) + .subscriptionType(SubscriptionType.Shared) + .ackTimeout(1, TimeUnit.SECONDS) + .deadLetterPolicy(DeadLetterPolicy.builder() + .maxRedeliverCount(maxRedeliveryCount) + .initialSubscriptionName(dlqInitialSub) + .build()) + .receiverQueueSize(100) + .subscriptionInitialPosition(SubscriptionInitialPosition.Earliest) + .subscribe(); + + Consumer otherConsumer = pulsarClient.newConsumer(Schema.BYTES) + .topic(topic) + .subscriptionName(subscriptionName) + .subscriptionType(SubscriptionType.Shared) + .ackTimeout(1, TimeUnit.SECONDS) + .deadLetterPolicy(DeadLetterPolicy.builder() + .maxRedeliverCount(maxRedeliveryCount) + .initialSubscriptionName(dlqInitialSub) + .build()) + .receiverQueueSize(100) + .subscriptionInitialPosition(SubscriptionInitialPosition.Earliest) + .subscribe(); + + @Cleanup + PulsarClient newPulsarClient = newPulsarClient(lookupUrl.toString(), 0);// Creates new client connection + + Producer producer = pulsarClient.newProducer(Schema.BYTES) + .topic(topic) + .create(); + + for (int i = 0; i < sendMessages; i++) { + producer.send(String.format("Hello Pulsar [%d]", i).getBytes()); + } + + producer.close(); + + final AtomicInteger totalReceived = new AtomicInteger(0); + CompletableFuture.allOf(consumerReceiveForDLQ(consumer, totalReceived, sendMessages, maxRedeliveryCount), + consumerReceiveForDLQ(otherConsumer, totalReceived, sendMessages, maxRedeliveryCount)) + .get(10, TimeUnit.SECONDS); + + String deadLetterTopic = "persistent://my-property/my-ns/dead-letter-topic-my-subscription-DLQ"; + Awaitility.await().atMost(Duration.ofSeconds(10)).pollInterval(Duration.ofSeconds(1)).untilAsserted(() -> { + assertTrue(admin.namespaces().getTopics("my-property/my-ns").contains(deadLetterTopic)); + assertTrue(admin.topics().getSubscriptions(deadLetterTopic).contains(dlqInitialSub)); + }); + + Consumer deadLetterConsumer = newPulsarClient.newConsumer(Schema.BYTES) + .topic(deadLetterTopic) + .subscriptionName(dlqInitialSub) + .subscribe(); + + int totalInDeadLetter = 0; + do { + Message message = deadLetterConsumer.receive(10, TimeUnit.SECONDS); + assertNotNull(message, "Dead letter consumer can not receive messages."); + log.info("dead letter consumer received message : {} {}", message.getMessageId(), + new String(message.getData())); + deadLetterConsumer.acknowledge(message); + totalInDeadLetter++; + } while (totalInDeadLetter < sendMessages); + + deadLetterConsumer.close(); + otherConsumer.close(); + consumer.close(); + } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerCreationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerCreationTest.java index 0d24a4b624875..3e43e72f47ec2 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerCreationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerCreationTest.java @@ -18,6 +18,8 @@ */ package org.apache.pulsar.client.api; +import org.apache.pulsar.client.admin.PulsarAdminException; +import org.apache.pulsar.client.impl.ProducerBuilderImpl; import org.apache.pulsar.client.impl.ProducerImpl; import org.apache.pulsar.common.naming.TopicDomain; import org.apache.pulsar.common.naming.TopicName; @@ -93,4 +95,76 @@ public void testGeneratedNameProducerReconnect(TopicDomain domain) throws Pulsar Assert.assertEquals(producer.getConnectionHandler().getEpoch(), 1); Assert.assertTrue(producer.isConnected()); } + + @Test(dataProvider = "topicDomainProvider") + public void testInitialSubscriptionCreation(TopicDomain domain) throws PulsarClientException, PulsarAdminException { + final String initialSubscriptionName = "init-sub"; + final TopicName topic = TopicName.get(domain.value(), "public", "default", "testInitialSubscriptionCreation"); + + // Should not create initial subscription when the initialSubscriptionName is null or empty + Producer nullInitSubProducer = ((ProducerBuilderImpl) pulsarClient.newProducer()) + .initialSubscriptionName(null) + .topic(topic.toString()) + .create(); + nullInitSubProducer.close(); + Assert.assertFalse(admin.topics().getSubscriptions(topic.toString()).contains(initialSubscriptionName)); + + Producer emptyInitSubProducer = ((ProducerBuilderImpl) pulsarClient.newProducer()) + .initialSubscriptionName("") + .topic(topic.toString()) + .create(); + emptyInitSubProducer.close(); + Assert.assertFalse(admin.topics().getSubscriptions(topic.toString()).contains(initialSubscriptionName)); + + Producer producer = ((ProducerBuilderImpl) pulsarClient.newProducer()) + .initialSubscriptionName(initialSubscriptionName) + .topic(topic.toString()) + .create(); + producer.close(); + + // Initial subscription will only be created if the topic is persistent + Assert.assertEquals(topic.isPersistent(), + admin.topics().getSubscriptions(topic.toString()).contains(initialSubscriptionName)); + + // Existing subscription should not fail the producer creation. + Producer otherProducer = ((ProducerBuilderImpl) pulsarClient.newProducer()) + .initialSubscriptionName(initialSubscriptionName) + .topic(topic.toString()) + .create(); + otherProducer.close(); + + Assert.assertEquals(topic.isPersistent(), + admin.topics().getSubscriptions(topic.toString()).contains(initialSubscriptionName)); + } + + @Test + public void testCreateInitialSubscriptionOnPartitionedTopic() throws PulsarAdminException, PulsarClientException { + final TopicName topic = + TopicName.get("persistent", "public", "default", "testCreateInitialSubscriptionOnPartitionedTopic"); + final String initialSubscriptionName = "init-sub"; + admin.topics().createPartitionedTopic(topic.toString(), 10); + Producer producer = ((ProducerBuilderImpl) pulsarClient.newProducer()) + .initialSubscriptionName(initialSubscriptionName) + .topic(topic.toString()) + .create(); + producer.close(); + + Assert.assertTrue(admin.topics().getSubscriptions(topic.toString()).contains(initialSubscriptionName)); + } + + @Test + public void testCreateInitialSubscriptionWhenExisting() throws PulsarClientException, PulsarAdminException { + final TopicName topic = + TopicName.get("persistent", "public", "default", "testCreateInitialSubscriptionWhenExisting"); + final String initialSubscriptionName = "init-sub"; + admin.topics().createNonPartitionedTopic(topic.toString()); + admin.topics().createSubscription(topic.toString(), initialSubscriptionName, MessageId.earliest); + Producer producer = ((ProducerBuilderImpl) pulsarClient.newProducer()) + .initialSubscriptionName(initialSubscriptionName) + .topic(topic.toString()) + .create(); + producer.close(); + + Assert.assertTrue(admin.topics().getSubscriptions(topic.toString()).contains(initialSubscriptionName)); + } } diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/DeadLetterPolicy.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/DeadLetterPolicy.java index 279a5c508f5a1..351768f0e7bfd 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/DeadLetterPolicy.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/DeadLetterPolicy.java @@ -49,4 +49,9 @@ public class DeadLetterPolicy { */ private String deadLetterTopic; + /** + * Name of the initial subscription name of the dead letter topic. + * If this field is not set, the initial subscription for the dead letter topic will not be created. + */ + private String initialSubscriptionName; } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java index bfa81b6e1c3e6..60f9ca7922a8a 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java @@ -353,6 +353,11 @@ protected ConsumerImpl(PulsarClientImpl client, String topic, ConsumerConfigurat topic, subscription)); } + if (StringUtils.isNotBlank(conf.getDeadLetterPolicy().getInitialSubscriptionName())) { + this.deadLetterPolicy.setInitialSubscriptionName( + conf.getDeadLetterPolicy().getInitialSubscriptionName()); + } + } else { deadLetterPolicy = null; possibleSendToDeadLetterTopicMessages = null; @@ -1897,10 +1902,12 @@ private void initDeadLetterProducerIfNeeded() { createProducerLock.writeLock().lock(); try { if (deadLetterProducer == null) { - deadLetterProducer = client.newProducer(Schema.AUTO_PRODUCE_BYTES(schema)) - .topic(this.deadLetterPolicy.getDeadLetterTopic()) - .blockIfQueueFull(false) - .createAsync(); + deadLetterProducer = + ((ProducerBuilderImpl) client.newProducer(Schema.AUTO_PRODUCE_BYTES(schema))) + .initialSubscriptionName(this.deadLetterPolicy.getInitialSubscriptionName()) + .topic(this.deadLetterPolicy.getDeadLetterTopic()) + .blockIfQueueFull(false) + .createAsync(); } } finally { createProducerLock.writeLock().unlock(); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBuilderImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBuilderImpl.java index 989ae51d80a51..bd533fa9f576a 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBuilderImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBuilderImpl.java @@ -325,6 +325,19 @@ public ProducerBuilder enableLazyStartPartitionedProducers(boolean lazyStartP return this; } + /** + * Use this config to automatically create an initial subscription when creating the topic. + * If this field is not set, the initial subscription will not be created. + * This method is limited to internal use + * + * @param initialSubscriptionName Name of the initial subscription of the topic. + * @return the producer builder implementation instance + */ + public ProducerBuilderImpl initialSubscriptionName(String initialSubscriptionName) { + conf.setInitialSubscriptionName(initialSubscriptionName); + return this; + } + private void setMessageRoutingMode() throws PulsarClientException { if (conf.getMessageRoutingMode() == null && conf.getCustomMessageRouter() == null) { messageRoutingMode(MessageRoutingMode.RoundRobinPartition); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java index 38333b1e6742c..da3786db3ac8e 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java @@ -1532,8 +1532,9 @@ public void connectionOpened(final ClientCnx cnx) { cnx.sendRequestWithId( Commands.newProducer(topic, producerId, requestId, producerName, conf.isEncryptionEnabled(), metadata, - schemaInfo, epoch, userProvidedProducerName, - conf.getAccessMode(), topicEpoch, client.conf.isEnableTransaction()), + schemaInfo, epoch, userProvidedProducerName, + conf.getAccessMode(), topicEpoch, client.conf.isEnableTransaction(), + conf.getInitialSubscriptionName()), requestId).thenAccept(response -> { String producerName = response.getProducerName(); long lastSequenceId = response.getLastSequenceId(); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java index be12433e05294..958c6ce844d34 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java @@ -101,6 +101,8 @@ public class ProducerConfigurationData implements Serializable, Cloneable { private SortedMap properties = new TreeMap<>(); + private String initialSubscriptionName = null; + /** * * Returns true if encryption keys are added. 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 3a130aeae4d0e..12bfba127bce3 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 @@ -22,6 +22,7 @@ import static com.scurrilous.circe.checksum.Crc32cIntChecksum.resumeChecksum; import static java.nio.charset.StandardCharsets.UTF_8; import com.google.common.annotations.VisibleForTesting; +import com.google.common.base.Strings; import io.netty.buffer.ByteBuf; import io.netty.buffer.CompositeByteBuf; import io.netty.buffer.Unpooled; @@ -764,9 +765,19 @@ private static void convertSchema(SchemaInfo schemaInfo, Schema schema) { } public static ByteBuf newProducer(String topic, long producerId, long requestId, String producerName, - boolean encrypted, Map metadata, SchemaInfo schemaInfo, - long epoch, boolean userProvidedProducerName, - ProducerAccessMode accessMode, Optional topicEpoch, boolean isTxnEnabled) { + boolean encrypted, Map metadata, SchemaInfo schemaInfo, + long epoch, boolean userProvidedProducerName, + ProducerAccessMode accessMode, Optional topicEpoch, boolean isTxnEnabled) { + return newProducer(topic, producerId, requestId, producerName, encrypted, metadata, schemaInfo, epoch, + userProvidedProducerName, accessMode, topicEpoch, isTxnEnabled, null); + + } + + public static ByteBuf newProducer(String topic, long producerId, long requestId, String producerName, + boolean encrypted, Map metadata, SchemaInfo schemaInfo, + long epoch, boolean userProvidedProducerName, + ProducerAccessMode accessMode, Optional topicEpoch, boolean isTxnEnabled, + String initialSubscriptionName) { BaseCommand cmd = localCmd(Type.PRODUCER); CommandProducer producer = cmd.setProducer() .setTopic(topic) @@ -792,6 +803,11 @@ public static ByteBuf newProducer(String topic, long producerId, long requestId, } topicEpoch.ifPresent(producer::setTopicEpoch); + + if (!Strings.isNullOrEmpty(initialSubscriptionName)) { + producer.setInitialSubscriptionName(initialSubscriptionName); + } + return serializeWithSize(cmd); } diff --git a/pulsar-common/src/main/proto/PulsarApi.proto b/pulsar-common/src/main/proto/PulsarApi.proto index 12367a2236082..b0e343cb41267 100644 --- a/pulsar-common/src/main/proto/PulsarApi.proto +++ b/pulsar-common/src/main/proto/PulsarApi.proto @@ -491,6 +491,10 @@ message CommandProducer { optional uint64 topic_epoch = 11; optional bool txn_enabled = 12 [default = false]; + + // Name of the initial subscription of the topic. + // If this field is not set, the initial subscription will not be created. + optional string initial_subscription_name = 13; } message CommandSend {