diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java index ea50e99971..7e2d0053fb 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java @@ -202,11 +202,13 @@ public static class OffsetAndTopicListener implements NamespaceBundleOwnershipLi final NamespaceName kafkaTopicNs; final GroupCoordinator groupCoordinator; final String brokerUrl; + final String tenant; public OffsetAndTopicListener(BrokerService service, String tenant, KafkaServiceConfiguration kafkaConfig, GroupCoordinator groupCoordinator) { + this.tenant = tenant; this.service = service; this.kafkaMetaNs = NamespaceName .get(tenant, kafkaConfig.getKafkaMetadataNamespace()); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java index fedfdf95cd..486785a866 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java @@ -1619,7 +1619,6 @@ protected void handleJoinGroupRequest(KafkaHeaderAndRequest joinGroup, log.trace("Sending join group response {} for correlation id {} to client {}.", response, joinGroup.getHeader().correlationId(), joinGroup.getHeader().clientId()); } - resultFuture.complete(response); }); } @@ -2002,7 +2001,9 @@ protected void handleAddOffsetsToTxn(KafkaHeaderAndRequest kafkaHeaderAndRequest CompletableFuture response) { AddOffsetsToTxnRequest request = (AddOffsetsToTxnRequest) kafkaHeaderAndRequest.getRequest(); int partition = getGroupCoordinator().partitionFor(request.consumerGroupId()); - String offsetTopicName = getGroupCoordinator().getGroupManager().getOffsetConfig().offsetsTopicName(); + String currentTenant = getCurrentTenant(); + String offsetTopicName = getGroupCoordinator().getGroupManager() + .getOffsetConfig().offsetsTopicName(); TransactionCoordinator transactionCoordinator = getTransactionCoordinator(); Set topicPartitions = Collections.singleton(new TopicPartition(offsetTopicName, partition)); transactionCoordinator.handleAddPartitionsToTransaction( diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupCoordinator.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupCoordinator.java index 7569549d42..b282d5730d 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupCoordinator.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupCoordinator.java @@ -71,6 +71,34 @@ @Slf4j public class GroupCoordinator { + private static final String NoState = ""; + private static final String NoProtocolType = ""; + static final String NoProtocol = ""; + private static final String NoLeader = ""; + private static final int NoGeneration = -1; + private static final String NoMemberId = ""; + private static final GroupSummary EmptyGroup = new GroupSummary( + NoState, + NoProtocolType, + NoProtocol, + Collections.emptyList() + ); + static final GroupSummary DeadGroup = new GroupSummary( + Dead.toString(), + NoProtocolType, + NoProtocol, + Collections.emptyList() + ); + + private final String tenant; + private final AtomicBoolean isActive = new AtomicBoolean(false); + private final GroupConfig groupConfig; + private final GroupMetadataManager groupManager; + private final DelayedOperationPurgatory heartbeatPurgatory; + private final DelayedOperationPurgatory joinPurgatory; + private final Time time; + + public static GroupCoordinator of( String tenant, SystemTopicClient client, @@ -86,6 +114,7 @@ public static GroupCoordinator of( .build(); GroupMetadataManager metadataManager = new GroupMetadataManager( + tenant, offsetConfig, client.newProducerBuilder(), client.newReaderBuilder(), @@ -95,17 +124,18 @@ public static GroupCoordinator of( ); DelayedOperationPurgatory joinPurgatory = DelayedOperationPurgatory.builder() - .purgatoryName("group-coordinator-delayed-join") - .timeoutTimer(timer) - .build(); - - DelayedOperationPurgatory heartbeatPurgatory = - DelayedOperationPurgatory.builder() - .purgatoryName("group-coordinator-delayed-heartbeat") + .purgatoryName("group-coordinator-delayed-join") .timeoutTimer(timer) .build(); + DelayedOperationPurgatory heartbeatPurgatory = + DelayedOperationPurgatory.builder() + .purgatoryName("group-coordinator-delayed-heartbeat") + .timeoutTimer(timer) + .build(); + return new GroupCoordinator( + tenant, groupConfig, metadataManager, heartbeatPurgatory, @@ -114,39 +144,15 @@ public static GroupCoordinator of( ); } - static final String NoState = ""; - static final String NoProtocolType = ""; - static final String NoProtocol = ""; - static final String NoLeader = ""; - static final int NoGeneration = -1; - static final String NoMemberId = ""; - static final GroupSummary EmptyGroup = new GroupSummary( - NoState, - NoProtocolType, - NoProtocol, - Collections.emptyList() - ); - static final GroupSummary DeadGroup = new GroupSummary( - Dead.toString(), - NoProtocolType, - NoProtocol, - Collections.emptyList() - ); - - private final AtomicBoolean isActive = new AtomicBoolean(false); - private final GroupConfig groupConfig; - private final GroupMetadataManager groupManager; - private final DelayedOperationPurgatory heartbeatPurgatory; - private final DelayedOperationPurgatory joinPurgatory; - private final Time time; - public GroupCoordinator( + String tenant, GroupConfig groupConfig, GroupMetadataManager groupManager, DelayedOperationPurgatory heartbeatPurgatory, DelayedOperationPurgatory joinPurgatory, Time time ) { + this.tenant = tenant; this.groupConfig = groupConfig; this.groupManager = groupManager; this.heartbeatPurgatory = heartbeatPurgatory; @@ -158,7 +164,7 @@ public GroupCoordinator( * Startup logic executed at the same time when the server starts up. */ public void startup(boolean enableMetadataExpiration) { - log.info("Starting up group coordinator."); + log.info("Starting up group coordinator for tenant {}...", tenant); groupManager.startup(enableMetadataExpiration); isActive.set(true); log.info("Group coordinator started."); @@ -169,7 +175,7 @@ public void startup(boolean enableMetadataExpiration) { * Ordering of actions should be reversed from the startup process. */ public void shutdown() { - log.info("Shutting down group coordinator ..."); + log.info("Shutting down group coordinator for tenant {}...", tenant); isActive.set(false); groupManager.shutdown(); heartbeatPurgatory.shutdown(); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupMetadataManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupMetadataManager.java index 542eeac104..0411fe7743 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupMetadataManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupMetadataManager.java @@ -93,11 +93,13 @@ @Slf4j public class GroupMetadataManager { - private final byte magicValue = RecordBatch.CURRENT_MAGIC_VALUE; + private static final byte MAGIC_VALUE = RecordBatch.CURRENT_MAGIC_VALUE; private final CompressionType compressionType; @Getter private final OffsetConfig offsetConfig; + private final String tenant; private final String namespacePrefix; + private final ConcurrentMap groupMetadataCache; /* lock protecting access to loading and owned partition sets */ private final ReentrantLock partitionLock = new ReentrantLock(); @@ -202,13 +204,16 @@ public String toString() { } - public GroupMetadataManager(OffsetConfig offsetConfig, + + public GroupMetadataManager(String tenant, + OffsetConfig offsetConfig, ProducerBuilder metadataTopicProducerBuilder, ReaderBuilder metadataTopicReaderBuilder, ScheduledExecutorService scheduler, String namespacePrefixForMetadata, Time time) { - this(offsetConfig, + this(tenant, + offsetConfig, metadataTopicProducerBuilder, metadataTopicReaderBuilder, scheduler, @@ -224,13 +229,15 @@ public static int getPartitionId(String groupId, int offsetsTopicNumPartitions) return MathUtils.signSafeMod(groupId.hashCode(), offsetsTopicNumPartitions); } - GroupMetadataManager(OffsetConfig offsetConfig, + GroupMetadataManager(String tenant, + OffsetConfig offsetConfig, ProducerBuilder metadataTopicProducerBuilder, ReaderBuilder metadataTopicConsumerBuilder, ScheduledExecutorService scheduler, Time time, Function partitioner, String namespacePrefix) { + this.tenant = tenant; this.namespacePrefix = namespacePrefix; this.offsetConfig = offsetConfig; this.compressionType = offsetConfig.offsetsTopicCompressionType(); @@ -398,13 +405,13 @@ public CompletableFuture storeGroup(GroupMetadata group, // construct the record ByteBuffer buffer = ByteBuffer.allocate(AbstractRecords.estimateSizeInBytes( - magicValue, + MAGIC_VALUE, compressionType, Lists.newArrayList(new SimpleRecord(timestamp, key, value)) )); MemoryRecordsBuilder recordsBuilder = MemoryRecords.builder( buffer, - magicValue, + MAGIC_VALUE, compressionType, timestampType, 0L @@ -509,12 +516,12 @@ public CompletableFuture> storeOffsets( ByteBuffer buffer = ByteBuffer.allocate( AbstractRecords.estimateSizeInBytes( - magicValue, compressionType, records + MAGIC_VALUE, compressionType, records ) ); MemoryRecordsBuilder builder = MemoryRecords.builder( - buffer, magicValue, compressionType, + buffer, MAGIC_VALUE, compressionType, timestampType, 0L, timestamp, producerId, producerEpoch, @@ -712,7 +719,8 @@ private boolean validateOffsetMetadataLength(String metadata) { public CompletableFuture scheduleLoadGroupAndOffsets(int offsetsPartition, Consumer onGroupLoaded) { - String topicPartition = offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + offsetsPartition; + String topicPartition = offsetConfig.offsetsTopicName() + + PARTITIONED_TOPIC_SUFFIX + offsetsPartition; if (addLoadingPartition(offsetsPartition)) { log.info("Scheduling loading of offsets and group metadata from {}", topicPartition); long startMs = time.milliseconds(); @@ -1248,7 +1256,7 @@ CompletableFuture cleanGroupMetadata(Stream groups, if (!tombstones.isEmpty()) { MemoryRecords records = MemoryRecords.withRecords( - magicValue, 0L, compressionType, + MAGIC_VALUE, 0L, compressionType, timestampType, tombstones.toArray(new SimpleRecord[tombstones.size()]) ); @@ -1265,7 +1273,8 @@ CompletableFuture cleanGroupMetadata(Stream groups, log.error("Failed to append {} tombstones to topic {} for expired/deleted " + "offsets and/or metadata for group {}", tombstones.size(), - offsetConfig.offsetsTopicName() + '-' + partitioner.apply(group.groupId()), + offsetConfig.offsetsTopicName() + + '-' + partitioner.apply(group.groupId()), group.groupId(), cause); // ignore and continue return 0; @@ -1366,7 +1375,7 @@ CompletableFuture> getOffsetsTopicProducer(String groupId) partitionId -> { if (log.isDebugEnabled()) { log.debug("Created Partitioned producer: {} for consumer group: {}", - offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + partitionId, + offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + partitionId, groupId); } return metadataTopicProducerBuilder.clone() @@ -1380,7 +1389,7 @@ CompletableFuture> getOffsetsTopicProducer(int partitionId) id -> { if (log.isDebugEnabled()) { log.debug("Will create Partitioned producer: {}", - offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + id); + offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + id); } return metadataTopicProducerBuilder.clone() .topic(offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + id) @@ -1393,7 +1402,7 @@ CompletableFuture> getOffsetsTopicReader(int partitionId) { id -> { if (log.isDebugEnabled()) { log.debug("Will create Partitioned reader: {}", - offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + id); + offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + id); } return metadataTopicReaderBuilder.clone() .topic(offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + partitionId) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/OffsetConfig.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/OffsetConfig.java index f76cfdcb54..3a8ef14ba5 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/OffsetConfig.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/OffsetConfig.java @@ -17,6 +17,7 @@ import lombok.Builder; import lombok.Builder.Default; import lombok.Data; +import lombok.NonNull; import lombok.experimental.Accessors; import org.apache.kafka.common.record.CompressionType; @@ -31,11 +32,11 @@ public class OffsetConfig { public static final int DefaultMaxMetadataSize = 4096; public static final long DefaultOffsetsRetentionMs = 24 * 60 * 60 * 1000L; public static final long DefaultOffsetsRetentionCheckIntervalMs = 600000L; - public static final String DefaultOffsetsTopicName = "public/default/__consumer_offsets"; public static final int DefaultOffsetsNumPartitions = KafkaServiceConfiguration.DefaultOffsetsTopicNumPartitions; @Default - private String offsetsTopicName = DefaultOffsetsTopicName; + @NonNull + private String offsetsTopicName = null; @Default private int maxMetadataSize = DefaultMaxMetadataSize; @Default diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/DifferentNamespaceTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/DifferentNamespaceTest.java index f138002416..ced92fe501 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/DifferentNamespaceTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/DifferentNamespaceTest.java @@ -41,6 +41,7 @@ import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.internals.Topic; +import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.policies.data.TenantInfo; import org.testng.annotations.AfterClass; import org.testng.annotations.BeforeClass; @@ -94,11 +95,21 @@ protected void setup() throws Exception { .allowedClusters(Collections.singleton(configClusterName)) .build()); admin.namespaces().createNamespace(NOT_ALLOWED_TENANT + "/" + NOT_ALLOWED_NAMESPACE); + + // ensure that the 'public' tenant does not exist + try { + admin.tenants().deleteTenant("public", true); + } catch (PulsarAdminException.NotFoundException ok) { + } } @AfterClass @Override protected void cleanup() throws Exception { + + // ensure that nobody tried to create the "public" tenant + assertTrue(!admin.tenants().getTenants().contains("public")); + super.internalCleanup(); } diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupCoordinatorTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupCoordinatorTest.java index 56a9d6a78d..2471a43ae3 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupCoordinatorTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupCoordinatorTest.java @@ -133,6 +133,7 @@ protected void setUp() throws PulsarClientException { otherGroupId = "otherGroupId"; offsetConfig.offsetsTopicNumPartitions(4); groupMetadataManager = spy(new GroupMetadataManager( + conf.getKafkaMetadataTenant(), offsetConfig, producerBuilder, readerBuilder, @@ -164,11 +165,12 @@ protected void setUp() throws PulsarClientException { .build(); groupCoordinator = new GroupCoordinator( - groupConfig, - groupMetadataManager, - heartbeatPurgatory, - joinPurgatory, - timer.time() + tenant, + groupConfig, + groupMetadataManager, + heartbeatPurgatory, + joinPurgatory, + timer.time() ); // start the group coordinator diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupMetadataManagerTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupMetadataManagerTest.java index 979b38aa8e..dacf8b6297 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupMetadataManagerTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupMetadataManagerTest.java @@ -153,6 +153,7 @@ protected void setUp() throws PulsarClientException, PulsarAdminException { .numThreads(1) .build(); groupMetadataManager = new GroupMetadataManager( + tenant, offsetConfig, producerBuilder, readerBuilder, @@ -882,6 +883,7 @@ public void testGroupLoadWithConsumerAndTransactionalOffsetCommitsTransactionWin @Test public void testGroupNotExits() { groupMetadataManager = new GroupMetadataManager( + conf.getKafkaMetadataTenant(), offsetConfig, producerBuilder, readerBuilder, @@ -1229,6 +1231,7 @@ public void testLoadGroupAndOffsetsFromDifferentSegments() throws Exception { @Test public void testAddGroup() { groupMetadataManager = new GroupMetadataManager( + conf.getKafkaMetadataTenant(), offsetConfig, producerBuilder, readerBuilder,