• pykafka的基本使用及统计kafka消息总数


    目录

    1、pykafka的安装及连接kafka

    2、获取topic并给kafka的目标topic里面写入数据:

    3、对写入的消息进行消费:

    4、统计kafka消息总数:

    kafka简单说明:

    作为一个消息队列,主要就是由生产者Producer、消费者Consumer这两种角色进行队列的写入和队列的消费:

    1、pykafka的安装及连接kafka

    1. 安装:
    2. pip install pykafka
    3. #支持python3和python2
    4. 连接本地kafka:
    5. >>> from pykafka import KafkaClient
    6. # kafka在服务器上则进行修改:
    7. >>> client = KafkaClient(hosts="127.0.0.1:9092")

    2、获取topic并给kafka的目标topic里面写入数据:

    1. 获取所有topic:
    2. >>> client.topics
    3. {'my.test': 0x19bc8c0 (name=my.test)>}
    4. >>> topic = client.topics['my.test']
    5. 创建producer对象,并对目标topic进行消息写入,也就是生产者写入数据:
    6. >>> with topic.get_producer() as producer:
    7. ... for i in range(4):
    8. ... producer.produce('test message ' + i ** 2)

    3、对写入的消息进行消费:

    1. 创建消费者对象进行消费消息,可以指定group组也可以不指定,两个相同的组会消费相同的数据:
    2. >>> consumer = topic.get_simple_consumer(consumer_group='testwtgroup',
    3. auto_commit_enable=True)
    4. >>> for message in consumer:
    5. if message is not None:
    6. print message.offset, message.value
    7. 结果:
    8. 0 test message 0
    9. 1 test message 1
    10. 2 test message 4
    11. 3 test message 9
    12. # 消费者另一种方式,会根据指定的groupid进行动态分配,保证相同的组不会消费到相同数据:
    13. >>> balanced_consumer = topic.get_balanced_consumer(
    14. consumer_group='testgroup',
    15. auto_commit_enable=True,
    16. zookeeper_connect='myZkClusterNode1.com:2181,myZkClusterNode2.com:2181/myZkChroot'
    17. )

    4、统计kafka消息总数:

    1. # 统计kafka消息总数,找了一圈都没找到相关实现,最后发现其实就是用最后的偏移量的数据-最初偏移量的数据
    2. # 原理参考kafka的官网,通过命令行实现统计总的消息数据:
    3. # kafka-topics.sh --bootstrap-server {IP:port} --list
    4. # 查看kafka的数据 --time-1 表示要获取指定topic所有分区当前的最大位移,--time-2 表示获取当前最早位移
    5. # 下面用参数time -2可进行替换,[--time-1] - [--time-2] 就是partition里面当前实际的数据(相减),
    6. # kafka-run-class.sh kafka.tools.GetOffsetShell --broker-list {IP:port} --topic {target_topic} --time -1
    7. offsets = topic.earliest_available_offsets()
    8. offsets2 = topic.latest_available_offsets()
    9. # print(offsets2)
    10. # print(offsets[0][0][0])
    11. # print(offsets2[0][0][0])
    12. a = offsets2[0][0][0] - offsets[0][0][0]
    13. b = offsets2[1][0][0] - offsets[1][0][0]
    14. c = offsets2[2][0][0] - offsets[2][0][0]
    15. print(a+b+c)
    16. 结果:
    17. 44064

    pykafka官方地址:

    pykafka · PyPI

  • 相关阅读:
    HTTP 请求方法
    三层架构与web结合图解
    IO 多路复用
    Eureka服务发现深度配置:实例ID与租约续期策略
    D. Bicolored RBS.
    直播实时数仓基于DataLeap开放平台在发布管控场景的业务实践
    (附源码)计算机毕业设计Java巴音学院学生资料管理系统
    学习鸿蒙基础(11)
    set() 函数 | Python
    【Linux】使用 Alist 实现阿里云盘4K播放
  • 原文地址:https://blog.csdn.net/m0_37570494/article/details/127677941