通过pykafka接收Kafka消息队列的方法
没有Kafka环境,所以也没有进行验证。感觉今后应该能用到,所以借抄在此,备查。
pykafka使用示例,自动消费最新消息,不重复消费:
#-*coding:utf8*-
frompykafkaimportKafkaClient
host='192.168.200.38'
client=KafkaClient(hosts="%s:9092"%host)
printclient.topics
#生产者
#topicdocu=client.topics['task_pull']
#producer=topicdocu.get_producer()
#foriinrange(4):
#printi
#producer.produce('testmessage'+str(i**2))
#producer.stop()
#消费者
topic=client.topics['task_push']
consumer=topic.get_simple_consumer(consumer_group='test',auto_commit_enable=True,consumer_id='test')
formessageinconsumer:
ifmessageisnotNone:
printmessage.offset,message.value
以上这篇通过pykafka接收Kafka消息队列的方法就是小编分享给大家的全部内容了,希望能给大家一个参考,也希望大家多多支持毛票票。