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

多进程异步处理数据,需保证Parquet文件顺序写入的问题求助

问题描述

处理300GB级超大日志文件时,现有脚本通过分块读取文件,利用multiprocessing进程池异步解析日志键值对并写入Parquet文件,但存在处理后的块乱序写入的问题。需要实现异步处理数据的同时,保证按文件读取顺序写入Parquet,同时寻求更高效的实现方案。

现有代码如下:

def process_line(line: bytes):
    """Parse values with regex from the line."""
    data ....
    return data


def process_and_save(chunk: list[bytes], lock: multiprocessing.Lock):
    """Process chunk of lines and save the result to parquet file."""
    result = [process_line(line) for line in chunk]
    lock.acquire()
    try:
        write(.....)
    finally:
        lock.release()
    return


def main():
    """Process log file, parse data, save to parquet file."""
    # create a pool of processes
    mamanager = multiprocessing.Manager()
    # to asure saving to parquet happens only one at a time
    lock = mamanager.Lock()
    pool = multiprocessing.Pool(processes=NUM_PROCESSES)
    # delete file if exists
    if os.path.exists(PARQUET_FILE_PATH):
        os.remove(PARQUET_FILE_PATH)
    with lzma.open(LOG_FILE_PATH, 'rb') as file:
        while True:
            # readlines wont cut the rows in half!
            chunk = file.readlines(BYTES_PER_CHUNK)
            if not chunk:
                break
            pool.apply_async(process_and_save, (chunk, lock))

    # close the pool of processes
    pool.close()
    pool.join()

if __name__ == '__main__':
    main()
解决方案

方案1:基于块序号的有序写入

核心思路是给每个读取的块分配唯一序号,异步处理完成后,仅当当前块是待写入的下一个序号时才执行写入操作,确保顺序一致。同时替换开销较高的multiprocessing.Manager锁,改用进程内锁降低通信成本。

修改后代码

import multiprocessing
import os
import lzma
from typing import List, Dict

def process_line(line: bytes):
    """Parse values with regex from the line."""
    # 保留原解析逻辑
    data = {}  # 替换为实际解析结果
    return data

def process_chunk(chunk: List[bytes], chunk_idx: int):
    """仅处理块,不直接写入,返回处理结果和块序号"""
    result = [process_line(line) for line in chunk]
    return chunk_idx, result

def main():
    PARQUET_FILE_PATH = "output.parquet"
    LOG_FILE_PATH = "large_log.xz"
    NUM_PROCESSES = multiprocessing.cpu_count()
    BYTES_PER_CHUNK = 1024 * 1024 * 64  # 64MB块大小,可根据内存调整

    # 初始化状态:已处理的块、当前待写入的序号、写入锁
    processed_chunks: Dict[int, List] = {}
    current_write_idx = 0
    write_lock = multiprocessing.Lock()

    def write_callback(result):
        """异步处理完成后的回调函数,负责有序写入"""
        nonlocal current_write_idx
        chunk_idx, parsed_data = result

        # 先把处理结果存入字典
        with write_lock:
            processed_chunks[chunk_idx] = parsed_data

        # 检查是否可以连续写入多个块
        while True:
            with write_lock:
                if current_write_idx in processed_chunks:
                    data_to_write = processed_chunks.pop(current_write_idx)
                    current_write_idx += 1
                else:
                    break

            # 写入Parquet(替换为实际写入逻辑,支持append模式)
            write_mode = 'w' if current_write_idx == 1 else 'a'
            write(data_to_write, PARQUET_FILE_PATH, mode=write_mode)

    # 清理已有文件
    if os.path.exists(PARQUET_FILE_PATH):
        os.remove(PARQUET_FILE_PATH)

    # 创建进程池并提交任务
    with multiprocessing.Pool(processes=NUM_PROCESSES) as pool:
        chunk_idx = 0
        with lzma.open(LOG_FILE_PATH, 'rb') as file:
            while True:
                chunk = file.readlines(BYTES_PER_CHUNK)
                if not chunk:
                    break
                pool.apply_async(process_chunk, args=(chunk, chunk_idx), callback=write_callback)
                chunk_idx += 1

        pool.close()
        pool.join()

if __name__ == '__main__':
    main()

关键改进点

  • 拆分处理与写入逻辑:process_chunk仅负责解析,写入逻辑由回调函数统一控制
  • 用块序号维护顺序:通过current_write_idx跟踪下一个需要写入的块,确保按读取顺序输出
  • 替换Manager.Lock为进程内Lock:减少跨进程通信开销
  • 支持连续写入:当多个块提前处理完成时,可连续写入,避免等待

方案2:更高效的大文件处理方案

针对300GB级文件,可结合以下优化手段进一步提升性能:

  • 使用Dask替代原生multiprocessing:Dask天然支持并行数据处理和有序输出,自动处理分块、并行计算和结果合并,无需手动维护块序号
    import dask.bag as db
    import pandas as pd
    
    def process_line(line: bytes):
        # 原解析逻辑
        return data
    
    # 读取压缩日志文件,指定块大小
    b = db.read_text(LOG_FILE_PATH, compression='xz', blocksize='64MB')
    # 并行解析
    parsed = b.map(process_line)
    # 转换为DataFrame并写入Parquet(自动按顺序合并)
    parsed.to_dataframe().to_parquet(PARQUET_FILE_PATH, write_index=False)
    
  • 调整块大小:根据内存和CPU核心数调整BYTES_PER_CHUNK,避免块过小导致进程切换频繁,或块过大导致内存压力
  • 预编译正则表达式:如果process_line中使用正则,提前编译正则对象,避免重复编译开销
  • 使用更快的Parquet写入库:比如pyarrow替代pandas默认引擎,提升写入速度

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 20:55:33