From 87166ab7f44aa2d8ec64392588b37d06de6288a2 Mon Sep 17 00:00:00 2001 From: Balaji Seshadri Date: Thu, 9 Oct 2014 15:25:16 -0600 Subject: [PATCH] KAFKA-1476 Get a list of consumer groups --- .../main/scala/kafka/tools/ConsumerCommand.scala | 78 ++++++++++++++++++++++ 1 file changed, 78 insertions(+) create mode 100644 core/src/main/scala/kafka/tools/ConsumerCommand.scala diff --git a/core/src/main/scala/kafka/tools/ConsumerCommand.scala b/core/src/main/scala/kafka/tools/ConsumerCommand.scala new file mode 100644 index 0000000..9e67c8d --- /dev/null +++ b/core/src/main/scala/kafka/tools/ConsumerCommand.scala @@ -0,0 +1,78 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You 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 kafka.tools + +import joptsimple.OptionParser +import org.I0Itec.zkclient.ZkClient +import kafka.network.BlockingChannel +import kafka.utils.{ZKStringSerializer, CommandLineUtils} +import collection.JavaConversions._ + +object ConsumerCommand { + + def main(args: Array[String]) { + + val parser = new OptionParser() + + val zkConnectOpt = parser.accepts("zookeeper", "ZooKeeper connect string."). + withRequiredArg().defaultsTo("localhost:2181").ofType(classOf[String]) + val topicsOpt = parser.accepts("topic", + "Comma-separated list of consumer topics (all topics if absent).") + val listGroupOpt = parser.accepts("list", " List Consumer groups.") + + parser.accepts("help", "Print this message.") + + if(args.length == 0) + CommandLineUtils.printUsageAndDie(parser, "List consumer groups or check offset lag or describe group.") + + val options = parser.parse(args: _*) + + if(options.has("help")) { + parser.printHelpOn(System.out) + System.exit(0) + } + + CommandLineUtils.checkRequiredArgs(parser, options, listGroupOpt, zkConnectOpt) + + val zkConnect = options.valueOf(zkConnectOpt) + var zkClient: ZkClient = null + var channel: BlockingChannel = null + + try { + zkClient = new ZkClient(zkConnect, 30000, 30000, ZKStringSerializer) + if(options.has(listGroupOpt)) { + val list = getConsumerGroups(zkClient).toList + printConsumerGroups(list) + } + } + catch { + case ex : Throwable => + println("Exiting due to: %s.".format(ex.getMessage)) + } + } + + def getConsumerGroups(zkClient: ZkClient) = { + zkClient.getChildren("/consumers") + } + + def printConsumerGroups(list: List[String]) = { + println("CONSUMER GROUPS") + println("*****************") + list.foreach(println) + } +} -- 1.8.5.2 (Apple Git-48)