【问题标题】:how to check if certain consumer is connected to Kafka 0.9.0.x using java?如何检查某些消费者是否使用 java 连接到 Kafka 0.9.0.x?
【发布时间】:2016-10-18 12:11:17
【问题描述】:

如何在 kafka 上获取已连接的消费者列表? 由于消费者是在代理上连接的,是否有任何 Java 实用程序(如 ZkClient/ZkUtils)来获取 Kafka 0.9.0.x 中的已连接消费者列表?就像我们使用以下实用程序获取经纪人列表一样:

        ZkClient zkClient = new ZkClient(endpoint.getZookeeperConnect(), 60000);

        if(zkClient!=null){
            List<String> brokerIds = zkClient.getChildren(ZkUtils.BrokerIdsPath());
            if(CollectionUtils.isNotEmpty(brokerIds) &&  brokerIds.contains(brokerId)){
                logger.debug("Broker:{{}} is connected to Zookeeper.",brokerId);
                flag = true;    
            }
            else{
                logger.error("ERROR:Broker:{{}} is not connected to Zookeeper.",brokerId);
            }
            zkClient.close();
        }

我正在使用 Kafka 0.9.0.x 以及来自 maven 的以下 java lib:

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka_2.11</artifactId>
    <version>0.9.0.1</version>
</dependency>

更新:

我打开了一个“kafka-console-consumer.bat”并运行了一次,然后穿过了 cmd 提示符。然后继续“zookeeper-shell.bat”和ls /consumers然后显示[console-consumer-6008],但我的程序没有显示消费者。使用 zkClient.getChildren(ZkUtils.ConsumersPath()) 我现在只能查看提到的消费者。

【问题讨论】:

    标签: java apache-kafka apache-zookeeper


    【解决方案1】:

    不确定您究竟需要什么信息,但我做了一个示例程序,它提供的信息与 kafka-consumer-groups.sh --describe 相同。

    要使用此代码,请将此依赖项添加到您的 pom。

    <dependency>
        <groupId>org.apache.kafka</groupId>
        <artifactId>kafka_2.11</artifactId>
        <version>0.9.0.1</version>
    </dependency>
    

    然后:

    import org.apache.kafka.clients.CommonClientConfigs;
    import org.apache.kafka.clients.consumer.ConsumerConfig;
    import org.apache.kafka.clients.consumer.KafkaConsumer;
    import org.apache.kafka.common.TopicPartition;
    import kafka.admin.AdminClient;
    import kafka.coordinator.GroupOverview;
    
    Properties props = new Properties();
    props.put(CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, "kafka-1:9092");
    AdminClient adminClient = AdminClient.create(props);
    
    List<GroupOverview> groups =  scala.collection.JavaConversions.seqAsJavaList(
            adminClient.listAllConsumerGroupsFlattened());
    for (GroupOverview group : groups) {
        String groupId = group.groupId();
    
        Properties consProps = new Properties();
        consProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-1:9092");
        consProps.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        consProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
        consProps.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "30000");
        consProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
        consProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
        KafkaConsumer consumer = new KafkaConsumer(consProps);
    
        List<AdminClient.ConsumerSummary> groupSummaries = scala.collection.JavaConversions.seqAsJavaList(
                adminClient.describeConsumerGroup(groupId));
    
        System.out.println("GROUP, TOPIC, PARTITION, CURRENT OFFSET, LOG END OFFSET, LAG, OWNER");
    
        for (AdminClient.ConsumerSummary summary : groupSummaries) {
            String owner = summary.clientId() + "_" + summary.clientHost();
            List<TopicPartition> topicPartitions = scala.collection.JavaConversions.seqAsJavaList(
                    summary.assignment());
            for (TopicPartition tp : topicPartitions) {
    
                // Get current offset
                long currentOffset = consumer.committed(tp).offset();
    
                // get log end offset
                consumer.assign(Arrays.asList(tp));
                consumer.seekToEnd();
                long logEndOffset = consumer.position(tp);
    
                long lag = logEndOffset - currentOffset;
    
                System.out.println(groupId + ", " + tp.topic() + ", " + tp.partition() + ", " +
                        currentOffset + ", " + logEndOffset + ", " + lag + ", " + owner);
            }
        }
    }
    

    【讨论】:

    • 谢谢,正是我只需要获取跑步消费者列表。这是使用“AdminClient”+“listAllConsumerGroupsFlattened()”方法实现的。事情在kafak中还很隐蔽。
    【解决方案2】:

    对于 0.9.x 新消费者并列出所有活跃的消费者组:

    1. 查找所有代理并向每个代理发送“ListGroups”请求并获取所有组信息;

    详情可以参考$KAFKA_HOME/bin/kafka-consumer-groups.sh(kafka.admin.ConsumerGroupCommand.KafkaConsumerGroupService.list())

    对于0.9.x的新消费者并描述某些消费者群体的详细信息:

    1. 找到消费者组协调器并向其发送“DescribeGroups”请求,并获取所有组成员信息和分区分配信息;
    2. 调用 KafkaConsumer.committed(TopicPartition partition) 以获取给定分区的最后提交偏移量。

    详情可以参考$KAFKA_HOME/bin/kafka-consumer-groups.sh(kafka.admin.ConsumerGroupCommand.KafkaConsumerGroupService.describe())

    请注意,旧消费者和新消费者对此有完全不同的实现。(这两个逻辑都在 kafka.admin.ConsumerGroupCommand 中实现。

    【讨论】:

    • 这个编辑是一个错误必须添加到我自己的问题中,很抱歉它不知道如何丢弃它。
    • 顺便说一句,我在 Windows 上使用 java 代码,Kafka Windows 没有 'kafka-consumer-groups.bat' 现在该怎么办。
    【解决方案3】:

    几乎相同,但您必须检查 ZkUtils.ConsumersPath (= /consumers)。

    Zookeeper 中的消费者结构是下一个 /consumers/[groupId]/ids/[consumerId],因此您可以通过导航获得每个组的组和消费者。

    【讨论】:

    • ZkUtils.ConsumersPath (/consumers) 总是返回 [ ]。我认为消费者群体信息现在保存在 kafka 上。我已经通过这部分来验证消费者列表。
    • 在 0.9.x 和 0.10.x 中仍然保留消费者组和消费者。你可以在代码中查看 ZkUtils.getConsumers 获取 ConsumersPath 的孩子。 github.com/apache/kafka/blob/trunk/core/src/main/scala/kafka/…
    • zkClient.getChildren(ZkUtils.ConsumersPath()) 返回空 [ ]。
    • 试试shell,打开一个控制台生产者和控制台消费者。然后产生一条消息。最后使用工具 zookeeper shell ls /consumers = [console-consumer-44669] 检查是否正常,一切正常,错误将出现在未注册的消费者身上。我刚试过,我在 zk /consumers 节点看到了消费者。
    • 我打开了一个“kafka-console-consumer.bat”并运行了一次,然后穿过了 cmd 提示符。然后继续“zookeeper-shell.bat”和 ls /consumers 然后显示 [console-consumer-6008],但我的编程消费者没有显示。使用 'zkClient.getChildren(ZkUtils.ConsumersPath())' 我现在只能查看提到的消费者。
    猜你喜欢
    • 1970-01-01
    • 2016-07-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-03-14
    • 2017-01-31
    相关资源
    最近更新 更多