通过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消息队列的方法就是小编分享给大家的全部内容了,希望能给大家一个参考,也希望大家多多支持毛票票。