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

Confluent Kafka Consumer批量消费剩余消息处理方案咨询

Kafka批量消费完整处理方案

问题核心

原代码无法处理最后一批不足batch_size的消息,且存在文件覆盖风险,原因如下:

  1. 内层循环持续阻塞,仅当消息数达标时才触发处理,未考虑“无新消息但有剩余数据”的场景。
  2. 所有批次共用同一日期命名的JSON文件,后续批次会覆盖前序数据。

修改方案

  1. 增加超时处理逻辑:当连续一段时间没有新消息时,自动处理剩余未达标的消息。
  2. 生成唯一批次文件:为每个批次的JSON文件加入时间戳和批次ID,避免数据覆盖。
  3. 显式提交偏移量:处理完每个批次后提交偏移量,保证消费的幂等性。

修改后的完整代码

#!/usr/bin/env python

import os
from argparse import ArgumentParser, FileType
from configparser import ConfigParser
from confluent_kafka import Consumer, OFFSET_BEGINNING
import pandas as pd
import datetime
from datetime import timedelta
import json
import time


def parse_args():
    """解析命令行参数"""
    parser = ArgumentParser()
    parser.add_argument("config_file", type=FileType("r"))
    parser.add_argument("--reset", action="store_true")
    parser.add_argument("topic_name", help="Kafka Topic名称")
    return parser.parse_args()


def parse_config(args):
    """解析配置文件"""
    config_parser = ConfigParser()
    config_parser.read_file(args.config_file)
    config = dict(config_parser["default"])
    config.update(config_parser["consumer"])
    return config


def create_consumer(config):
    """创建并返回Consumer实例"""
    consumer = Consumer(config)
    return consumer


def reset_offset(consumer, partitions, reset):
    """根据reset标志设置消息偏移量"""
    if reset:
        for p in partitions:
            p.offset = OFFSET_BEGINNING
        consumer.assign(partitions)


def get_file_name(topic, batch_id=None):
    """生成唯一的文件名,包含批次标识"""
    folder_name = "trimet_raw_data"
    if not os.path.exists(folder_name):
        os.makedirs(folder_name)
    timestamp = datetime.datetime.now().strftime('%Y-%m-%d_%H%M%S')
    if batch_id:
        return os.path.join(
            folder_name, f"{topic}_{timestamp}_batch{batch_id}.json"
        )
    return os.path.join(
        folder_name, f"{topic}_{timestamp}.json"
    )


def write_message_to_file(f, topic, key, value):
    """将消息写入文件"""
    if key is None:
        key = ""
    else:
        key = key.decode("utf-8")
    if value is not None:
        value = value.decode("utf-8")
    f.write(f"{value}\n")

        
def consume_messages(consumer, topic, reset, batch_size, idle_timeout=10):
    """批量消费Kafka消息,处理所有批次包括最后不足批量的消息"""
    batch_id = 0
    while True:
        try:
            message_count = 0
            data_list = []
            batch_id += 1
            file_name = get_file_name(topic, batch_id)
            f = open(file_name, "w")
            idle_count = 0
            while True:
                msg = consumer.poll(1.0)
                if msg is None:
                    idle_count += 1
                    print(f"无新消息,等待中... ({idle_count}/{idle_timeout})")
                    if idle_count >= idle_timeout:
                        # 超时无新消息,处理剩余数据
                        break
                    time.sleep(1)
                    continue
                idle_count = 0  # 收到新消息重置超时计数器
                
                if msg.error():
                    print(f"ERROR: {msg.error()}")
                    continue
                
                key = msg.key()
                value = msg.value()
                try:
                    data_dict = json.loads(value)
                    data_list.append(data_dict)
                    message_count += 1
                    if message_count % 10000 == 0:
                        print(f"已处理 {message_count} 条消息")
                except ValueError:
                    print("错误:消息格式不是JSON")
                    continue
                
                write_message_to_file(f, topic, key, value)
                
                if message_count % batch_size == 0:
                    print(f"已处理 {message_count} 条消息,开始处理当前批次...")
                    f.close()
                    read_raw_data(file_name)
                    # 提交偏移量
                    consumer.commit(msg)
                    message_count = 0
                    data_list = []
                    # 为下一批次创建新文件
                    batch_id += 1
                    file_name = get_file_name(topic, batch_id)
                    f = open(file_name, "w")
            
            # 处理最后一批不足batch_size的消息
            if message_count > 0:
                print(f"处理最后一批 {message_count} 条消息...")
                f.close()
                read_raw_data(file_name)
                # 提交最后一批的偏移量
                consumer.commit()
            
            print("当前所有消息已处理完成,等待新消息...")
            
        except KeyboardInterrupt:
            print("关闭Consumer")
            print(f"累计处理消息数: {message_count}")
            consumer.close()
            break
        except Exception as e:
            print(f"发生错误: {str(e)}")
            continue


def data_transform(df):
    """数据转换逻辑"""
    if df["OPD_DATE"].str.contains(r"\d{2}-[A-Za-z]{3}-\d{2}").any():
        filtered_df = df.copy()
        filtered_df.rename(
            columns={
                "EVENT_NO_TRIP": "trip_id",
                "OPD_DATE": "tstamp",
                "VELOCITY": "longitude",
                "DIRECTION": "latitude",
                "RADIO_QUALITY": "gps_satellites",
                "GPS_LONGITUDE": "gps_hdop",
            },
            inplace=True,
        )
        filtered_df.columns = filtered_df.columns.str.lower()
    else:
        filtered_df = df.copy()
        filtered_df.rename(
            columns={
                "EVENT_NO_TRIP": "trip_id",
                "OPD_DATE": "tstamp",
                "GPS_LONGITUDE": "longitude",
                "GPS_LATITUDE": "latitude",
            },
            inplace=True,
        )
        filtered_df.columns = filtered_df.columns.str.lower()

    filtered_df["tstamp"] = filtered_df["tstamp"].apply(
        lambda value: pd.to_datetime(value, format="%d-%b-%y", errors="coerce")
        if len(value) <= 11
        else pd.to_datetime(value, format="%d%b%Y:%H:%M:%S", errors="coerce")
    )
    filtered_df["act_time"] = pd.to_numeric(filtered_df["act_time"], errors="coerce")
    filtered_df["tstamp"] = filtered_df.apply(
        lambda row: row["tstamp"] + timedelta(seconds=row["act_time"])
        if pd.notnull(row["tstamp"])
        else "",
        axis=1,
    )
    filtered_df = filtered_df.sort_values(["trip_id", "tstamp"])
    filtered_df["dmeters"] = filtered_df.groupby(["trip_id"])["meters"].diff()
    filtered_df["dtimestamp"] = filtered_df.groupby(["trip_id"])["tstamp"].diff()
    filtered_df["speed"] = filtered_df.apply(
        lambda row: round(row["dmeters"] / row["dtimestamp"].total_seconds(), 2)
        if row["dtimestamp"].total_seconds() != 0
        else 0,
        axis=1,
    )

    filtered_df["speed"] = filtered_df.groupby(["trip_id"])["speed"].fillna(
        method="bfill"
    )
    filtered_df["service_key"] = filtered_df["tstamp"].dt.dayofweek.apply(
        lambda day: "Weekday" if day < 5 else ("Saturday" if day == 5 else "Sunday")
    )
    return filtered_df


def read_raw_data(file_path):
    """读取原始JSON文件,转换后写入CSV"""
    csv_filename = "test_csv.csv"
    
    with open(file_path, "r") as f:
        df = pd.read_json(f, lines=True)
    transformed_df = data_transform(df)
    
    if not os.path.isfile(csv_filename):
        print("创建CSV文件")
        transformed_df.to_csv(csv_filename, index=False)
    else:
        print("追加到CSV文件")
        transformed_df.to_csv(csv_filename, mode='a', index=False, header=False)


def main():
    """主函数"""
    args = parse_args()
    config = parse_config(args)
    consumer = create_consumer(config)
    topic = args.topic_name
    consumer.subscribe([topic])
    # 批量大小和空闲超时时间(秒)
    batch_size = 100000
    consume_messages(consumer, topic, args.reset, batch_size)


if __name__ == "__main__":
    main()

关键修改说明

  1. 文件名生成:get_file_name函数增加了batch_id参数,每个批次生成唯一的JSON文件,彻底避免数据覆盖问题。
  2. 超时处理:新增idle_timeout参数,当连续指定时间无新消息时,自动跳出循环处理剩余数据,保证最后一批消息被处理。
  3. 批次处理逻辑:每次达到batch_size或超时无新消息时,都会触发read_raw_data处理当前批次,并提交偏移量,确保消息不会重复消费。
  4. 偏移量提交:显式提交每个批次的偏移量,配合唯一文件名,实现消费的幂等性。

内容的提问来源于stack exchange,提问作者Alice

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 02:52:37