【发布时间】:2016-10-02 07:12:34
【问题描述】:
我是 Twisted 的新手,这是我的第一个程序。
我找不到使用 kafka-python 库中的 KafkaConsumer 并使用 treq 触发对 elasticsearch 的发布请求的方法。
我可以将问题分解成小块: 创建一个kafka消费者迭代器并从中读取数据(主题可能很大)
def consumeKafka():
consumer = KafkaConsumer(bootstrap_servers="kafka:9092", auto_offset_reset='earliest')
consumer.subscribe(['kafkapipeline'])
for v in consumer:
v.value
使用 treq 发布到 elasticsearch
def post(self):
d = treq.post('http://es:9200/pro/pr/', self.data)
d.addCallbacks(lambda x: print(x), lambda x: print("error %s " % x))
启动反应堆
from twisted.internet import reactor
reactor.callWhenRunning(consumeKafka)
reactor.run()
知道如何进行这项工作吗?
【问题讨论】:
标签: python asynchronous twisted