diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java index b1ba940220..8d58ae80fb 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java @@ -21,7 +21,9 @@ import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperation; import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperationPurgatory; import io.streamnative.pulsar.handlers.kop.utils.timer.SystemTimer; +import java.util.Collection; import java.util.Collections; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Optional; @@ -29,8 +31,10 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.locks.ReentrantReadWriteLock; import java.util.stream.Collectors; import lombok.extern.slf4j.Slf4j; +import org.apache.kafka.common.Node; import org.apache.kafka.common.config.ConfigResource; import org.apache.kafka.common.errors.InvalidRequestException; import org.apache.kafka.common.errors.TopicExistsException; @@ -52,6 +56,9 @@ class AdminManager { private final PulsarAdmin admin; private final int defaultNumPartitions; + private volatile Set brokersCache = new HashSet<>(); + private final ReentrantReadWriteLock brokersCacheLock = new ReentrantReadWriteLock(); + public AdminManager(PulsarAdmin admin, KafkaServiceConfiguration conf) { this.admin = admin; @@ -221,4 +228,17 @@ public Map deleteTopics(Set topicsToDelete) { }); return result; } + + public Collection getBrokers() { + return brokersCache; + } + + public void setBrokers(Set newBrokers) { + brokersCacheLock.writeLock().lock(); + try { + this.brokersCache = newBrokers; + } finally { + brokersCacheLock.writeLock().unlock(); + } + } } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/ChildChangeHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/ChildChangeHandler.java new file mode 100644 index 0000000000..7543f91c72 --- /dev/null +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/ChildChangeHandler.java @@ -0,0 +1,57 @@ +/** + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package io.streamnative.pulsar.handlers.kop; + +public interface ChildChangeHandler { + String path(); + + void handleChildChange(); +} + +class DeletionTopicsHandler implements ChildChangeHandler { + private final KopEventManager kopEventManager; + + public DeletionTopicsHandler(KopEventManager kopEventManager) { + this.kopEventManager = kopEventManager; + } + + @Override + public String path() { + return KopEventManager.getBrokersChangePath(); + } + + @Override + public void handleChildChange() { + kopEventManager.put(kopEventManager.getDeleteTopicEvent()); + } +} + +class BrokersChangeHandler implements ChildChangeHandler { + private final KopEventManager kopEventManager; + + public BrokersChangeHandler(KopEventManager kopEventManager) { + this.kopEventManager = kopEventManager; + } + + @Override + public String path() { + return KopEventManager.getBrokersChangePath(); + } + + @Override + public void handleChildChange() { + kopEventManager.put(kopEventManager.getBrokersChangeEvent()); + } + +} 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 6c27b4a048..278a689eff 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 @@ -95,6 +95,8 @@ public class KafkaProtocolHandler implements ProtocolHandler { private GroupCoordinator groupCoordinator; @Getter private TransactionCoordinator transactionCoordinator; + @Getter + private KopEventManager kopEventManager; /** * Listener for the changing of topic that stores offsets of consumer group. @@ -285,6 +287,12 @@ public void start(BrokerService service) { ZooKeeperUtils.tryCreatePath(brokerService.pulsar().getZkClient(), kafkaConfig.getGroupIdZooKeeperPath(), new byte[0]); + ZooKeeperUtils.tryCreatePath(brokerService.pulsar().getZkClient(), + KopEventManager.getKopPath(), new byte[0]); + + ZooKeeperUtils.tryCreatePath(brokerService.pulsar().getZkClient(), + KopEventManager.getDeleteTopicsPath(), new byte[0]); + PulsarAdmin pulsarAdmin; try { pulsarAdmin = brokerService.getPulsar().getAdminClient(); @@ -322,6 +330,12 @@ public void start(BrokerService service) { // init and start group coordinator startGroupCoordinator(pulsarClient); + // init KopEventManager + kopEventManager = new KopEventManager(groupCoordinator, + adminManager, + brokerService.getPulsar().getLocalMetadataStore()); + kopEventManager.start(); + // and listener for Offset topics load/unload brokerService.pulsar() .getNamespaceService() @@ -425,13 +439,13 @@ public void startGroupCoordinator(PulsarClient pulsarClient) { .build(); this.groupCoordinator = GroupCoordinator.of( - (PulsarClientImpl) pulsarClient, - groupConfig, - offsetConfig, - SystemTimer.builder() - .executorName("group-coordinator-timer") - .build(), - Time.SYSTEM + (PulsarClientImpl) pulsarClient, + groupConfig, + offsetConfig, + SystemTimer.builder() + .executorName("group-coordinator-timer") + .build(), + Time.SYSTEM ); // always enable metadata expiration this.groupCoordinator.startup(true); @@ -498,4 +512,5 @@ private void loadTxnLogTopics(TransactionCoordinator txnCoordinator) throws Exce public static @NonNull LookupClient getLookupClient(final PulsarService pulsarService) { return LOOKUP_CLIENT_MAP.computeIfAbsent(pulsarService, ignored -> new LookupClient(pulsarService)); } + } 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 cb10d044b2..b01350b9bd 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 @@ -492,6 +492,8 @@ protected void handleTopicMetadataRequest(KafkaHeaderAndRequest metadataHar, // Command response for all topics List allTopicMetadata = Collections.synchronizedList(Lists.newArrayList()); List allNodes = Collections.synchronizedList(Lists.newArrayList()); + // Get all kop brokers in local cache + allNodes.addAll(adminManager.getBrokers()); List topics = metadataRequest.topics(); // topics in format : persistent://%s/%s/abc-partition-x, will be grouped by as: @@ -672,13 +674,12 @@ protected void handleTopicMetadataRequest(KafkaHeaderAndRequest metadataHar, if (e != null) { log.warn("[{}] Request {}: Exception fetching metadata, will return null Response", ctx.channel(), metadataHar.getHeader(), e); - allNodes.add(newSelfNode()); MetadataResponse finalResponse = - new MetadataResponse( - allNodes, - clusterName, - controllerId, - Collections.emptyList()); + new MetadataResponse( + allNodes, + clusterName, + controllerId, + Collections.emptyList()); resultFuture.complete(finalResponse); return; } @@ -687,13 +688,12 @@ protected void handleTopicMetadataRequest(KafkaHeaderAndRequest metadataHar, if (topicsNumber == 0) { // no topic partitions added, return now. - allNodes.add(newSelfNode()); MetadataResponse finalResponse = - new MetadataResponse( - allNodes, - clusterName, - controllerId, - allTopicMetadata); + new MetadataResponse( + allNodes, + clusterName, + controllerId, + allTopicMetadata); resultFuture.complete(finalResponse); return; } @@ -2049,7 +2049,18 @@ protected void handleDeleteTopics(KafkaHeaderAndRequest deleteTopics, checkArgument(deleteTopics.getRequest() instanceof DeleteTopicsRequest); DeleteTopicsRequest request = (DeleteTopicsRequest) deleteTopics.getRequest(); Set topicsToDelete = request.topics(); - resultFuture.complete(new DeleteTopicsResponse(adminManager.deleteTopics(topicsToDelete))); + Map deleteTopicsResponse = adminManager.deleteTopics(topicsToDelete); + + // create topic znode to trigger the coordinator DeleteTopicsEvent event + deleteTopicsResponse.forEach((topic, errors) -> { + if (errors == Errors.NONE) { + ZooKeeperUtils.tryCreatePath(pulsarService.getZkClient(), + KopEventManager.getDeleteTopicsPath() + "/" + topic, + new byte[0]); + } + }); + + resultFuture.complete(new DeleteTopicsResponse(deleteTopicsResponse)); } /** diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopEventManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopEventManager.java new file mode 100644 index 0000000000..8a384015ab --- /dev/null +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopEventManager.java @@ -0,0 +1,308 @@ +/** + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package io.streamnative.pulsar.handlers.kop; + +import static com.google.common.base.Preconditions.checkState; + +import com.google.common.annotations.VisibleForTesting; +import com.google.common.collect.Sets; +import com.google.gson.JsonElement; +import com.google.gson.JsonObject; +import com.google.gson.JsonParser; +import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupCoordinator; +import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupMetadata; +import io.streamnative.pulsar.handlers.kop.utils.KopTopic; +import io.streamnative.pulsar.handlers.kop.utils.ShutdownableThread; +import java.nio.charset.StandardCharsets; +import java.util.Collection; +import java.util.HashSet; +import java.util.List; +import java.util.Optional; +import java.util.Set; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.locks.ReentrantLock; +import java.util.regex.Matcher; +import java.util.regex.Pattern; +import java.util.stream.Collectors; +import lombok.extern.slf4j.Slf4j; +import org.apache.kafka.common.Node; +import org.apache.kafka.common.TopicPartition; +import org.apache.pulsar.broker.loadbalance.LoadManager; +import org.apache.pulsar.common.util.Murmur3_32Hash; +import org.apache.pulsar.metadata.api.MetadataStore; +import org.apache.pulsar.metadata.api.Notification; + +@Slf4j +public class KopEventManager { + private static final String REGEX = "^(.*)://\\[?([0-9a-zA-Z\\-%._:]*)\\]?:(-?[0-9]+)"; + private static final Pattern PATTERN = Pattern.compile(REGEX); + + private static final String kopEventThreadName = "kop-event-thread"; + private final ReentrantLock putLock = new ReentrantLock(); + private static final LinkedBlockingQueue queue = + new LinkedBlockingQueue<>(); + private final KopEventThread thread = + new KopEventThread(kopEventThreadName); + private final GroupCoordinator coordinator; + private final AdminManager adminManager; + private final DeletionTopicsHandler deletionTopicsHandler; + private final BrokersChangeHandler brokersChangeHandler; + private final MetadataStore metadataStore; + + public KopEventManager(GroupCoordinator coordinator, + AdminManager adminManager, + MetadataStore metadataStore) { + this.coordinator = coordinator; + this.adminManager = adminManager; + this.deletionTopicsHandler = new DeletionTopicsHandler(this); + this.brokersChangeHandler = new BrokersChangeHandler(this); + this.metadataStore = metadataStore; + } + + public void start() { + registerChildChangeHandler(); + thread.start(); + } + + public void close() { + try { + thread.shutdown(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + log.error("Interrupted at shutting down {}", kopEventThreadName); + } + + } + + + public void put(KopEvent event) { + putLock.lock(); + try { + queue.put(event); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + log.error("Error put event {} to coordinator event queue", event, e); + } finally { + putLock.unlock(); + } + } + + public void clearAndPut(KopEvent event) { + putLock.lock(); + try { + queue.clear(); + put(event); + } finally { + putLock.unlock(); + } + } + + static class KopEventThread extends ShutdownableThread { + + public KopEventThread(String name) { + super(name); + } + + @Override + protected void doWork() { + KopEvent event = null; + try { + event = queue.take(); + event.process(); + } catch (InterruptedException e) { + log.error("Error processing event {}", event, e); + } + } + + } + + private void registerChildChangeHandler() { + metadataStore.registerListener(this::handleChildChangePathNotification); + + // Really register ChildChange notification. + metadataStore.getChildren(getDeleteTopicsPath()); + // init local kop brokers cache + getBrokers(metadataStore.getChildren(getBrokersChangePath()).join()); + } + + private void handleChildChangePathNotification(Notification notification) { + if (notification.getPath().equals(LoadManager.LOADBALANCE_BROKERS_ROOT)) { + this.brokersChangeHandler.handleChildChange(); + } else if (notification.getPath().equals(getDeleteTopicsPath())) { + this.deletionTopicsHandler.handleChildChange(); + } + } + + private void getBrokers(List pulsarBrokers) { + final Set kopBrokers = Sets.newConcurrentHashSet(); + final AtomicInteger pendingBrokers = new AtomicInteger(pulsarBrokers.size()); + + pulsarBrokers.forEach(broker -> { + metadataStore.get(getBrokersChangePath() + "/" + broker).whenComplete( + (brokerData, e) -> { + if (e != null) { + log.error("Get broker {} path data failed which have an error", broker, e); + return; + } + + if (brokerData.isPresent()) { + JsonObject jsonObject = parseJsonObject( + new String(brokerData.get().getValue(), StandardCharsets.UTF_8)); + JsonObject protocols = jsonObject.getAsJsonObject("protocols"); + JsonElement element = protocols.get("kafka"); + + if (element != null) { + String kopBrokerStr = element.getAsString(); + Node kopNode = getNode(kopBrokerStr); + kopBrokers.add(kopNode); + } else { + if (log.isDebugEnabled()) { + log.debug("Get broker {} path currently not a kop broker, skip it.", broker); + } + } + } else { + if (log.isDebugEnabled()) { + log.debug("Get broker {} path data empty.", broker); + } + } + + if (pendingBrokers.decrementAndGet() == 0) { + Collection oldKopBrokers = adminManager.getBrokers(); + adminManager.setBrokers(kopBrokers); + log.info("Refresh kop brokers new cache {}, old brokers cache {}", + adminManager.getBrokers(), oldKopBrokers); + } + } + ); + }); + + } + + private JsonObject parseJsonObject(String info) { + JsonParser parser = new JsonParser(); + return parser.parse(info).getAsJsonObject(); + } + + @VisibleForTesting + public static Node getNode(String kopBrokerStr) { + final String errorMessage = "kopBrokerStr " + kopBrokerStr + " is invalid"; + final Matcher matcher = PATTERN.matcher(kopBrokerStr); + checkState(matcher.find(), errorMessage); + checkState(matcher.groupCount() == 3, errorMessage); + String host = matcher.group(2); + String port = matcher.group(3); + + return new Node( + Murmur3_32Hash.getInstance().makeHash((host + port).getBytes(StandardCharsets.UTF_8)), + host, + Integer.parseInt(port)); + } + + + interface KopEvent { + void process(); + } + + class DeleteTopicsEvent implements KopEvent { + + @Override + public void process() { + if (!coordinator.isActive()) { + return; + } + + try { + List topicsDeletions = metadataStore.getChildren(getDeleteTopicsPath()).get(); + + HashSet topicsFullNameDeletionsSets = Sets.newHashSet(); + HashSet kopTopicsSet = Sets.newHashSet(); + topicsDeletions.forEach(topic -> { + KopTopic kopTopic = new KopTopic(topic); + kopTopicsSet.add(kopTopic); + topicsFullNameDeletionsSets.add(kopTopic.getFullName()); + }); + + log.debug("Delete topics listener fired for topics {} to be deleted", topicsDeletions); + Iterable groupMetadataIterable = coordinator.getGroupManager().currentGroups(); + HashSet topicPartitionsToBeDeletions = Sets.newHashSet(); + + groupMetadataIterable.forEach(groupMetadata -> { + topicPartitionsToBeDeletions.addAll( + groupMetadata.collectPartitionsWithTopics(topicsFullNameDeletionsSets)); + }); + + Set deletedTopics = Sets.newHashSet(); + if (!topicPartitionsToBeDeletions.isEmpty()) { + coordinator.handleDeletedPartitions(topicPartitionsToBeDeletions); + Set collectDeleteTopics = topicPartitionsToBeDeletions + .stream() + .map(TopicPartition::topic) + .collect(Collectors.toSet()); + + deletedTopics = kopTopicsSet.stream().filter( + kopTopic -> collectDeleteTopics.contains(kopTopic.getFullName()) + ).map(KopTopic::getOriginalName).collect(Collectors.toSet()); + + deletedTopics.forEach(deletedTopic -> { + metadataStore.delete( + getDeleteTopicsPath() + "/" + deletedTopic, Optional.of((long) -1)); + }); + } + + log.info("GroupMetadata delete topics {}, no matching topics {}", + deletedTopics, Sets.difference(topicsFullNameDeletionsSets, deletedTopics)); + + } catch (ExecutionException | InterruptedException e) { + log.error("DeleteTopicsEvent process have an error", e); + } + } + } + + class BrokersChangeEvent implements KopEvent { + @Override + public void process() { + metadataStore.getChildren(getBrokersChangePath()).whenComplete( + (brokers, e) -> { + if (e != null) { + log.error("BrokersChangeEvent process have an error", e); + return; + } + getBrokers(brokers); + }); + } + } + + public DeleteTopicsEvent getDeleteTopicEvent() { + return new DeleteTopicsEvent(); + } + + public BrokersChangeEvent getBrokersChangeEvent() { + return new BrokersChangeEvent(); + } + + public static String getKopPath() { + return "/kop"; + } + + public static String getDeleteTopicsPath() { + return getKopPath() + "/delete_topics"; + } + + public static String getBrokersChangePath() { + return LoadManager.LOADBALANCE_BROKERS_ROOT; + } + +} 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 0dc6627266..982a03c40a 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 @@ -115,11 +115,11 @@ public static GroupCoordinator of( .build(); return new GroupCoordinator( - groupConfig, - metadataManager, - heartbeatPurgatory, - joinPurgatory, - time + groupConfig, + metadataManager, + heartbeatPurgatory, + joinPurgatory, + time ); } @@ -154,7 +154,8 @@ public GroupCoordinator( GroupMetadataManager groupManager, DelayedOperationPurgatory heartbeatPurgatory, DelayedOperationPurgatory joinPurgatory, - Time time) { + Time time + ) { this.groupConfig = groupConfig; this.groupManager = groupManager; this.heartbeatPurgatory = heartbeatPurgatory; @@ -894,7 +895,7 @@ public KeyValue handleDescribeGroup(String groupId) { ); } - public CompletableFuture handleDeletedPartitions(List topicPartitions) { + public CompletableFuture handleDeletedPartitions(Set topicPartitions) { return groupManager.cleanGroupMetadata(groupManager.currentGroupsStream(), group -> group.removeOffsets(topicPartitions.stream()) ).thenApply(offsetsRemoved -> { @@ -1297,4 +1298,7 @@ private boolean isCoordinatorLoadInProgress(String groupId) { return groupManager.isGroupLoading(groupId); } + public boolean isActive() { + return isActive.get(); + } } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupMetadata.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupMetadata.java index c70b1bbf56..a2b478e24b 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupMetadata.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupMetadata.java @@ -21,12 +21,14 @@ import com.google.common.base.MoreObjects; import com.google.common.base.MoreObjects.ToStringHelper; import com.google.common.base.Supplier; +import com.google.common.collect.Lists; import com.google.common.collect.Sets; import io.streamnative.pulsar.handlers.kop.coordinator.group.MemberMetadata.MemberSummary; import io.streamnative.pulsar.handlers.kop.exceptions.KoPTopicException; import io.streamnative.pulsar.handlers.kop.offset.OffsetAndMetadata; import io.streamnative.pulsar.handlers.kop.utils.CoreUtils; import io.streamnative.pulsar.handlers.kop.utils.KopTopic; +import java.util.ArrayList; import java.util.Collections; import java.util.Comparator; import java.util.HashMap; @@ -588,6 +590,13 @@ public Map removeOffsets(Stream( + topicPartition, + OffsetAndMetadata.apply(0) + ); + } return new KeyValue<>( topicPartition, removedOffset.offsetAndMetadata() @@ -598,6 +607,26 @@ public Map removeOffsets(Stream collectPartitionsWithTopics(Set topics) { + ArrayList topicPartitions = Lists.newArrayList(); + + topicPartitions.addAll(pendingOffsetCommits.keySet().stream().filter( + topicPartition -> topics.contains(topicPartition.topic()) + ).collect(Collectors.toList())); + + pendingTransactionalOffsetCommits.values().stream().map(Map::keySet) + .collect(Collectors.toList()).forEach(partitionSet -> { + topicPartitions.addAll(partitionSet.stream().filter( + topicPartition -> topics.contains(topicPartition.topic())) + .collect(Collectors.toList())); + }); + + topicPartitions.addAll(offsets.keySet().stream().filter( + topicPartition -> topics.contains(topicPartition.topic()) + ).collect(Collectors.toList())); + return topicPartitions; + } + public Map removeExpiredOffsets(long startMs) { Map expiredOffsets = offsets.entrySet().stream() .filter(e -> diff --git a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/KopEventManagerTest.java b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/KopEventManagerTest.java new file mode 100644 index 0000000000..805fce2396 --- /dev/null +++ b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/KopEventManagerTest.java @@ -0,0 +1,33 @@ +/** + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package io.streamnative.pulsar.handlers.kop; + +import org.apache.kafka.common.Node; +import org.testng.Assert; +import org.testng.annotations.Test; + +public class KopEventManagerTest { + + @Test + public void testGetNode() { + final String host = "localhost"; + final int port = 9120; + final String securityProtocol = "SASL_PLAINTEXT"; + final String brokerStr = securityProtocol + "://" + host + ":" + port; + Node node = KopEventManager.getNode(brokerStr); + Assert.assertEquals(node.host(), host); + Assert.assertEquals(node.port(), port); + } + +} diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java index 9b374f0616..a445d92be9 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java @@ -710,4 +710,34 @@ public void testMetadataForNonPartitionedTopic(short version) throws Exception { assertEquals(response.topicMetadata().size(), 1); assertEquals(response.errors().size(), 0); } + + @Test(timeOut = 10000) + public void testDeleteTopicsAndCheckChildPath() throws Exception { + Properties props = new Properties(); + props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:" + getKafkaBrokerPort()); + + @Cleanup + AdminClient kafkaAdmin = AdminClient.create(props); + Map topicToNumPartitions = new HashMap(){{ + put("testCreateTopics-0", 1); + put("testCreateTopics-1", 3); + put("my-tenant/my-ns/testCreateTopics-2", 1); + put("persistent://my-tenant/my-ns/testCreateTopics-3", 5); + }}; + // create + createTopicsByKafkaAdmin(kafkaAdmin, topicToNumPartitions); + verifyTopicsCreatedByPulsarAdmin(topicToNumPartitions); + // delete + deleteTopicsByKafkaAdmin(kafkaAdmin, topicToNumPartitions.keySet()); + verifyTopicsDeletedByPulsarAdmin(topicToNumPartitions); + // check deleted topics path + List deletedTopics = handler.getPulsarService() + .getBrokerService() + .getPulsar() + .getLocalMetadataStore() + .getChildren(KopEventManager.getDeleteTopicsPath()) + .join(); + + assertTrue(topicToNumPartitions.keySet().containsAll(deletedTopics)); + } } diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KopEventManagerTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KopEventManagerTest.java new file mode 100644 index 0000000000..1f94b197c7 --- /dev/null +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KopEventManagerTest.java @@ -0,0 +1,206 @@ +/** + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package io.streamnative.pulsar.handlers.kop; + +import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertTrue; + +import com.google.common.collect.Lists; +import java.time.Duration; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.Properties; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.TimeUnit; +import org.apache.kafka.clients.admin.AdminClient; +import org.apache.kafka.clients.admin.AdminClientConfig; +import org.apache.kafka.clients.admin.ConsumerGroupDescription; +import org.apache.kafka.clients.admin.NewTopic; +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.consumer.ConsumerRecords; +import org.apache.kafka.clients.consumer.KafkaConsumer; +import org.apache.kafka.clients.producer.KafkaProducer; +import org.apache.kafka.clients.producer.ProducerConfig; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.ConsumerGroupState; +import org.apache.kafka.common.serialization.StringDeserializer; +import org.apache.kafka.common.serialization.StringSerializer; +import org.testng.annotations.AfterMethod; +import org.testng.annotations.BeforeMethod; +import org.testng.annotations.Test; + +public class KopEventManagerTest extends KopProtocolHandlerTestBase { + private AdminClient adminClient; + private KafkaProducer kafkaProducer; + private String broker; + private final String topic1 = "test-topic1"; + private final String topic2 = "test-topic2"; + private final String topic3 = "test-topic3"; + + @BeforeMethod + @Override + protected void setup() throws Exception { + super.internalSetup(); + final EndPoint plainEndPoint = getPlainEndPoint(); + this.broker = plainEndPoint.getHostname() + ":" + plainEndPoint.getPort(); + Properties adminPro = new Properties(); + adminPro.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, broker); + this.adminClient = AdminClient.create(adminPro); + final Properties producerPro = new Properties(); + producerPro.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, broker); + producerPro.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); + producerPro.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); + this.kafkaProducer = new KafkaProducer<>(producerPro); + } + + @AfterMethod + @Override + protected void cleanup() throws Exception { + adminClient.close(); + kafkaProducer.close(); + super.internalCleanup(); + } + + @Test(timeOut = 6000) + public void testOneTopicGroupState() throws Exception { + // 1. create topics + createTopics(Collections.singletonList(topic1)); + // 2. send messages + sendOneMessages(topic1); + // 3. check group state which only consumed one topic + final Properties properties = new Properties(); + properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, broker); + properties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + properties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + final String groupId1 = "test-group1"; + properties.put(ConsumerConfig.GROUP_ID_CONFIG, groupId1); + + final KafkaConsumer kafkaConsumer1 = new KafkaConsumer<>(properties); + kafkaConsumer1.subscribe(Collections.singletonList(topic1)); + ConsumerRecords records = kafkaConsumer1.poll(Duration.ofMillis(500)); + assertEquals(records.count(), 1); + + // 4. check group state must be Stable + Map describeGroup1 = + adminClient.describeConsumerGroups(Collections.singletonList(groupId1)) + .all() + .get(1000, TimeUnit.MILLISECONDS); + assertTrue(describeGroup1.containsKey(groupId1)); + assertEquals(ConsumerGroupState.STABLE, describeGroup1.get(groupId1).state()); + // 5. close consumer1 + kafkaConsumer1.close(); + // 6. check group state must be Empty + Map describeGroup2 = + adminClient.describeConsumerGroups(Collections.singletonList(groupId1)) + .all() + .get(1000, TimeUnit.MILLISECONDS); + assertTrue(describeGroup1.containsKey(groupId1)); + assertEquals(ConsumerGroupState.EMPTY, describeGroup2.get(groupId1).state()); + // 7. delete topic1 + adminClient.deleteTopics(Collections.singletonList(topic1)); + // 8. describe group who only consume topic1 which have been deleted + // check group state must be Dead + retryUntilStateDead(groupId1, 5); + } + + @Test(timeOut = 6000) + public void testTwoTopicsGroupState() throws Exception { + // 1. create topics + createTopics(Arrays.asList(topic2, topic3)); + // 2. send messages + sendOneMessages(topic2); + sendOneMessages(topic3); + + // 3. check group state which consumed two topics + final Properties properties = new Properties(); + properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, broker); + properties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + properties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + final String groupId2 = "test-group2"; + properties.put(ConsumerConfig.GROUP_ID_CONFIG, groupId2); + final KafkaConsumer kafkaConsumer2 = new KafkaConsumer<>(properties); + kafkaConsumer2.subscribe(Arrays.asList(topic2, topic3)); + int consumeCount = 0; + while (consumeCount < 2) { + ConsumerRecords records = kafkaConsumer2.poll(Duration.ofMillis(500)); + consumeCount += records.count(); + } + // 4. check group state must be Stable + Map describeGroup1 = + adminClient.describeConsumerGroups(Collections.singletonList(groupId2)) + .all() + .get(1000, TimeUnit.MILLISECONDS); + assertTrue(describeGroup1.containsKey(groupId2)); + assertEquals(ConsumerGroupState.STABLE, describeGroup1.get(groupId2).state()); + + // 5. close consumer2 + kafkaConsumer2.close(); + + // 6. check group state must be Empty + Map describeGroup2 = + adminClient.describeConsumerGroups(Collections.singletonList(groupId2)) + .all() + .get(1000, TimeUnit.MILLISECONDS); + assertTrue(describeGroup2.containsKey(groupId2)); + assertEquals(ConsumerGroupState.EMPTY, describeGroup2.get(groupId2).state()); + + // 7. delete topic2 and topic3 + List deleteTopics = Lists.newArrayList(); + deleteTopics.add(topic2); + deleteTopics.add(topic3); + adminClient.deleteTopics(deleteTopics); + // 8. check group state must be Dead + retryUntilStateDead(groupId2, 5); + } + + private void createTopics(List topics) throws ExecutionException, InterruptedException { + List topicsList = Lists.newArrayList(); + topics.forEach( + topic -> { + NewTopic newTopic = new NewTopic(topic, 1, (short) 1); + topicsList.add(newTopic); + } + ); + adminClient.createTopics(topicsList).all().get(); + } + + private void sendOneMessages(String topic) { + kafkaProducer.send(new ProducerRecord<>(topic, null, "test-value")); + } + + private void retryUntilStateDead(String groupId, int timeOutSec) throws Exception { + long startTimeMs = System.currentTimeMillis(); + long deadTimeMs = startTimeMs + timeOutSec * 1000L; + + Map describeGroup = null; + + while (System.currentTimeMillis() < deadTimeMs) { + describeGroup = adminClient.describeConsumerGroups(Collections.singletonList(groupId)) + .all() + .get(1000, TimeUnit.MILLISECONDS); + assertTrue(describeGroup.containsKey(groupId)); + if (describeGroup.get(groupId).state().name().equals("DEAD")) { + break; + } + } + + assertEquals(ConsumerGroupState.DEAD, describeGroup.get(groupId).state()); + + } + +}