You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何修复Kafka消费者批量写入JSONL文件的循环与批处理逻辑

问题修复:Kafka消费者批量写入JSONL文件的逻辑问题

问题场景与症状

  • 需求:将Kafka消息写入JSONL文件,每个文件固定条数(示例为2条);生产者一次性发送3条消息,预期生成2个JSONL文件(分别含2条、1条记录),消费者需在无消息可消费时自动停止
  • 现有代码问题:第三条消息丢失,需手动中断程序才会写入剩余记录,无法触发剩余记录写入分支及自动停止逻辑

问题根源

  1. KafkaConsumer的迭代器是阻塞式的,for message in self.consumer会一直阻塞等待新消息,永远无法退出循环,导致后续处理剩余记录的代码无法执行
  2. write_to_jsonl方法中未定义now变量,会触发NameError
  3. 未设置消费超时,无法判断是否还有新消息到来

修复后的完整代码

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.15 14:16:36