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

如何通过多进程将超大型SAS数据集加载到Python中?

超大型SAS数据集多进程加载方案

原有代码错误原因

  • 多进程拥有独立的内存地址空间,子进程无法直接读写主进程的全局变量dfs,且主进程未初始化dfs变量,直接触发未定义报错
  • 入口判断逻辑错误,正确写法为if __name__ == '__main__':,原有代码的判断永远无法触发主进程执行逻辑
  • pd.read_sas返回的迭代器是单进程串行读取生成的,将迭代器传入pool.map仅会实现处理阶段并行,读取阶段仍是串行,无法实现全链路加速

可行实现方案

方案1:基于pandas原生接口的多进程读取

将分块读取逻辑放到子进程执行,每个子进程读取独立的文件区间,真正实现IO并行,代码如下:

import pandas as pd
from multiprocessing import Pool

SAS_FILE_PATH = "替换为你的SAS文件路径"
ENCODING = "ISO-8859-1"
CHUNKSIZE = 100000
# 进程数根据CPU核心数和磁盘IO能力调整,机械盘建议4-8,SSD建议8-16
PROCESS_NUM = 8

def get_total_rows(path: str) -> int:
    """仅读取元数据获取SAS文件总行数,不加载全量数据"""
    chunk_iter = pd.read_sas(path, chunksize=1)
    total = 0
    for _ in chunk_iter:
        total += 1
    return total

def read_chunk(param: dict) -> pd.DataFrame:
    """子进程读取指定区间的分块数据"""
    return pd.read_sas(
        SAS_FILE_PATH,
        encoding=ENCODING,
        skiprows=param["skip"],
        nrows=param["nrows"]
    )

if __name__ == "__main__":
    total_rows = get_total_rows(SAS_FILE_PATH)
    # 生成每个子进程的读取参数
    chunk_params = []
    for skip in range(0, total_rows, CHUNKSIZE):
        chunk_params.append({
            "skip": skip,
            "nrows": min(CHUNKSIZE, total_rows - skip)
        })
    # 多进程读取分块
    with Pool(PROCESS_NUM) as pool:
        chunk_list = pool.map(read_chunk, chunk_params)
    # 合并为全量数据集,内存不足时可跳过该步骤直接处理分块
    full_df = pd.concat(chunk_list, ignore_index=True)

方案2:基于pyreadstat的高性能读取(推荐)

pyreadstat是专门读取统计软件格式文件的高性能库,读取SAS文件的速度是pandas原生接口的3-5倍,适配多进程更简单,首先执行pip install pyreadstat安装依赖,代码如下:

import pyreadstat
import pandas as pd
from multiprocessing import Pool

SAS_FILE_PATH = "替换为你的SAS文件路径"
ENCODING = "ISO-8859-1"
CHUNKSIZE = 100000
PROCESS_NUM = 8

def read_chunk(offset: int) -> pd.DataFrame:
    chunk, _ = pyreadstat.read_sas7bdat(
        SAS_FILE_PATH,
        row_offset=offset,
        row_limit=CHUNKSIZE,
        encoding=ENCODING
    )
    return chunk

if __name__ == "__main__":
    # 直接读取元数据拿总行数,不需要遍历分块
    _, meta = pyreadstat.read_sas7bdat(SAS_FILE_PATH, metadataonly=True)
    total_rows = meta.number_rows
    offsets = list(range(0, total_rows, CHUNKSIZE))
    
    with Pool(PROCESS_NUM) as pool:
        chunk_list = pool.map(read_chunk, offsets)
    full_df = pd.concat(chunk_list, ignore_index=True)

优化建议

  • 如果内存不足以容纳全量数据,可以在read_chunk函数中新增列筛选、行过滤逻辑,仅返回需要的数据,不需要合并全量数据集
  • 进程数不要超过磁盘IO上限,过多进程会导致磁盘随机IO增加反而降低读取速度
  • 如果读取后不需要全量聚合,可以直接在子进程中完成分块处理逻辑,进一步降低主进程内存消耗

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 08:45:05