使用python操作kafka

Posted 三度

tags:

篇首语:本文由小常识网(cha138.com)小编为大家整理,主要介绍了使用python操作kafka相关的知识,希望对你有一定的参考价值。

使用python操作kafka目前比较常用的库是kafka-python库

安装kafka-python

pip3 install kafka-python

生产者

producer_test.py

from kafka import KafkaProducer

producer = KafkaProducer(bootstrap_servers='192.168.0.121:9092')  # 连接kafka

msg = "Hello World".encode('utf-8')  # 发送内容,必须是bytes类型
producer.send('test', msg)  # 发送的topic为test
producer.close()

执行此程序,它没有输出!这个是正常的

消费者

from kafka import KafkaConsumer

consumer = KafkaConsumer('test', bootstrap_servers=['192.168.0.121:9092'])
for msg in consumer:
    recv = "%s:%d:%d: key=%s value=%s" % (msg.topic, msg.partition, msg.offset, msg.key, msg.value)
    print(recv)

执行此程序,此时会hold住,因为它在等待生产者发送消息!

再次执行生产者,此时会输出:

test:0:9: key=None value=b'Hello World'

以上是关于使用python操作kafka的主要内容,如果未能解决你的问题,请参考以下文章

python 使用 kafka

学习笔记:python3,代码片段(2017)

python操作kafka

kafka-python消息读写操作kafka,python,Windows

MySQL系列:kafka停止命令

Kafka-python 客户端导致的 cpu 使用过高,且无法消费消息的问题