如何修复Kafka消费者批量写入JSONL文件的循环与批处理逻辑
问题修复:Kafka消费者批量写入JSONL文件的逻辑问题
问题场景与症状
- 需求:将Kafka消息写入JSONL文件,每个文件固定条数(示例为2条);生产者一次性发送3条消息,预期生成2个JSONL文件(分别含2条、1条记录),消费者需在无消息可消费时自动停止
- 现有代码问题:第三条消息丢失,需手动中断程序才会写入剩余记录,无法触发剩余记录写入分支及自动停止逻辑
问题根源
- KafkaConsumer的迭代器是阻塞式的,
for message in self.consumer会一直阻塞等待新消息,永远无法退出循环,导致后续处理剩余记录的代码无法执行 write_to_jsonl方法中未定义now变量,会触发NameError- 未设置消费超时,无法判断是否还有新消息到来
修复后的完整代码
import os import json from datetime import datetime from kafka import KafkaConsumer class WikiConsumer: def __init__(self, topic: str, server: str, group_id: str, output_dir: str) -> None: self.consumer = KafkaConsumer( topic, bootstrap_servers=server, group_id=group_id, consumer_timeout_ms=1000 # 新增:1秒无消息则退出迭代 ) self.running = True self.output_dir = output_dir # 确保输出目录存在 os.makedirs(self.output_dir, exist_ok=True) def consume_messages(self, batch_size: int): records = [] while self.running: has_new_messages = False # 迭代消费消息,超时后自动退出循环 for message in self.consumer: has_new_messages = True if message.value != b'""': try: json_str = message.value.decode("utf-8") # 适配消息格式:处理外层可能的字符串包裹 json_obj = json.loads(json_str) if json_str.startswith("{") else json.loads(json.loads(json_str)) records.append(json_obj) print(f"已收集记录: {len(records)}条") # 达到批量大小则写入文件 if len(records) >= batch_size: self.write_to_jsonl(records) records = [] except Exception as e: print(f"解析消息失败: {e}") continue # 处理剩余未写入的记录 if records: self.write_to_jsonl(records) print("剩余记录已写入") records = [] # 如果本轮未收到新消息,停止运行 if not has_new_messages: self.running = False print("无新消息,消费者停止") def write_to_jsonl(self, records): # 修复:定义now变量 now = datetime.now() timestamp = now.strftime("%Y-%m-%d-%H-%M-%S") filename = f"data_{timestamp}.jsonl" with open( os.path.join(self.output_dir, filename), "w" # 用w模式避免重复追加 ) as f: for record in records: json.dump(record, f) f.write("\n") print(f"文件 {filename} 已写入 {len(records)} 条记录") def run(self) -> None: self.consume_messages(2) if __name__ == "__main__": output_directory = "my_dir/output/" consumer = WikiConsumer( "my_project", "localhost:9092", "project-group", output_directory ) consumer.run()
修复说明
- 新增消费超时:设置
consumer_timeout_ms=1000,当1秒内没有新消息时,消费者迭代器会自动退出,让程序有机会处理剩余记录并判断是否停止 - 调整循环逻辑:新增
has_new_messages标记,本轮未收到新消息则停止运行,实现自动退出 - 修复变量错误:在
write_to_jsonl中用datetime.now()定义now变量,解决NameError - 优化文件写入:将文件打开模式从
a改为w,避免同一批次记录重复写入;同时确保输出目录存在 - 增加异常处理:添加消息解析的异常捕获,避免单条坏消息导致程序崩溃
内容的提问来源于stack exchange,提问作者Omega
相关产品推荐
相关产品推荐

