diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java index 47da95b34ac3f..1ee3f513ef288 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java @@ -58,6 +58,7 @@ import org.apache.pulsar.common.policies.data.ClusterData.ClusterUrl; import org.apache.pulsar.common.policies.data.stats.ConsumerStatsImpl; import org.apache.pulsar.common.protocol.Commands; +import org.apache.pulsar.common.schema.SchemaType; import org.apache.pulsar.common.stats.Rate; import org.apache.pulsar.common.util.DateFormatter; import org.apache.pulsar.common.util.FutureUtil; @@ -142,12 +143,24 @@ public class Consumer { private long negtiveUnackedMsgsTimestamp; + @Getter + private final SchemaType schemaType; + public Consumer(Subscription subscription, SubType subType, String topicName, long consumerId, int priorityLevel, String consumerName, boolean isDurable, TransportCnx cnx, String appId, Map metadata, boolean readCompacted, KeySharedMeta keySharedMeta, MessageId startMessageId, long consumerEpoch) { + this(subscription, subType, topicName, consumerId, priorityLevel, consumerName, isDurable, cnx, appId, + metadata, readCompacted, keySharedMeta, startMessageId, consumerEpoch, null); + } + public Consumer(Subscription subscription, SubType subType, String topicName, long consumerId, + int priorityLevel, String consumerName, + boolean isDurable, TransportCnx cnx, String appId, + Map metadata, boolean readCompacted, + KeySharedMeta keySharedMeta, MessageId startMessageId, + long consumerEpoch, SchemaType schemaType) { this.subscription = subscription; this.subType = subType; this.topicName = topicName; @@ -204,6 +217,8 @@ public Consumer(Subscription subscription, SubType subType, String topicName, lo this.consumerEpoch = consumerEpoch; this.isAcknowledgmentAtBatchIndexLevelEnabled = subscription.getTopic().getBrokerService() .getPulsar().getConfiguration().isAcknowledgmentAtBatchIndexLevelEnabled(); + + this.schemaType = schemaType; } @VisibleForTesting @@ -231,6 +246,7 @@ public Consumer(Subscription subscription, SubType subType, String topicName, lo this.clientAddress = null; this.startMessageId = null; this.isAcknowledgmentAtBatchIndexLevelEnabled = false; + this.schemaType = null; MESSAGE_PERMITS_UPDATER.set(this, availablePermits); } 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 5b9c7d39ff12f..80b5d24d17543 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 @@ -1195,8 +1195,9 @@ protected void handleSubscribe(final CommandSubscribe subscribe) { .replicatedSubscriptionStateArg(isReplicated).keySharedMeta(keySharedMeta) .subscriptionProperties(subscriptionProperties) .consumerEpoch(consumerEpoch) + .schemaType(schema == null ? null : schema.getType()) .build(); - if (schema != null) { + if (schema != null && schema.getType() != SchemaType.AUTO_CONSUME) { return topic.addSchemaIfIdleOrCheckCompatible(schema) .thenCompose(v -> topic.subscribe(option)); } else { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SubscriptionOption.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SubscriptionOption.java index d375c539e550e..af56d023616b4 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SubscriptionOption.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SubscriptionOption.java @@ -29,6 +29,7 @@ import org.apache.pulsar.common.api.proto.CommandSubscribe; import org.apache.pulsar.common.api.proto.KeySharedMeta; import org.apache.pulsar.common.api.proto.KeyValue; +import org.apache.pulsar.common.schema.SchemaType; @Getter @Builder @@ -49,6 +50,7 @@ public class SubscriptionOption { private KeySharedMeta keySharedMeta; private Optional> subscriptionProperties; private long consumerEpoch; + private SchemaType schemaType; public static Optional> getPropertiesMap(List list) { if (list == null) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java index 3b046570d732e..a0a8462a22753 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java @@ -87,6 +87,7 @@ import org.apache.pulsar.common.policies.data.stats.PublisherStatsImpl; import org.apache.pulsar.common.policies.data.stats.SubscriptionStatsImpl; import org.apache.pulsar.common.protocol.schema.SchemaData; +import org.apache.pulsar.common.schema.SchemaType; import org.apache.pulsar.common.util.FutureUtil; import org.apache.pulsar.common.util.collections.ConcurrentOpenHashMap; import org.apache.pulsar.metadata.api.MetadataStoreException; @@ -255,7 +256,8 @@ public CompletableFuture subscribe(SubscriptionOption option) { option.isDurable(), option.getStartMessageId(), option.getMetadata(), option.isReadCompacted(), option.getStartMessageRollbackDurationSec(), option.isReplicatedSubscriptionStateArg(), - option.getKeySharedMeta(), option.getSubscriptionProperties().orElse(null)); + option.getKeySharedMeta(), option.getSubscriptionProperties().orElse(null), + option.getSchemaType()); } @Override @@ -268,7 +270,7 @@ public CompletableFuture subscribe(final TransportCnx cnx, String subs KeySharedMeta keySharedMeta) { return internalSubscribe(cnx, subscriptionName, consumerId, subType, priorityLevel, consumerName, isDurable, startMessageId, metadata, readCompacted, resetStartMessageBackInSec, - replicateSubscriptionState, keySharedMeta, null); + replicateSubscriptionState, keySharedMeta, null, null); } private CompletableFuture internalSubscribe(final TransportCnx cnx, String subscriptionName, @@ -279,7 +281,8 @@ private CompletableFuture internalSubscribe(final TransportCnx cnx, St long resetStartMessageBackInSec, boolean replicateSubscriptionState, KeySharedMeta keySharedMeta, - Map subscriptionProperties) { + Map subscriptionProperties, + SchemaType schemaType) { return brokerService.checkTopicNsOwnership(getName()).thenCompose(__ -> { final CompletableFuture future = new CompletableFuture<>(); @@ -321,8 +324,8 @@ private CompletableFuture internalSubscribe(final TransportCnx cnx, St name -> new NonPersistentSubscription(this, subscriptionName, isDurable, subscriptionProperties)); Consumer consumer = new Consumer(subscription, subType, topic, consumerId, priorityLevel, consumerName, - false, cnx, cnx.getAuthRole(), metadata, readCompacted, keySharedMeta, - MessageId.latest, DEFAULT_CONSUMER_EPOCH); + false, cnx, cnx.getAuthRole(), metadata, readCompacted, keySharedMeta, MessageId.latest, + DEFAULT_CONSUMER_EPOCH, schemaType); if (isMigrated()) { consumer.topicMigrated(getClusterMigrationUrl()); } @@ -1162,12 +1165,14 @@ public CompletableFuture getLastMessageId() { @Override public CompletableFuture addSchemaIfIdleOrCheckCompatible(SchemaData schema) { return hasSchema().thenCompose((hasSchema) -> { - int numActiveConsumers = subscriptions.values().stream() - .mapToInt(subscription -> subscription.getConsumers().size()) + int numActiveConsumersWithoutAutoSchema = subscriptions.values().stream() + .mapToInt(subscription -> subscription.getConsumers().stream() + .filter(consumer -> consumer.getSchemaType() != SchemaType.AUTO_CONSUME) + .toList().size()) .sum(); if (hasSchema || (!producers.isEmpty()) - || (numActiveConsumers != 0) + || (numActiveConsumersWithoutAutoSchema != 0) || ENTRIES_ADDED_COUNTER_UPDATER.get(this) != 0) { return checkSchemaCompatibleForConsumer(schema); } else { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 4277edb074c5e..fd0b069421275 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -155,6 +155,7 @@ import org.apache.pulsar.common.policies.data.stats.TopicStatsImpl; import org.apache.pulsar.common.protocol.Commands; import org.apache.pulsar.common.protocol.schema.SchemaData; +import org.apache.pulsar.common.schema.SchemaType; import org.apache.pulsar.common.topics.TopicCompactionStrategy; import org.apache.pulsar.common.util.Codec; import org.apache.pulsar.common.util.DateFormatter; @@ -727,7 +728,8 @@ public CompletableFuture subscribe(SubscriptionOption option) { option.getStartMessageId(), option.getMetadata(), option.isReadCompacted(), option.getInitialPosition(), option.getStartMessageRollbackDurationSec(), option.isReplicatedSubscriptionStateArg(), option.getKeySharedMeta(), - option.getSubscriptionProperties().orElse(Collections.emptyMap()), option.getConsumerEpoch()); + option.getSubscriptionProperties().orElse(Collections.emptyMap()), + option.getConsumerEpoch(), option.getSchemaType()); } private CompletableFuture internalSubscribe(final TransportCnx cnx, String subscriptionName, @@ -740,7 +742,8 @@ private CompletableFuture internalSubscribe(final TransportCnx cnx, St boolean replicatedSubscriptionStateArg, KeySharedMeta keySharedMeta, Map subscriptionProperties, - long consumerEpoch) { + long consumerEpoch, + SchemaType schemaType) { if (readCompacted && !(subType == SubType.Failover || subType == SubType.Exclusive)) { return FutureUtil.failedFuture(new NotAllowedException( "readCompacted only allowed on failover or exclusive subscriptions")); @@ -828,7 +831,7 @@ private CompletableFuture internalSubscribe(final TransportCnx cnx, St CompletableFuture future = subscriptionFuture.thenCompose(subscription -> { Consumer consumer = new Consumer(subscription, subType, topic, consumerId, priorityLevel, consumerName, isDurable, cnx, cnx.getAuthRole(), metadata, - readCompacted, keySharedMeta, startMessageId, consumerEpoch); + readCompacted, keySharedMeta, startMessageId, consumerEpoch, schemaType); return addConsumerToSubscription(subscription, consumer).thenCompose(v -> { if (subscription instanceof PersistentSubscription persistentSubscription) { @@ -907,7 +910,7 @@ public CompletableFuture subscribe(final TransportCnx cnx, String subs KeySharedMeta keySharedMeta) { return internalSubscribe(cnx, subscriptionName, consumerId, subType, priorityLevel, consumerName, isDurable, startMessageId, metadata, readCompacted, initialPosition, startMessageRollbackDurationSec, - replicatedSubscriptionStateArg, keySharedMeta, null, DEFAULT_CONSUMER_EPOCH); + replicatedSubscriptionStateArg, keySharedMeta, null, DEFAULT_CONSUMER_EPOCH, null); } private CompletableFuture getDurableSubscription(String subscriptionName, @@ -3107,21 +3110,22 @@ public synchronized OffloadProcessStatus offloadStatus() { @Override public CompletableFuture addSchemaIfIdleOrCheckCompatible(SchemaData schema) { - return hasSchema() - .thenCompose((hasSchema) -> { - int numActiveConsumers = subscriptions.values().stream() - .mapToInt(subscription -> subscription.getConsumers().size()) - .sum(); - if (hasSchema - || (!producers.isEmpty()) - || (numActiveConsumers != 0) - || (ledger.getTotalSize() != 0)) { - return checkSchemaCompatibleForConsumer(schema); - } else { - return addSchema(schema).thenCompose(schemaVersion -> - CompletableFuture.completedFuture(null)); - } - }); + return hasSchema().thenCompose((hasSchema) -> { + int numActiveConsumersWithoutAutoSchema = subscriptions.values().stream() + .mapToInt(subscription -> subscription.getConsumers().stream() + .filter(consumer -> consumer.getSchemaType() != SchemaType.AUTO_CONSUME) + .toList().size()) + .sum(); + if (hasSchema + || (!producers.isEmpty()) + || (numActiveConsumersWithoutAutoSchema != 0) + || (ledger.getTotalSize() != 0)) { + return checkSchemaCompatibleForConsumer(schema); + } else { + return addSchema(schema).thenCompose(schemaVersion -> + CompletableFuture.completedFuture(null)); + } + }); } public synchronized void checkReplicatedSubscriptionControllerState() { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleSchemaTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleSchemaTest.java index 21d6a7ed89a56..c8c7c3b2ccc38 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleSchemaTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleSchemaTest.java @@ -1239,6 +1239,79 @@ public void testAutoCreatedSchema(String domain) throws Exception { Assert.assertEquals(admin.schemas().getSchemaInfo(topic2).getType(), SchemaType.STRING); } + @Test(dataProvider = "topicDomain") + public void testSubscribeWithSchemaAfterAutoConsumeNewTopic(String domain) throws Exception { + final String topic = domain + "my-property/my-ns/testSubscribeWithSchemaAfterAutoConsume-1"; + + @Cleanup + Consumer autoConsumer1 = pulsarClient.newConsumer(Schema.AUTO_CONSUME()) + .topic(topic) + .subscriptionType(SubscriptionType.Shared) + .subscriptionName("sub0") + .consumerName("autoConsumer1") + .subscribe(); + @Cleanup + Consumer autoConsumer2 = pulsarClient.newConsumer(Schema.AUTO_CONSUME()) + .topic(topic) + .subscriptionType(SubscriptionType.Shared) + .subscriptionName("sub0") + .consumerName("autoConsumer2") + .subscribe(); + @Cleanup + Consumer autoConsumer3 = pulsarClient.newConsumer(Schema.AUTO_CONSUME()) + .topic(topic) + .subscriptionType(SubscriptionType.Shared) + .subscriptionName("sub1") + .consumerName("autoConsumer3") + .subscribe(); + @Cleanup + Consumer autoConsumer4 = pulsarClient.newConsumer(Schema.AUTO_CONSUME()) + .topic(topic) + .subscriptionType(SubscriptionType.Shared) + .subscriptionName("sub1") + .consumerName("autoConsumer4") + .subscribe(); + try { + log.info("The autoConsumer1 isConnected: " + autoConsumer1.isConnected()); + log.info("The autoConsumer2 isConnected: " + autoConsumer2.isConnected()); + log.info("The autoConsumer3 isConnected: " + autoConsumer3.isConnected()); + log.info("The autoConsumer4 isConnected: " + autoConsumer4.isConnected()); + admin.schemas().getSchemaInfo(topic); + fail("The schema of topic should not exist"); + } catch (PulsarAdminException e) { + assertEquals(e.getStatusCode(), 404); + } + + @Cleanup + Consumer consumerWithSchema1 = pulsarClient.newConsumer(Schema.AVRO(V1Data.class)) + .topic(topic) + .subscriptionType(SubscriptionType.Shared) + .subscriptionName("sub0") + .consumerName("consumerWithSchema-1") + .subscribe(); + @Cleanup + Consumer consumerWithSchema2 = pulsarClient.newConsumer(Schema.AVRO(V1Data.class)) + .topic(topic) + .subscriptionType(SubscriptionType.Shared) + .subscriptionName("sub0") + .consumerName("consumerWithSchema-2") + .subscribe(); + @Cleanup + Consumer consumerWithSchema3 = pulsarClient.newConsumer(Schema.AVRO(V1Data.class)) + .topic(topic) + .subscriptionType(SubscriptionType.Shared) + .subscriptionName("sub1") + .consumerName("consumerWithSchema-3") + .subscribe(); + @Cleanup + Consumer consumerWithSchema4 = pulsarClient.newConsumer(Schema.AVRO(V1Data.class)) + .topic(topic) + .subscriptionType(SubscriptionType.Shared) + .subscriptionName("sub1") + .consumerName("consumerWithSchema-4") + .subscribe(); + } + @DataProvider(name = "keyEncodingType") public static Object[] keyEncodingType() { return new Object[] { KeyValueEncodingType.SEPARATED, KeyValueEncodingType.INLINE }; 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 08a6bb15807c9..196b0600c0bf6 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 @@ -84,6 +84,7 @@ import org.apache.pulsar.client.api.transaction.TxnID; import org.apache.pulsar.client.impl.conf.ConsumerConfigurationData; import org.apache.pulsar.client.impl.crypto.MessageCryptoBc; +import org.apache.pulsar.client.impl.schema.AutoConsumeSchema; import org.apache.pulsar.client.impl.transaction.TransactionImpl; import org.apache.pulsar.client.util.ExecutorProvider; import org.apache.pulsar.client.util.RetryMessageUtil; @@ -811,6 +812,11 @@ public void connectionOpened(final ClientCnx cnx) { if (si != null && (SchemaType.BYTES == si.getType() || SchemaType.NONE == si.getType())) { // don't set schema for Schema.BYTES si = null; + } else { + if (schema instanceof AutoConsumeSchema + && Commands.peerSupportsCarryAutoConsumeSchemaToBroker(cnx.getRemoteEndpointProtocolVersion())) { + si = AutoConsumeSchema.SCHEMA_INFO; + } } // startMessageRollbackDurationInSec should be consider only once when consumer connects to first time long startMessageRollbackDuration = (startMessageRollbackDurationInSec > 0 diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AutoConsumeSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AutoConsumeSchema.java index 33fcd18876be6..82a3b69da20b6 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AutoConsumeSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AutoConsumeSchema.java @@ -57,6 +57,12 @@ public class AutoConsumeSchema implements Schema { private SchemaInfoProvider schemaInfoProvider; + public static final SchemaInfo SCHEMA_INFO = SchemaInfoImpl.builder() + .name("AutoConsume") + .type(SchemaType.AUTO_CONSUME) + .schema(new byte[0]) + .build(); + private ConcurrentMap> initSchemaMap() { ConcurrentMap> schemaMap = new ConcurrentHashMap<>(); // The Schema.BYTES will not be uploaded to the broker and store in the schema storage, 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 081dfe4275b24..8a5684cf676b0 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 @@ -771,7 +771,9 @@ public static ByteBuf newProducer(String topic, long producerId, long requestId, } private static Schema.Type getSchemaType(SchemaType type) { - if (type.getValue() < 0) { + if (type == SchemaType.AUTO_CONSUME) { + return Schema.Type.AutoConsume; + } else if (type.getValue() < 0) { return Schema.Type.None; } else { return Schema.Type.valueOf(type.getValue()); @@ -779,7 +781,9 @@ private static Schema.Type getSchemaType(SchemaType type) { } public static SchemaType getSchemaType(Schema.Type type) { - if (type.getValue() < 0) { + if (type == Schema.Type.AutoConsume) { + return SchemaType.AUTO_CONSUME; + } else if (type.getValue() < 0) { // this is unexpected return SchemaType.NONE; } else { @@ -1965,6 +1969,10 @@ public static boolean peerSupportsAckReceipt(int peerVersion) { return peerVersion >= ProtocolVersion.v17.getValue(); } + public static boolean peerSupportsCarryAutoConsumeSchemaToBroker(int peerVersion) { + return peerVersion >= ProtocolVersion.v21.getValue(); + } + private static org.apache.pulsar.common.api.proto.ProducerAccessMode convertProducerAccessMode( ProducerAccessMode accessMode) { switch (accessMode) { diff --git a/pulsar-common/src/main/proto/PulsarApi.proto b/pulsar-common/src/main/proto/PulsarApi.proto index acf75eab85826..d9c41eeec9740 100644 --- a/pulsar-common/src/main/proto/PulsarApi.proto +++ b/pulsar-common/src/main/proto/PulsarApi.proto @@ -45,6 +45,7 @@ message Schema { LocalTime = 18; LocalDateTime = 19; ProtobufNative = 20; + AutoConsume = 21; } required string name = 1; @@ -263,6 +264,7 @@ enum ProtocolVersion { v18 = 18; // Add client support for broker entry metadata v19 = 19; // Add CommandTcClientConnectRequest and CommandTcClientConnectResponse v20 = 20; // Add client support for topic migration redirection CommandTopicMigrated + v21 = 21; // Carry the AUTO_CONSUME schema to the Broker after this version } message CommandConnect {