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

如何快速批量反序列化JSON文件中的dataclass提升处理速度

问题描述

你定义了一组用于后续数据处理的dataclass结构,通过from_dict方法将序列化后的JSON数据还原为类实例,类结构定义如下:

class Currency:
    USD = "USD"
    EUR = "EUR"

@dataclass
class CurrencyPosition:
    currency: Currency
    balance: float

@dataclass
class StockPosition:
    ticker: str
    name: str
    balance: int

@dataclass
class Portfolio:
    currencies: List[CurrencyPosition]
    stocks: List[StockPosition]
   
    def to_dict(self) -> dict: 
         return {
             "currencies": [x.to_json() for x in self.currencies]
             "stocks": [x.to_json() for x in self.stocks]
         }

    @classmethod
    def from_dict(cls, data) -> "Portfolio":
        return cls(
            currencies=[CurrencyPosition.from_dict(x) for x in data["currencies"],
            stocks=[StockPositon.from_dict(x) for x in data["stocks"]
        )

当前存储投资组合数据的目录共包含200万个JSON文件,目录结构如下:

portfolios
├── portfolio_1.json
├── portfolio_2.json
├── .
├── .
├── .
└── portfolio_2000000.json

你需要尽可能快地将全部200万个JSON文件读取为Portfolio实例列表,出于内存占用考虑优先使用生成器实现。目前已测试基于joblib的多线程、多进程实现与串行实现,三者性能差异极小,处理速度均为30±5个/秒,处理完全部文件需要约18小时,无法满足效率要求。当前实现代码如下:

from joblib import delayed, Parallel
from pathlib import Path
from my_lib import Portfolio

def portfolio_from_path(path: Path) -> Portfolio:
    with open(str(path), "r") as f:
        data = json.load(f)
    return Portfolio.from_dict(data)

def serial(paths: List[Path]) -> List[Portfolio]:
    return [portfolio_from_path(x) for x in paths]

def multi_processing(paths: List[Path]) -> List[Portfolio]:
    with Parallel(n_jobs=-1) as parallel:
        results = parallel(delayed(portfolio_from_path)(x) for x in paths)
    return results

def multi_threading(paths: List[Path]) -> List[Portfolio]:
    # n_jobs from ThreadPoolExecutor min(32, (os.cpu_count() or 1) + 4)
    # 8 cores on my machine
    with Parallel(n_jobs=12, prefer="threads") as parallel:
        results = parallel(delayed(portfolio_from_path)(x) for x in paths)
    return results

if __name__ == "__main__":
    paths = list(Path("./portfolios").glob("*.json"))
    serialized_portfolios = multi_processing(paths)
优化方案
  • 先修复现有代码的显性bug
    你当前的类定义存在多处语法错误,先修复再做性能优化:

    1. Portfolio.to_dict方法里两个字典键值对之间缺少逗号
    2. Portfolio.from_dict方法里两个列表推导式的末尾缺少闭合括号
    3. StockPosition在from_dict调用中拼写错误(写为StockPositon)
    4. 子dataclass(CurrencyPosition、StockPosition)未实现from_dict和to_json方法,当前代码无法直接运行
  • 解决底层IO瓶颈 (这是你当前并行无提速的核心原因之一)
    单目录存放200万个文件本身就会导致极高的文件系统元数据查询开销:

    1. 不要用Path.glob一次性加载所有文件路径到内存,换用os.scandir做惰性遍历,天生支持生成器模式,内存占用极低,遍历速度比Path.glob快2-3倍
    2. 读取文件时直接用二进制模式rb读取字节流,跳过Python默认的文本编码转换步骤,给高性能JSON解析器使用
    3. 如果文件系统支持,开启目录索引(ext4的dir_index、NTFS默认索引);如果可以调整存储结构,将200万个文件按哈希拆分到128/256个子目录中,可将文件打开、遍历的速度提升10倍以上
  • 替换慢的JSON解析与对象构造逻辑
    标准库json的解析性能、手动写循环构造dataclass的纯Python开销是第二大性能瓶颈:

    1. 替换标准库json为orjson,高性能JSON解析库,解析速度是标准库的3-10倍,原生支持bytes输入,无需转码
    2. 去掉手动实现的from_dict逐层构造逻辑,换用msgspec做结构化反序列化,可直接将JSON字节流一次性映射为dataclass实例,跳过中间dict构造、列表遍历的Python层开销,这一步可将对象构造速度提升10-20倍
  • 重构并行逻辑,规避joblib的开销
    你之前用joblib多进程无提速的核心原因是任务粒度过细,200万个单文件任务的进程调度、进程间pickle序列化开销远大于实际处理逻辑的开销:

    1. 直接用标准库multiprocessing.Pool实现多进程,进程数设置为物理CPU核心数即可,不需要开超线程
    2. 用imap_unordered方法做惰性迭代,天生符合生成器的内存要求,边处理边返回结果,不会一次性把所有实例加载到内存
    3. 对文件路径做分块处理,每个进程一次性处理1000-5000个文件,批量回传结果,将进程间通信的次数降低几个数量级
  • 终极优化方案
    如果以上优化后仍不能满足性能要求,直接放弃200万个小文件的存储模式:将所有投资组合数据合并存储为Parquet、MsgPack等批量序列化格式,可彻底消除百万级文件open/close的系统调用开销,整体处理速度可再提升一个数量级。

按以上方案优化后,在普通SSD上处理速度可达到3000~10000个/秒,200万个文件仅需数分钟即可处理完成,且全程内存占用稳定。参考实现如下:

import os
import msgspec
from multiprocessing import Pool, cpu_count
from typing import Iterator

# 用msgspec定义结构化类型,替换原生dataclass,解析性能更高
class CurrencyPosition(msgspec.Struct):
    currency: str
    balance: float

class StockPosition(msgspec.Struct):
    ticker: str
    name: str
    balance: int

class Portfolio(msgspec.Struct):
    currencies: list[CurrencyPosition]
    stocks: list[StockPosition]

def _process_file_chunk(path_chunk: list[str]) -> list[Portfolio]:
    # 每个进程独立初始化解码器,避免重复创建开销
    decoder = msgspec.json.Decoder(Portfolio)
    res = []
    for p in path_chunk:
        with open(p, "rb") as f:
            # 直接读取二进制流,跳过编码转换、中间dict构造步骤
            res.append(decoder.decode(f.read()))
    return res

def _chunk_generator(iterable, chunk_size):
    chunk = []
    for item in iterable:
        chunk.append(item)
        if len(chunk) >= chunk_size:
            yield chunk
            chunk = []
    if chunk:
        yield chunk

def iter_portfolios(root_dir: str, chunk_size: int = 2000, workers: int = None) -> Iterator[Portfolio]:
    if workers is None:
        workers = cpu_count(logical=False)
    
    # 惰性遍历目录下所有json文件,不一次性加载全部路径
    def path_generator():
        for entry in os.scandir(root_dir):
            if entry.is_file() and entry.name.endswith(".json"):
                yield entry.path
    
    with Pool(processes=workers) as pool:
        # 分块处理+惰性返回结果,内存全程稳定
        for chunk_res in pool.imap_unordered(_process_file_chunk, _chunk_generator(path_generator(), chunk_size)):
            yield from chunk_res

if __name__ == "__main__":
    count = 0
    for portfolio in iter_portfolios("./portfolios"):
        # 逐份处理业务逻辑,无需将所有实例加载到内存
        count += 1
    print(f"共处理{count}份投资组合")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 11:03:15