如何快速批量反序列化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
你当前的类定义存在多处语法错误,先修复再做性能优化:Portfolio.to_dict方法里两个字典键值对之间缺少逗号Portfolio.from_dict方法里两个列表推导式的末尾缺少闭合括号StockPosition在from_dict调用中拼写错误(写为StockPositon)- 子dataclass(
CurrencyPosition、StockPosition)未实现from_dict和to_json方法,当前代码无法直接运行
解决底层IO瓶颈 (这是你当前并行无提速的核心原因之一)
单目录存放200万个文件本身就会导致极高的文件系统元数据查询开销:- 不要用
Path.glob一次性加载所有文件路径到内存,换用os.scandir做惰性遍历,天生支持生成器模式,内存占用极低,遍历速度比Path.glob快2-3倍 - 读取文件时直接用二进制模式
rb读取字节流,跳过Python默认的文本编码转换步骤,给高性能JSON解析器使用 - 如果文件系统支持,开启目录索引(ext4的dir_index、NTFS默认索引);如果可以调整存储结构,将200万个文件按哈希拆分到128/256个子目录中,可将文件打开、遍历的速度提升10倍以上
- 不要用
替换慢的JSON解析与对象构造逻辑
标准库json的解析性能、手动写循环构造dataclass的纯Python开销是第二大性能瓶颈:- 替换标准库json为
orjson,高性能JSON解析库,解析速度是标准库的3-10倍,原生支持bytes输入,无需转码 - 去掉手动实现的
from_dict逐层构造逻辑,换用msgspec做结构化反序列化,可直接将JSON字节流一次性映射为dataclass实例,跳过中间dict构造、列表遍历的Python层开销,这一步可将对象构造速度提升10-20倍
- 替换标准库json为
重构并行逻辑,规避joblib的开销
你之前用joblib多进程无提速的核心原因是任务粒度过细,200万个单文件任务的进程调度、进程间pickle序列化开销远大于实际处理逻辑的开销:- 直接用标准库
multiprocessing.Pool实现多进程,进程数设置为物理CPU核心数即可,不需要开超线程 - 用
imap_unordered方法做惰性迭代,天生符合生成器的内存要求,边处理边返回结果,不会一次性把所有实例加载到内存 - 对文件路径做分块处理,每个进程一次性处理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

