From 515df6892377b2db3329a522098624d65d8ede33 Mon Sep 17 00:00:00 2001 From: aloyszhang Date: Sat, 9 Oct 2021 16:38:02 +0800 Subject: [PATCH 1/9] fix offset data loss --- .../handlers/kop/KafkaProtocolHandler.java | 34 ++++++++++++++++--- .../group/GroupMetadataManagerTest.java | 20 +++++++++++ 2 files changed, 50 insertions(+), 4 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 abb1f9281a..d28e4fb098 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 @@ -60,6 +60,7 @@ import org.apache.pulsar.broker.service.BrokerService; import org.apache.pulsar.client.admin.Lookup; import org.apache.pulsar.client.admin.PulsarAdmin; +import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.naming.NamespaceBundle; import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.naming.TopicName; @@ -101,6 +102,10 @@ public GroupCoordinator getGroupCoordinator(String tenant) { return groupCoordinatorsByTenant.computeIfAbsent(tenant, this::createAndBootGroupCoordinator); } + public Map getGroupCoordinator() { + return groupCoordinatorsByTenant; + } + @Override public TransactionCoordinator getTransactionCoordinator(String tenant) { return transactionCoordinatorByTenant.computeIfAbsent(tenant, this::createAndBootTransactionCoordinator); @@ -485,11 +490,32 @@ private GroupCoordinator startGroupCoordinator(String tenant, SystemTopicClient kafkaConfig.getGroupInitialRebalanceDelayMs() ); + String topicName = tenant + "/" + kafkaConfig.getKafkaMetadataNamespace() + + "/" + Topic.GROUP_METADATA_TOPIC_NAME; + + PulsarAdmin pulsarAdmin; + int existedOffsetTopicNumPartitions = 0; + try { + pulsarAdmin = brokerService.getPulsar().getAdminClient(); + existedOffsetTopicNumPartitions = pulsarAdmin.topics().getPartitionedTopicMetadata(topicName).partitions; + } catch (PulsarServerException | PulsarAdminException e) { + log.error("Failed to get offset topic partition metadata .", e); + throw new IllegalStateException(e); + } + + int offsetTopicNumPartitions; + if (existedOffsetTopicNumPartitions == 0) { + log.info("Not existed offset topic number partitions found, use the default {}", + kafkaConfig.getOffsetsTopicNumPartitions()); + offsetTopicNumPartitions = kafkaConfig.getOffsetsTopicNumPartitions(); + } else { + log.info("Existed offset topic number partitions found {}", existedOffsetTopicNumPartitions); + offsetTopicNumPartitions = existedOffsetTopicNumPartitions; + } + OffsetConfig offsetConfig = OffsetConfig.builder() - .offsetsTopicName(tenant + "/" - + kafkaConfig.getKafkaMetadataNamespace() - + "/" + Topic.GROUP_METADATA_TOPIC_NAME) - .offsetsTopicNumPartitions(kafkaConfig.getOffsetsTopicNumPartitions()) + .offsetsTopicName(topicName) + .offsetsTopicNumPartitions(offsetTopicNumPartitions) .offsetsTopicCompressionType(CompressionType.valueOf(kafkaConfig.getOffsetsTopicCompressionCodec())) .maxMetadataSize(kafkaConfig.getOffsetMetadataMaxSize()) .offsetsRetentionCheckIntervalMs(kafkaConfig.getOffsetsRetentionCheckIntervalMs()) 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 8697823655..bf70d7c54a 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 @@ -45,6 +45,7 @@ import io.streamnative.pulsar.handlers.kop.offset.OffsetAndMetadata; import io.streamnative.pulsar.handlers.kop.utils.KopTopic; import io.streamnative.pulsar.handlers.kop.utils.timer.MockTime; +import java.lang.reflect.Field; import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.Collections; @@ -286,6 +287,25 @@ private int completeTransactionalOffsetCommit(ByteBuffer buffer, return 1; } + @Test + public void testOffsetTopicNumPartitionsModify() throws Exception { + int consumerGroupPartitionId = + GroupMetadataManager.getPartitionId(groupId, conf.getOffsetsTopicNumPartitions()); + Field partitionField = conf.getClass().getDeclaredField("offsetsTopicNumPartitions"); + partitionField.setAccessible(true); + partitionField.set(conf, 100); + + KafkaProtocolHandler handler = (KafkaProtocolHandler) pulsar.getProtocolHandlers().protocol("kafka"); + // remove here to trigger a new creating for GroupCoordinator + handler.getGroupCoordinator().remove(conf.getKafkaMetadataTenant()); + GroupMetadataManager newMetaManager = + handler.getGroupCoordinator(conf.getKafkaMetadataTenant()).getGroupManager(); + + int newPartitionsId = + GroupMetadataManager.getPartitionId(groupId, newMetaManager.offsetConfig().offsetsTopicNumPartitions()); + assertEquals(consumerGroupPartitionId, newPartitionsId); + } + @Test public void testLoadOffsetsWithoutGroup() throws Exception { Map committedOffsets = new HashMap<>(); From da42f73a9aed4d688b922f7219ab0fb6d96c6aa8 Mon Sep 17 00:00:00 2001 From: aloyszhang Date: Sat, 9 Oct 2021 17:11:28 +0800 Subject: [PATCH 2/9] apply comment --- .../pulsar/handlers/kop/KafkaProtocolHandler.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) 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 d28e4fb098..daa1d80ede 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 @@ -18,6 +18,7 @@ import static io.streamnative.pulsar.handlers.kop.utils.TopicNameUtils.getKafkaTopicNameFromPulsarTopicname; import static org.apache.pulsar.common.naming.TopicName.PARTITIONED_TOPIC_SUFFIX; +import com.google.common.annotations.VisibleForTesting; import com.google.common.collect.ImmutableMap; import io.netty.channel.ChannelInitializer; import io.netty.channel.socket.SocketChannel; @@ -102,7 +103,8 @@ public GroupCoordinator getGroupCoordinator(String tenant) { return groupCoordinatorsByTenant.computeIfAbsent(tenant, this::createAndBootGroupCoordinator); } - public Map getGroupCoordinator() { + @VisibleForTesting + public Map getGroupCoordinator() { return groupCoordinatorsByTenant; } From d415b93dab5e010f0aa1588c60b05bf820a5ed92 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Sat, 9 Oct 2021 17:12:46 +0800 Subject: [PATCH 3/9] Remove redundant blank --- .../streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 daa1d80ede..abe08cec52 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 @@ -499,7 +499,7 @@ private GroupCoordinator startGroupCoordinator(String tenant, SystemTopicClient int existedOffsetTopicNumPartitions = 0; try { pulsarAdmin = brokerService.getPulsar().getAdminClient(); - existedOffsetTopicNumPartitions = pulsarAdmin.topics().getPartitionedTopicMetadata(topicName).partitions; + existedOffsetTopicNumPartitions = pulsarAdmin.topics().getPartitionedTopicMetadata(topicName).partitions; } catch (PulsarServerException | PulsarAdminException e) { log.error("Failed to get offset topic partition metadata .", e); throw new IllegalStateException(e); From 767c00512b1fc6b422d74337e2882d1d0584d76b Mon Sep 17 00:00:00 2001 From: aloyszhang Date: Sat, 9 Oct 2021 17:17:21 +0800 Subject: [PATCH 4/9] apply comment --- .../streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java | 2 +- .../kop/coordinator/group/GroupMetadataManagerTest.java | 2 +- 2 files changed, 2 insertions(+), 2 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 abe08cec52..32976f0b56 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 @@ -104,7 +104,7 @@ public GroupCoordinator getGroupCoordinator(String tenant) { } @VisibleForTesting - public Map getGroupCoordinator() { + public Map getGroupCoordinators() { return groupCoordinatorsByTenant; } 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 bf70d7c54a..a5affee66d 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 @@ -297,7 +297,7 @@ public void testOffsetTopicNumPartitionsModify() throws Exception { KafkaProtocolHandler handler = (KafkaProtocolHandler) pulsar.getProtocolHandlers().protocol("kafka"); // remove here to trigger a new creating for GroupCoordinator - handler.getGroupCoordinator().remove(conf.getKafkaMetadataTenant()); + handler.getGroupCoordinators().remove(conf.getKafkaMetadataTenant()); GroupMetadataManager newMetaManager = handler.getGroupCoordinator(conf.getKafkaMetadataTenant()).getGroupManager(); From 7d6d61c6fc1c3aea2f167c852c9961b6437d34b7 Mon Sep 17 00:00:00 2001 From: aloyszhang Date: Sat, 9 Oct 2021 19:08:21 +0800 Subject: [PATCH 5/9] apply comment --- .../handlers/kop/KafkaProtocolHandler.java | 17 ++++++----------- 1 file changed, 6 insertions(+), 11 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 32976f0b56..89c1955af9 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 @@ -496,24 +496,19 @@ private GroupCoordinator startGroupCoordinator(String tenant, SystemTopicClient + "/" + Topic.GROUP_METADATA_TOPIC_NAME; PulsarAdmin pulsarAdmin; - int existedOffsetTopicNumPartitions = 0; + int offsetTopicNumPartitions; try { pulsarAdmin = brokerService.getPulsar().getAdminClient(); - existedOffsetTopicNumPartitions = pulsarAdmin.topics().getPartitionedTopicMetadata(topicName).partitions; + offsetTopicNumPartitions = pulsarAdmin.topics().getPartitionedTopicMetadata(topicName).partitions; + if (offsetTopicNumPartitions == 0) { + log.error("Offset topic should not be a non-partitioned topic."); + throw new IllegalStateException("Offset topic should not be a non-partitioned topic."); + } } catch (PulsarServerException | PulsarAdminException e) { log.error("Failed to get offset topic partition metadata .", e); throw new IllegalStateException(e); } - int offsetTopicNumPartitions; - if (existedOffsetTopicNumPartitions == 0) { - log.info("Not existed offset topic number partitions found, use the default {}", - kafkaConfig.getOffsetsTopicNumPartitions()); - offsetTopicNumPartitions = kafkaConfig.getOffsetsTopicNumPartitions(); - } else { - log.info("Existed offset topic number partitions found {}", existedOffsetTopicNumPartitions); - offsetTopicNumPartitions = existedOffsetTopicNumPartitions; - } OffsetConfig offsetConfig = OffsetConfig.builder() .offsetsTopicName(topicName) From 4b323945ef17bc9215a5c02a35ca4ff92d6c30ec Mon Sep 17 00:00:00 2001 From: aloyszhang Date: Sat, 9 Oct 2021 23:40:26 +0800 Subject: [PATCH 6/9] fix test --- .../kop/coordinator/group/GroupMetadataManager.java | 4 +++- .../coordinator/group/GroupMetadataManagerTest.java | 10 +++++----- 2 files changed, 8 insertions(+), 6 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 0a618b3cd9..0280818a9e 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 @@ -26,6 +26,7 @@ import static org.apache.kafka.common.internals.Topic.GROUP_METADATA_TOPIC_NAME; import static org.apache.pulsar.common.naming.TopicName.PARTITIONED_TOPIC_SUFFIX; +import com.google.common.annotations.VisibleForTesting; import com.google.common.collect.Lists; import com.google.common.collect.Sets; import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupMetadata.CommitRecordMetadataAndOffset; @@ -1364,7 +1365,8 @@ boolean addLoadingPartition(int partition) { * *

Visible for testing */ - boolean removeLoadingPartition(int partition) { + @VisibleForTesting + public boolean removeLoadingPartition(int partition) { return inLock(partitionLock, () -> loadingPartitions.remove(partition)); } 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 a5affee66d..2e423f3837 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 @@ -1052,17 +1052,17 @@ public void testOffsetWriteAfterGroupRemoved() throws Exception { ByteBuffer buffer = newMemoryRecordsBuffer(newOffsetCommitRecords); byte[] key = groupMetadataKey(groupId); - - Producer producer = groupMetadataManager.getOffsetsTopicProducer(groupPartitionId).get(); + int consumerGroupPartitionId = + GroupMetadataManager.getPartitionId(groupId, conf.getOffsetsTopicNumPartitions()); + Producer producer = groupMetadataManager.getOffsetsTopicProducer(consumerGroupPartitionId).get(); producer.newMessage() .keyBytes(key) .value(buffer) .eventTime(Time.SYSTEM.milliseconds()) .send(); - + groupMetadataManager.removeLoadingPartition(consumerGroupPartitionId); CompletableFuture onLoadedFuture = new CompletableFuture<>(); - groupMetadataManager.scheduleLoadGroupAndOffsets( - groupPartitionId, + groupMetadataManager.scheduleLoadGroupAndOffsets(consumerGroupPartitionId, groupMetadata -> onLoadedFuture.complete(groupMetadata) ).get(); GroupMetadata group = onLoadedFuture.get(); From 4bcb128c023a8739ef28001463c160cd9c617210 Mon Sep 17 00:00:00 2001 From: aloyszhang Date: Mon, 11 Oct 2021 13:51:43 +0800 Subject: [PATCH 7/9] fix test --- .../kop/coordinator/group/GroupMetadataManagerTest.java | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) 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 2e423f3837..4d8bb4f1d0 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 @@ -291,9 +291,7 @@ private int completeTransactionalOffsetCommit(ByteBuffer buffer, public void testOffsetTopicNumPartitionsModify() throws Exception { int consumerGroupPartitionId = GroupMetadataManager.getPartitionId(groupId, conf.getOffsetsTopicNumPartitions()); - Field partitionField = conf.getClass().getDeclaredField("offsetsTopicNumPartitions"); - partitionField.setAccessible(true); - partitionField.set(conf, 100); + conf.setOffsetsTopicNumPartitions(100); KafkaProtocolHandler handler = (KafkaProtocolHandler) pulsar.getProtocolHandlers().protocol("kafka"); // remove here to trigger a new creating for GroupCoordinator From e2ca1b7db68a4cd6bb85dc018ab0d9a5cf738dec Mon Sep 17 00:00:00 2001 From: aloyszhang Date: Mon, 11 Oct 2021 13:57:03 +0800 Subject: [PATCH 8/9] remove ununsed import --- .../handlers/kop/coordinator/group/GroupMetadataManagerTest.java | 1 - 1 file changed, 1 deletion(-) 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 4d8bb4f1d0..7d7b7e1e9a 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 @@ -45,7 +45,6 @@ import io.streamnative.pulsar.handlers.kop.offset.OffsetAndMetadata; import io.streamnative.pulsar.handlers.kop.utils.KopTopic; import io.streamnative.pulsar.handlers.kop.utils.timer.MockTime; -import java.lang.reflect.Field; import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.Collections; From 75342a20180600d64dfdefc2a4a5c9a629274144 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Mon, 11 Oct 2021 14:58:34 +0800 Subject: [PATCH 9/9] Remove exception signature Co-authored-by: Kai Wang --- .../kop/coordinator/group/GroupMetadataManagerTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 7d7b7e1e9a..10e72c1935 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 @@ -287,7 +287,7 @@ private int completeTransactionalOffsetCommit(ByteBuffer buffer, } @Test - public void testOffsetTopicNumPartitionsModify() throws Exception { + public void testOffsetTopicNumPartitionsModify() { int consumerGroupPartitionId = GroupMetadataManager.getPartitionId(groupId, conf.getOffsetsTopicNumPartitions()); conf.setOffsetsTopicNumPartitions(100);