【问题标题】:Python - Exit Kafka queue once all messages have been readPython - 读取所有消息后退出Kafka队列
【发布时间】:2021-03-17 20:57:23
【问题描述】:

我正在尝试使用 Python 读取 Kafka 队列的一些数据,如下代码所示:

from kafka import KafkaConsumer
import sys
import json 
import pandas as pd


bootstrap_servers = [localhost]
topicName = 'topic'
consumer = KafkaConsumer (topicName, group_id = 'topic',bootstrap_servers = bootstrap_servers, auto_offset_reset = 'earliest')

data_list = []
for message in consumer:
    print(message)
    data = json.loads(message.value)
    df = pd.json_normalize(data)
    data_list.append(df)

除非我终止连接,否则这似乎永远在循环中运行。在阅读完所有消息或队列中没有新消息后,有没有办法可以停止/退出此循环?

【问题讨论】:

    标签: python apache-kafka kafka-consumer-api


    【解决方案1】:

    poll 方法应该是您正在寻找的。​​p>

    请注意 max_records 参数默认为 max_poll_records,如果不更改,则为 500 条记录

    【讨论】:

      【解决方案2】:

      消息消费完成后,您可能需要close the consumer

      for message in consumer:
          print(message)
          data = json.loads(message.value)
          df = pd.json_normalize(data)
          data_list.append(df)
      
      consumer.close()
      

      【讨论】:

      • 锁在循环中,所以它永远不会退出循环
      猜你喜欢
      • 2019-09-27
      • 2022-01-04
      • 2012-09-03
      • 2016-04-13
      • 1970-01-01
      • 2021-03-01
      • 1970-01-01
      • 1970-01-01
      • 2014-10-28
      相关资源
      最近更新 更多