Python Multiprocessing实现多文件实时数据处理的技术缺陷排查
现有方案技术缺陷排查
核心效率缺陷(你提到的批量等待问题)
你遇到的批次等待问题根源是使用了阻塞的pool.map()方法:该方法必须等传入的所有任务全部处理完成才会返回,就算大部分进程早已空闲,也不会提前拉取下一批文件处理。同时你每次仅取10个文件提交给20进程的池子,天然就有一半进程闲置,再叠加文件大小差异大的场景,整体资源利用率非常低。
代码逻辑硬错误
现有代码直接运行会直接报错,存在多处低级错误:
- 依赖未导入:用到的
glob、gzip库,以及date_of自定义函数都没有声明/导入,运行直接抛出未定义异常 glob用法错误:glob.glob(Receive_Dir)只能拿到文件夹本身的路径,无法匹配目录下的文件,需要改为glob.glob(os.path.join(Receive_Dir, "*"))才能正常扫描待处理文件- 变量未定义:
parser函数接收的参数是infile,但代码里直接用了未声明的filename变量,运行直接报错 - 路径拼接错误:直接用字符串拼接路径,既跨平台兼容差,还少写了文件名和csv后缀之间的点,输出文件会变成
xxxcsv.gz而非xxx.csv.gz - 目录创建不安全:多进程同时判断目录不存在、同时发起创建请求时会抛出异常,应该用
os.makedirs(out_dir, exist_ok=True)保证线程安全 - 资源泄漏风险:
fout没有用上下文管理器包裹,如果处理过程中抛出异常,fout.close()不会执行,会导致文件句柄泄漏 - 无效代码:
pool.close()写在死循环后面,永远不会被执行
功能设计缺陷
- 没有写入状态判断:如果文件还在写入过程中就被扫描到,会读到不完整的数据,甚至处理失败
- 没有错误兜底:文件处理失败时没有异常捕获,会直接删除原文件,导致数据永久丢失
- 空转占用CPU:死循环没有休眠逻辑,扫不到文件时会持续空转,浪费CPU资源
优化方案(解决批量等待问题)
将阻塞的pool.map改为异步提交任务的apply_async,无需等待整批处理完成,有空闲进程就可以处理新扫描到的文件,参考修改后代码:
from multiprocessing import Pool import os import glob import gzip import time from datetime import datetime # 按文件修改时间生成输出目录,可根据需求调整 def date_of(filepath): mtime = os.path.getmtime(filepath) return datetime.fromtimestamp(mtime).strftime("%Y%m%d") PARSE_OUT = "/opt/out/" RECEIVE_DIR = "/opt/receive/" def parser(infile): try: # 校验文件是否写入完成:间隔1秒校验文件大小无变化 size1 = os.path.getsize(infile) time.sleep(1) size2 = os.path.getsize(infile) if size1 != size2: return False, f"跳过未写完的文件:{infile}" filename = os.path.basename(infile) date_str = date_of(infile) out_dir = os.path.join(PARSE_OUT, date_str) os.makedirs(out_dir, exist_ok=True) out_path = os.path.join(out_dir, f"{filename}.csv.gz") with gzip.open(out_path, 'wb') as fout, gzip.open(infile, 'rb') as fin: for line in fin: # 二进制模式下按字节切割,避免编码问题 data = line.split(b',') fout.write(b','.join(data)) os.remove(infile) return True, f"处理成功:{infile}" except Exception as e: return False, f"处理失败[{infile}]:{str(e)}" if __name__ == '__main__': pool = Pool(20) # 记录已经提交过的文件,避免重复处理 submitted_files = set() while True: # 扫描目录下所有未提交的文件 all_files = glob.glob(os.path.join(RECEIVE_DIR, "*")) target_files = [f for f in all_files if os.path.isfile(f) and f not in submitted_files] if not target_files: # 无文件时休眠2秒,降低CPU占用 time.sleep(2) continue # 异步提交所有待处理文件,不需要等上一批完成 for f in target_files: submitted_files.add(f) pool.apply_async(parser, args=(f,), callback=lambda res: print(res[1])) # 定期清理已处理完成的文件记录,避免集合过大(可根据自己的需求优化清理逻辑) time.sleep(0.5)
内容的提问来源于stack exchange,提问作者notilas
相关产品推荐
相关产品推荐

