【问题标题】:Is there a Python API for event-driven Kafka consumer?是否有用于事件驱动的 Kafka 消费者的 Python API?
【发布时间】:2019-01-25 21:01:08
【问题描述】:

我一直在尝试构建一个以 Kafka 作为唯一界面的 Flask 应用程序。出于这个原因,我希望有一个 Kafka 消费者,当相关主题的流中有新消息时触发,并通过将消息推送回 Kafka 流来响应。

我一直在寻找类似 Spring 的实现:

@KafkaListener(topics = "mytopic", groupId = "mygroup")
public void listen(String message) {
    System.out.println("Received Messasge in group mygroup: " + message);
}

我看过:

  1. kafka-python
  2. pykafka
  3. confluent-kafka

但我在 Python 中找不到与事件驱动的实现风格相关的任何内容。

【问题讨论】:

  • 消费者已经是“事件驱动的”。你还需要什么?
  • 如果你有Flask,那么HTTP也是一个接口,不仅仅是Kafka

标签: python events flask apache-kafka listener


【解决方案1】:

这是@MickaelMaison 的answer 给出的想法的实现。我用kafka-python

from kafka import KafkaConsumer
import threading

BOOTSTRAP_SERVERS = ['localhost:9092']

def register_kafka_listener(topic, listener):
# Poll kafka
    def poll():
        # Initialize consumer Instance
        consumer = KafkaConsumer(topic, bootstrap_servers=BOOTSTRAP_SERVERS)

        print("About to start polling for topic:", topic)
        consumer.poll(timeout_ms=6000)
        print("Started Polling for topic:", topic)
        for msg in consumer:
            print("Entered the loop\nKey: ",msg.key," Value:", msg.value)
            kafka_listener(msg)
    print("About to register listener to topic:", topic)
    t1 = threading.Thread(target=poll)
    t1.start()
    print("started a background thread")

def kafka_listener(data):
    print("Image Ratings:\n", data.value.decode("utf-8"))

register_kafka_listener('topic1', kafka_listener)

轮询在不同的线程中完成。收到消息后,通过传递从 Kafka 检索到的数据来调用侦听器。

【讨论】:

    【解决方案2】:

    Kafka Consumer 必须不断轮询以从代理检索数据。

    Spring 为您提供了这个花哨的 API,但在幕后,它在循环中调用 poll 并且仅在检索到记录时调用您的方法。

    您可以使用您提到的任何 Python 客户端轻松构建类似的东西。就像在 Java 中一样,这不是(大多数)Kafka 客户端直接公开的 API,而是由顶层提供的东西。这是需要构建的东西。

    【讨论】:

    • 在这种情况下使用套接字有意义吗?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-02-25
    • 2016-01-23
    • 1970-01-01
    • 2020-12-30
    • 2016-10-31
    • 1970-01-01
    相关资源
    最近更新 更多