【问题标题】:Twisted python to read from kafka and write to elasticsearch扭曲的python从kafka读取并写入elasticsearch
【发布时间】: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


    【解决方案1】:

    我根本不使用 Kafka,所以我不确定这是否适合你。另外,我假设您无法同时运行 Kafka 和 treq。我在 Twisted 中处理迭代器的一种通用方法是使用 inlineCallbacks 等待结果,然后对结果进行处理。

    from twisted.internet import defer
    
    @defer.inlineCallbacks
    def consumeKafka():
        consumer = KafkaConsumer(bootstrap_servers="kafka:9092", auto_offset_reset='earliest')
        consumer.subscribe(['kafkapipeline'])
        for v in consumer:
            value = yield v.value
            # do stuff with value
    

    然后你可以简单地调用这个函数,reactor 会处理剩下的事情。所以你的主要部分看起来像这样:

    consumeKafka()
    reactor.run()
    

    请注意,consumeKafka() 函数返回一个 Deferred,因此您可以根据需要添加回调和 errbacks。熟悉此模型后,请查看 Cooperator 对象以了解更多功能。

    【讨论】:

    • 感谢您的回答,我仍在努力使其工作,但我认为我需要先阅读扭曲的文档。这并不容易。
    猜你喜欢
    • 1970-01-01
    • 2018-10-24
    • 2018-01-31
    • 2023-03-07
    • 2018-05-07
    • 1970-01-01
    • 1970-01-01
    • 2023-03-14
    • 2021-04-28
    相关资源
    最近更新 更多