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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 18:15:01