From d25ce223b33fad57702f730d7b169406238ecac2 Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Mon, 25 Oct 2021 12:51:10 +0200 Subject: [PATCH 1/8] Multi tenant configuration - fix handling of group management topics tenant --- .../handlers/kop/KafkaProtocolHandler.java | 3 ++ .../handlers/kop/KafkaRequestHandler.java | 4 +-- .../coordinator/group/GroupCoordinator.java | 10 ++++-- .../group/GroupMetadataManager.java | 31 +++++++++++-------- .../kop/coordinator/group/OffsetConfig.java | 7 ++++- .../handlers/kop/DifferentNamespaceTest.java | 11 +++++++ .../group/GroupCoordinatorTest.java | 3 ++ .../group/GroupMetadataManagerTest.java | 2 ++ 8 files changed, 53 insertions(+), 18 deletions(-) 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 41b170f13e..60a6776f15 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 @@ -176,11 +176,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()); @@ -656,6 +658,7 @@ protected GroupCoordinator startGroupCoordinator(String tenant, SystemTopicClien .build(); GroupCoordinator groupCoordinator = GroupCoordinator.of( + tenant, client, groupConfig, offsetConfig, 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 c08a1f0447..5859c940e5 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 @@ -1811,7 +1811,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); }); } @@ -2200,7 +2199,8 @@ 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().getCurrentOffsetsTopicName(currentTenant); TransactionCoordinator transactionCoordinator = getTransactionCoordinator(); transactionCoordinator.handleAddPartitionsToTransaction( request.transactionalId(), 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 dd1922da5d..7b43fac151 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 @@ -73,6 +73,7 @@ public class GroupCoordinator { public static GroupCoordinator of( + String tenant, SystemTopicClient client, GroupConfig groupConfig, OffsetConfig offsetConfig, @@ -85,6 +86,7 @@ public static GroupCoordinator of( .build(); GroupMetadataManager metadataManager = new GroupMetadataManager( + tenant, offsetConfig, client.newProducerBuilder(), client.newReaderBuilder(), @@ -104,6 +106,7 @@ public static GroupCoordinator of( .build(); return new GroupCoordinator( + tenant, groupConfig, metadataManager, heartbeatPurgatory, @@ -131,6 +134,7 @@ public static GroupCoordinator of( Collections.emptyList() ); + private String tenant; private final AtomicBoolean isActive = new AtomicBoolean(false); private final GroupConfig groupConfig; private final GroupMetadataManager groupManager; @@ -139,12 +143,14 @@ public static GroupCoordinator of( 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; @@ -156,7 +162,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."); @@ -167,7 +173,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 0280818a9e..121e23f058 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 @@ -166,6 +166,7 @@ public String toString() { private final CompressionType compressionType; @Getter private final OffsetConfig offsetConfig; + private final String tenant; private final ConcurrentMap groupMetadataCache; /* lock protecting access to loading and owned partition sets */ private final ReentrantLock partitionLock = new ReentrantLock(); @@ -201,12 +202,13 @@ public String toString() { private final Time time; private final Function partitioner; - public GroupMetadataManager(OffsetConfig offsetConfig, + public GroupMetadataManager(String tenant, + OffsetConfig offsetConfig, ProducerBuilder metadataTopicProducerBuilder, ReaderBuilder metadataTopicReaderBuilder, ScheduledExecutorService scheduler, Time time) { - this( + this(tenant, offsetConfig, metadataTopicProducerBuilder, metadataTopicReaderBuilder, @@ -222,12 +224,14 @@ 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) { + this.tenant = tenant; this.offsetConfig = offsetConfig; this.compressionType = offsetConfig.offsetsTopicCompressionType(); this.groupMetadataCache = new ConcurrentHashMap<>(); @@ -318,11 +322,11 @@ public int partitionFor(String groupId) { } public String getTopicPartitionName() { - return offsetConfig.offsetsTopicName(); + return offsetConfig.getCurrentOffsetsTopicName(tenant); } public String getTopicPartitionName(int partitionId) { - return getTopicPartitionName(offsetConfig.offsetsTopicName(), partitionId); + return getTopicPartitionName(offsetConfig.getCurrentOffsetsTopicName(tenant), partitionId); } public static String getTopicPartitionName(String offsetsTopicName, int partitionId) { @@ -721,7 +725,7 @@ private boolean validateOffsetMetadataLength(String metadata) { public CompletableFuture scheduleLoadGroupAndOffsets(int offsetsPartition, Consumer onGroupLoaded) { - String topicPartition = offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + offsetsPartition; + String topicPartition = offsetConfig.getCurrentOffsetsTopicName(tenant) + PARTITIONED_TOPIC_SUFFIX + offsetsPartition; if (addLoadingPartition(offsetsPartition)) { log.info("Scheduling loading of offsets and group metadata from {}", topicPartition); long startMs = time.milliseconds(); @@ -1154,6 +1158,7 @@ public void removeGroupsForPartition(int offsetsPartition, TopicPartition topicPartition = new TopicPartition( GROUP_METADATA_TOPIC_NAME, offsetsPartition ); + log.info("removeGroupsForPartition {}", topicPartition); log.info("Scheduling unloading of offsets and group metadata from {}", topicPartition); scheduler.submit(() -> { AtomicInteger numOffsetsRemoved = new AtomicInteger(); @@ -1274,7 +1279,7 @@ 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.getCurrentOffsetsTopicName(tenant) + '-' + partitioner.apply(group.groupId()), group.groupId(), cause); // ignore and continue return 0; @@ -1375,11 +1380,11 @@ CompletableFuture> getOffsetsTopicProducer(String groupId) partitionId -> { if (log.isDebugEnabled()) { log.debug("Created Partitioned producer: {} for consumer group: {}", - offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + partitionId, + offsetConfig.getCurrentOffsetsTopicName(tenant) + PARTITIONED_TOPIC_SUFFIX + partitionId, groupId); } return metadataTopicProducerBuilder.clone() - .topic(offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + partitionId) + .topic(offsetConfig.getCurrentOffsetsTopicName(tenant) + PARTITIONED_TOPIC_SUFFIX + partitionId) .createAsync(); }); } @@ -1389,10 +1394,10 @@ CompletableFuture> getOffsetsTopicProducer(int partitionId) id -> { if (log.isDebugEnabled()) { log.debug("Will create Partitioned producer: {}", - offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + id); + offsetConfig.getCurrentOffsetsTopicName(tenant) + PARTITIONED_TOPIC_SUFFIX + id); } return metadataTopicProducerBuilder.clone() - .topic(offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + id) + .topic(offsetConfig.getCurrentOffsetsTopicName(tenant) + PARTITIONED_TOPIC_SUFFIX + id) .createAsync(); }); } @@ -1402,10 +1407,10 @@ CompletableFuture> getOffsetsTopicReader(int partitionId) { id -> { if (log.isDebugEnabled()) { log.debug("Will create Partitioned reader: {}", - offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + id); + offsetConfig.getCurrentOffsetsTopicName(tenant) + PARTITIONED_TOPIC_SUFFIX + id); } return metadataTopicReaderBuilder.clone() - .topic(offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + partitionId) + .topic(offsetConfig.getCurrentOffsetsTopicName(tenant) + PARTITIONED_TOPIC_SUFFIX + partitionId) .readCompacted(true) .createAsync(); }); 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..4262b04d96 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 @@ -31,11 +31,16 @@ 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 String DefaultOffsetsTopicName = "${tenant}/default/__consumer_offsets"; public static final int DefaultOffsetsNumPartitions = KafkaServiceConfiguration.DefaultOffsetsTopicNumPartitions; @Default private String offsetsTopicName = DefaultOffsetsTopicName; + + public String getCurrentOffsetsTopicName(String tenant) { + return offsetsTopicName.replace(KafkaServiceConfiguration.TENANT_PLACEHOLDER, tenant); + } + @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 3d61a89fda..62661d9b09 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 @@ -22,6 +22,7 @@ import com.google.common.collect.ImmutableMap; import com.google.common.collect.Lists; import com.google.common.collect.Sets; +import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; import io.streamnative.pulsar.handlers.kop.KopProtocolHandlerTestBase; import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupMetadata.GroupOverview; import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupMetadata.GroupSummary; @@ -134,6 +135,7 @@ public void setup() throws Exception { otherGroupId = "otherGroupId"; offsetConfig.offsetsTopicNumPartitions(4); groupMetadataManager = spy(new GroupMetadataManager( + conf.getKafkaMetadataTenant(), offsetConfig, producerBuilder, readerBuilder, @@ -164,6 +166,7 @@ public void setup() throws Exception { .build(); groupCoordinator = new GroupCoordinator( + tenant, groupConfig, groupMetadataManager, heartbeatPurgatory, 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 00f68c5e34..38b4617fab 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 @@ -823,6 +823,7 @@ public void testGroupLoadWithConsumerAndTransactionalOffsetCommitsTransactionWin @Test public void testGroupNotExits() { groupMetadataManager = new GroupMetadataManager( + conf.getKafkaMetadataTenant(), offsetConfig, producer, consumer, @@ -1165,6 +1166,7 @@ public void testLoadGroupAndOffsetsFromDifferentSegments() throws Exception { @Test public void testAddGroup() { groupMetadataManager = new GroupMetadataManager( + conf.getKafkaMetadataTenant(), offsetConfig, producer, consumer, From 83903d0c9aef5234e80c7a13c39d69f3762e8542 Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Mon, 25 Oct 2021 13:09:43 +0200 Subject: [PATCH 2/8] fix CI --- .../pulsar/handlers/kop/KafkaRequestHandler.java | 3 ++- .../kop/coordinator/group/GroupMetadataManager.java | 6 ++++-- .../kop/coordinator/group/GroupCoordinatorTest.java | 1 - 3 files changed, 6 insertions(+), 4 deletions(-) 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 5859c940e5..b489da7f54 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 @@ -2200,7 +2200,8 @@ protected void handleAddOffsetsToTxn(KafkaHeaderAndRequest kafkaHeaderAndRequest AddOffsetsToTxnRequest request = (AddOffsetsToTxnRequest) kafkaHeaderAndRequest.getRequest(); int partition = getGroupCoordinator().partitionFor(request.consumerGroupId()); String currentTenant = getCurrentTenant(); - String offsetTopicName = getGroupCoordinator().getGroupManager().getOffsetConfig().getCurrentOffsetsTopicName(currentTenant); + String offsetTopicName = getGroupCoordinator().getGroupManager() + .getOffsetConfig().getCurrentOffsetsTopicName(currentTenant); TransactionCoordinator transactionCoordinator = getTransactionCoordinator(); transactionCoordinator.handleAddPartitionsToTransaction( request.transactionalId(), 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 121e23f058..105ae6f397 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 @@ -725,7 +725,8 @@ private boolean validateOffsetMetadataLength(String metadata) { public CompletableFuture scheduleLoadGroupAndOffsets(int offsetsPartition, Consumer onGroupLoaded) { - String topicPartition = offsetConfig.getCurrentOffsetsTopicName(tenant) + PARTITIONED_TOPIC_SUFFIX + offsetsPartition; + String topicPartition = offsetConfig.getCurrentOffsetsTopicName(tenant) + + PARTITIONED_TOPIC_SUFFIX + offsetsPartition; if (addLoadingPartition(offsetsPartition)) { log.info("Scheduling loading of offsets and group metadata from {}", topicPartition); long startMs = time.milliseconds(); @@ -1279,7 +1280,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.getCurrentOffsetsTopicName(tenant) + '-' + partitioner.apply(group.groupId()), + offsetConfig.getCurrentOffsetsTopicName(tenant) + + '-' + partitioner.apply(group.groupId()), group.groupId(), cause); // ignore and continue return 0; 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 62661d9b09..bf6fcba2af 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 @@ -22,7 +22,6 @@ import com.google.common.collect.ImmutableMap; import com.google.common.collect.Lists; import com.google.common.collect.Sets; -import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; import io.streamnative.pulsar.handlers.kop.KopProtocolHandlerTestBase; import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupMetadata.GroupOverview; import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupMetadata.GroupSummary; From 35ecc015d26a9f125b5ecf4beca50eb770f5ffa5 Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Mon, 25 Oct 2021 14:13:38 +0200 Subject: [PATCH 3/8] Fix CI --- .../pulsar/handlers/kop/coordinator/group/GroupCoordinator.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 7b43fac151..d662fb6306 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 @@ -134,7 +134,7 @@ public static GroupCoordinator of( Collections.emptyList() ); - private String tenant; + private final String tenant; private final AtomicBoolean isActive = new AtomicBoolean(false); private final GroupConfig groupConfig; private final GroupMetadataManager groupManager; From 69d83ec95c3ce1fed290da564ee4e7edb226ecac Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Mon, 25 Oct 2021 14:29:58 +0200 Subject: [PATCH 4/8] Fix style --- .../coordinator/group/GroupCoordinator.java | 93 ++++++++++--------- .../group/GroupMetadataManager.java | 89 +++++++++--------- 2 files changed, 91 insertions(+), 91 deletions(-) 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 d662fb6306..6a428601b4 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 @@ -72,13 +72,41 @@ @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, - GroupConfig groupConfig, - OffsetConfig offsetConfig, - Timer timer, - Time time + String tenant, + SystemTopicClient client, + GroupConfig groupConfig, + OffsetConfig offsetConfig, + Timer timer, + Time time ) { ScheduledExecutorService coordinatorExecutor = OrderedScheduler.newSchedulerBuilder() .name("group-coordinator-executor") @@ -86,25 +114,25 @@ public static GroupCoordinator of( .build(); GroupMetadataManager metadataManager = new GroupMetadataManager( - tenant, - offsetConfig, - client.newProducerBuilder(), - client.newReaderBuilder(), - coordinatorExecutor, - time + tenant, + offsetConfig, + client.newProducerBuilder(), + client.newReaderBuilder(), + coordinatorExecutor, + time ); 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, @@ -115,33 +143,6 @@ 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 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 GroupCoordinator( String tenant, GroupConfig groupConfig, 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 105ae6f397..53462d38ee 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,6 +93,45 @@ @Slf4j public class GroupMetadataManager { + private final static byte MAGIC_VALUE = RecordBatch.CURRENT_MAGIC_VALUE; + private final CompressionType compressionType; + @Getter + private final OffsetConfig offsetConfig; + private final String tenant; + private final ConcurrentMap groupMetadataCache; + /* lock protecting access to loading and owned partition sets */ + private final ReentrantLock partitionLock = new ReentrantLock(); + /** + * partitions of consumer groups that are being loaded, its lock should + * be always called BEFORE the group lock if needed. + */ + private final Set loadingPartitions = new HashSet<>(); + /* partitions of consumer groups that are assigned, using the same loading partition lock */ + private final Set ownedPartitions = new HashSet<>(); + /* shutting down flag */ + private final AtomicBoolean shuttingDown = new AtomicBoolean(false); + private final int groupMetadataTopicPartitionCount; + + // Map of + private final ConcurrentMap>> offsetsProducers = + new ConcurrentHashMap<>(); + private final ConcurrentMap>> offsetsReaders = + new ConcurrentHashMap<>(); + + /* single-thread scheduler to handle offset/group metadata cache loading and unloading */ + private final ScheduledExecutorService scheduler; + /** + * The groups with open transactional offsets commits per producer. We need this because when the commit or abort + * marker comes in for a transaction, it is for a particular partition on the offsets topic and a particular + * producerId. We use this structure to quickly find the groups which need to be updated by the commit/abort + * marker. + */ + private final Map> openGroupsForProducer = new HashMap<>(); + + private final ProducerBuilder metadataTopicProducerBuilder; + private final ReaderBuilder metadataTopicReaderBuilder; + private final Time time; + private final Function partitioner; /** * The key interface. */ @@ -162,46 +201,6 @@ public String toString() { } - private final byte magicValue = RecordBatch.CURRENT_MAGIC_VALUE; - private final CompressionType compressionType; - @Getter - private final OffsetConfig offsetConfig; - private final String tenant; - private final ConcurrentMap groupMetadataCache; - /* lock protecting access to loading and owned partition sets */ - private final ReentrantLock partitionLock = new ReentrantLock(); - /** - * partitions of consumer groups that are being loaded, its lock should - * be always called BEFORE the group lock if needed. - */ - private final Set loadingPartitions = new HashSet<>(); - /* partitions of consumer groups that are assigned, using the same loading partition lock */ - private final Set ownedPartitions = new HashSet<>(); - /* shutting down flag */ - private final AtomicBoolean shuttingDown = new AtomicBoolean(false); - private final int groupMetadataTopicPartitionCount; - - // Map of - private final ConcurrentMap>> offsetsProducers = - new ConcurrentHashMap<>(); - private final ConcurrentMap>> offsetsReaders = - new ConcurrentHashMap<>(); - - /* single-thread scheduler to handle offset/group metadata cache loading and unloading */ - private final ScheduledExecutorService scheduler; - /** - * The groups with open transactional offsets commits per producer. We need this because when the commit or abort - * marker comes in for a transaction, it is for a particular partition on the offsets topic and a particular - * producerId. We use this structure to quickly find the groups which need to be updated by the commit/abort - * marker. - */ - private final Map> openGroupsForProducer = new HashMap<>(); - - private final ProducerBuilder metadataTopicProducerBuilder; - private final ReaderBuilder metadataTopicReaderBuilder; - private final Time time; - private final Function partitioner; - public GroupMetadataManager(String tenant, OffsetConfig offsetConfig, ProducerBuilder metadataTopicProducerBuilder, @@ -398,13 +397,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 +508,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, @@ -1263,7 +1262,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()]) ); From 897f37dd6852f3b036db8171de63b612a53a166e Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Mon, 25 Oct 2021 14:37:15 +0200 Subject: [PATCH 5/8] fix CI --- .../handlers/kop/coordinator/group/GroupMetadataManager.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 53462d38ee..912f65e28f 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,7 +93,7 @@ @Slf4j public class GroupMetadataManager { - private final static byte MAGIC_VALUE = RecordBatch.CURRENT_MAGIC_VALUE; + private static final byte MAGIC_VALUE = RecordBatch.CURRENT_MAGIC_VALUE; private final CompressionType compressionType; @Getter private final OffsetConfig offsetConfig; From 4ede02f4ea36ed37f3b436a9c2bd99614ffc4843 Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Mon, 29 Nov 2021 18:14:33 +0100 Subject: [PATCH 6/8] Address some review comments --- .../handlers/kop/coordinator/group/GroupMetadataManager.java | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) 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 507aae15b6..a1028c27af 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 @@ -405,13 +405,13 @@ public CompletableFuture storeGroup(GroupMetadata group, // construct the record ByteBuffer buffer = ByteBuffer.allocate(AbstractRecords.estimateSizeInBytes( - MAGIC_VALUE, + MAGIC_VALUE, compressionType, Lists.newArrayList(new SimpleRecord(timestamp, key, value)) )); MemoryRecordsBuilder recordsBuilder = MemoryRecords.builder( buffer, - MAGIC_VALUE, + MAGIC_VALUE, compressionType, timestampType, 0L @@ -1153,7 +1153,6 @@ public void removeGroupsForPartition(int offsetsPartition, TopicPartition topicPartition = new TopicPartition( GROUP_METADATA_TOPIC_NAME, offsetsPartition ); - log.info("removeGroupsForPartition {}", topicPartition); log.info("Scheduling unloading of offsets and group metadata from {}", topicPartition); scheduler.submit(() -> { AtomicInteger numOffsetsRemoved = new AtomicInteger(); From 8fdbd8427e08c16e8713c5def28bdd9f3a1da273 Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Mon, 29 Nov 2021 18:17:37 +0100 Subject: [PATCH 7/8] Fix build after conflict resolution --- .../handlers/kop/coordinator/group/GroupMetadataManagerTest.java | 1 + 1 file changed, 1 insertion(+) 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 b92d66936a..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, From e0c2c04ab758092e16d5959bffbc5d935a7311c5 Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Wed, 1 Dec 2021 08:14:55 +0100 Subject: [PATCH 8/8] Remove useless method --- .../handlers/kop/KafkaRequestHandler.java | 2 +- .../group/GroupMetadataManager.java | 20 +++++++++---------- .../kop/coordinator/group/OffsetConfig.java | 10 +++------- 3 files changed, 14 insertions(+), 18 deletions(-) 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 8b36acfacd..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 @@ -2003,7 +2003,7 @@ protected void handleAddOffsetsToTxn(KafkaHeaderAndRequest kafkaHeaderAndRequest int partition = getGroupCoordinator().partitionFor(request.consumerGroupId()); String currentTenant = getCurrentTenant(); String offsetTopicName = getGroupCoordinator().getGroupManager() - .getOffsetConfig().getCurrentOffsetsTopicName(currentTenant); + .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/GroupMetadataManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupMetadataManager.java index 582660d5b6..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 @@ -329,11 +329,11 @@ public int partitionFor(String groupId) { } public String getTopicPartitionName() { - return offsetConfig.getCurrentOffsetsTopicName(tenant); + return offsetConfig.offsetsTopicName(); } public String getTopicPartitionName(int partitionId) { - return getTopicPartitionName(offsetConfig.getCurrentOffsetsTopicName(tenant), partitionId); + return getTopicPartitionName(offsetConfig.offsetsTopicName(), partitionId); } public static String getTopicPartitionName(String offsetsTopicName, int partitionId) { @@ -719,7 +719,7 @@ private boolean validateOffsetMetadataLength(String metadata) { public CompletableFuture scheduleLoadGroupAndOffsets(int offsetsPartition, Consumer onGroupLoaded) { - String topicPartition = offsetConfig.getCurrentOffsetsTopicName(tenant) + String topicPartition = offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + offsetsPartition; if (addLoadingPartition(offsetsPartition)) { log.info("Scheduling loading of offsets and group metadata from {}", topicPartition); @@ -1273,7 +1273,7 @@ CompletableFuture cleanGroupMetadata(Stream groups, log.error("Failed to append {} tombstones to topic {} for expired/deleted " + "offsets and/or metadata for group {}", tombstones.size(), - offsetConfig.getCurrentOffsetsTopicName(tenant) + offsetConfig.offsetsTopicName() + '-' + partitioner.apply(group.groupId()), group.groupId(), cause); // ignore and continue @@ -1375,11 +1375,11 @@ CompletableFuture> getOffsetsTopicProducer(String groupId) partitionId -> { if (log.isDebugEnabled()) { log.debug("Created Partitioned producer: {} for consumer group: {}", - offsetConfig.getCurrentOffsetsTopicName(tenant) + PARTITIONED_TOPIC_SUFFIX + partitionId, + offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + partitionId, groupId); } return metadataTopicProducerBuilder.clone() - .topic(offsetConfig.getCurrentOffsetsTopicName(tenant) + PARTITIONED_TOPIC_SUFFIX + partitionId) + .topic(offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + partitionId) .createAsync(); }); } @@ -1389,10 +1389,10 @@ CompletableFuture> getOffsetsTopicProducer(int partitionId) id -> { if (log.isDebugEnabled()) { log.debug("Will create Partitioned producer: {}", - offsetConfig.getCurrentOffsetsTopicName(tenant) + PARTITIONED_TOPIC_SUFFIX + id); + offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + id); } return metadataTopicProducerBuilder.clone() - .topic(offsetConfig.getCurrentOffsetsTopicName(tenant) + PARTITIONED_TOPIC_SUFFIX + id) + .topic(offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + id) .createAsync(); }); } @@ -1402,10 +1402,10 @@ CompletableFuture> getOffsetsTopicReader(int partitionId) { id -> { if (log.isDebugEnabled()) { log.debug("Will create Partitioned reader: {}", - offsetConfig.getCurrentOffsetsTopicName(tenant) + PARTITIONED_TOPIC_SUFFIX + id); + offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + id); } return metadataTopicReaderBuilder.clone() - .topic(offsetConfig.getCurrentOffsetsTopicName(tenant) + PARTITIONED_TOPIC_SUFFIX + partitionId) + .topic(offsetConfig.offsetsTopicName() + PARTITIONED_TOPIC_SUFFIX + partitionId) .readCompacted(true) .createAsync(); }); 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 4262b04d96..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,16 +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 = "${tenant}/default/__consumer_offsets"; public static final int DefaultOffsetsNumPartitions = KafkaServiceConfiguration.DefaultOffsetsTopicNumPartitions; @Default - private String offsetsTopicName = DefaultOffsetsTopicName; - - public String getCurrentOffsetsTopicName(String tenant) { - return offsetsTopicName.replace(KafkaServiceConfiguration.TENANT_PLACEHOLDER, tenant); - } - + @NonNull + private String offsetsTopicName = null; @Default private int maxMetadataSize = DefaultMaxMetadataSize; @Default