【发布时间】:2022-10-21 21:42:05
【问题描述】:
我有一个系统,其中一台机器正在读取数据并不断附加一个 .txt 文件。该数据集通过 Kafka Connect 读入 Kafka 代理,然后使用一些 Python 代码进行预处理。机器大约每 5 分钟运行一次,因此我们预计数据会进入,然后空闲 5 分钟。直到下一批。 Kafka 设置很好,因此请假设此代码上游的所有内容都正常工作。
from confluent_kafka import Consumer
import json
KAFKA_BROKER_URL = 'localhost:9092'
live_data = []
def parse_poll_message(msg):
row = json.loads(msg)
split_msg = list(row['payload'].split('\t'))
return split_msg
consumer = Consumer({
'bootstrap.servers': KAFKA_BROKER_URL,
'group.id': 'mygroup',
'auto.offset.reset': 'earliest',
'enable.auto.commit': True
})
consumer.subscribe(['my_topic'])
while 1:
msg = consumer.poll()
if msg is None:
break
elif msg.error():
print("Consumer error: {}".format(msg.error()))
continue
else:
live_data.append(parse_poll_message(msg.value().decode('utf-8')))
consumer.close()
上面的代码只是演示了我在某个时间点会做什么。我想做的是每5分钟,收集当时所有的消息,转换成dataframe,进行一些计算,然后等待下一组消息。如何在正确的时间间隔内保留消息的同时保持此循环处于活动状态?
任何和所有的建议表示赞赏。谢谢!
【问题讨论】:
-
我的建议是使用(持续运行的)pyspark、Flink 或 Beam 作业,因为它们实际上支持这种带有水印的翻滚窗口函数来创建数据帧......否则,不清楚自从你在每个时间间隔之后是什么关闭了你的消费者存在消息时有一个无限循环(例如,假设您有一个非常大的延迟,需要超过 5 分钟才能阅读)
标签: python apache-kafka confluent-kafka-python