【问题标题】:How to see (new-) consumers connected to a specific Kafka topic如何查看连接到特定 Kafka 主题的(新)消费者
【发布时间】:2017-12-21 08:18:20
【问题描述】:

我有一个新消费者(python 消费者)列表。我可以使用以下命令检索组:

bin/kafka-consumer-groups.sh --new-consumer --bootstrap-server localhost:9092 --list

我可以为每个人获取他们所连接的主题

bin/kafka-consumer-groups.sh --new-consumer --bootstrap-server localhost:9092  --describe --group TheFoundGroupId
  1. 如何获取与某个主题相关的所有组(最好是所有消费者,即使不在一个组中)?
  2. 除了将其作为 shell 命令运行之外,还有其他方法可以从 python 访问它吗?

【问题讨论】:

    标签: python apache-kafka consumer


    【解决方案1】:

    感谢您提出这个问题。

    所有消费者配置,如消费者组ID,哪个消费者订阅了哪个主题存储在zookeeper中。

    运行以下命令以连接到 Zookeeper

    ./bin/zookeeper-shell localhost:2181

    然后运行

    ls /消费者

    您将获得所有存在的消费者群体。如果你不提供消费组。 Kafka 会随机分配消费者组。对于控制台消费者,它将分配 console-consumer-XXXXX id

    你可以从python sn-p下面获取所有的消费者组

    安装zookeeper python client

    from kazoo.client import KazooClient
    
    zk = KazooClient(hosts='127.0.0.1:2181')
    zk.start()
    
    
    # get all consumer groups
    consumer_groups = zk.get_children("/consumers")
    print("There are %s consumer group(s) with names %s" % (len(consumer_groups), consumer_groups))
    
    # get all consumers in group
    for consumer_group in consumer_groups:
        consumers = zk.get_children("/consumers/"+consumer_group)
        print("There are %s consumers in %s consumer group. consumer are : %s" % (len(consumers), consumer_group, consumers))
    

    获取连接到某个主题的消费者或消费者组。

    获取/consumers/consumergroup_id/ids/consumer_id/

    会给你像

    这样的输出
    {"version":1,"subscription":{"test":1},"pattern":"white_list","timestamp":"1514218381246"}
    

    订阅对象下消费者订阅的所有主题。根据您的用例实现逻辑

    谢谢

    【讨论】:

    • 不幸的是,这不适用于“新风格”消费者。该列表为空。
    • 哪个列表是空的?消费者还是消费群体? @colin-SBI
    【解决方案2】:

    这不是最好的解决方案,但由于似乎没有人有答案,这就是我最终解决它的方法(为消费者分配一个组并替换 ___YOURGROUP____):

        import subprocess
        import os
        if "KAFKA_HOME" in os.environ:
            kafkapath = os.environ["KAFKA_HOME"]
        else:
            kafkapath = oms_cfg.kafka_home
            # error("Please set up $KAFKA_HOME environment variable")
            # exit(-1)
    
        instances = []
        # cmd = kafkapath + '/bin/kafka-consumer-groups.sh --new-consumer --bootstrap-server {} --list'.format(oms_cfg.bootstrap_servers)
        # result = subprocess.run(cmd.split(' '), stdout=subprocess.PIPE, stderr=subprocess.PIPE)
        igr = ____YOURGROUP_____ # or run over all groups from the commented out command
        print("Checking topics of consumer group {}".format(igr))
        topic_cmd = kafkapath + '/bin/kafka-consumer-groups.sh --new-consumer --bootstrap-server ' + oms_cfg.bootstrap_servers + ' --describe --group {gr}'
        result = subprocess.run(topic_cmd.format(gr=igr).split(' '), stdout=subprocess.PIPE, stderr=subprocess.PIPE)
        table = result.stdout.split(b'\n')
        # You could add a loop over topic here
        for iline in table[1:]:
            iline = iline.split()
            if not len(iline):
                continue
            topic = iline[0]
            # we could check here for the topic. multiple consumers in same group -> only one will connect to each topic
            # if topic != oms_cfg.topic_in:
            #     continue
            client = iline[-1]
            instances.append(tuple([client, topic]))
            #    print("Client {} Topic {} is fine".format(client, topic))
        if len(instances):
            error("Cannot start. There are currently {} instances running. Client/topic {}".format(len(instances),
                                                                                                      instances))
            exit(-1)
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-12-21
      • 2020-04-02
      • 2019-09-20
      • 1970-01-01
      • 2018-06-26
      • 2017-10-17
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多