Confluent Kafka Consumer批量消费剩余消息处理方案咨询
Kafka批量消费完整处理方案
问题核心
原代码无法处理最后一批不足batch_size的消息,且存在文件覆盖风险,原因如下:
- 内层循环持续阻塞,仅当消息数达标时才触发处理,未考虑“无新消息但有剩余数据”的场景。
- 所有批次共用同一日期命名的JSON文件,后续批次会覆盖前序数据。
修改方案
- 增加超时处理逻辑:当连续一段时间没有新消息时,自动处理剩余未达标的消息。
- 生成唯一批次文件:为每个批次的JSON文件加入时间戳和批次ID,避免数据覆盖。
- 显式提交偏移量:处理完每个批次后提交偏移量,保证消费的幂等性。
修改后的完整代码
#!/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()
关键修改说明
- 文件名生成:
get_file_name函数增加了batch_id参数,每个批次生成唯一的JSON文件,彻底避免数据覆盖问题。 - 超时处理:新增
idle_timeout参数,当连续指定时间无新消息时,自动跳出循环处理剩余数据,保证最后一批消息被处理。 - 批次处理逻辑:每次达到
batch_size或超时无新消息时,都会触发read_raw_data处理当前批次,并提交偏移量,确保消息不会重复消费。 - 偏移量提交:显式提交每个批次的偏移量,配合唯一文件名,实现消费的幂等性。
内容的提问来源于stack exchange,提问作者Alice
相关产品推荐
相关产品推荐

