Pandas/Dask DataFrame内存优化 本地模拟扩容方案咨询
问题背景
搭建了一套针对患者开展实验的模拟流水线,全部基于Python技术栈(SciPy、NumPy、Pandas等)实现,运行在2017款MacBook Pro(配置为3.5 GHz 双核Intel Core i7处理器、16 GB 2133 MHz LPDDR3内存)的本地Docker容器中。
模拟流程最终需生成单观测维度的统计结果,前置步骤为基于compound字段对两个核心DataFrame做内连接(inner merge)。当前测试数据规模为patients表50000行、experiments表15000行,当尝试将实验规模扩大10倍(即experiments表增至150000行,patients表规模不变)时,模拟程序触发137内存溢出错误(已将Mac版Docker Desktop的内存配额调整至硬件上限)。要求在不采用Spark等分布式工作流、仅本地运行的前提下,给出可行的性能与扩展性优化方案。
核心DataFrame结构
1. experiments(实验表)
结构示例如下:
trial experiment observation compound value 0 1 10 1 4 7.578612 1 7 8 1 4 1.751288 2 11 7 1 4 0.754702 3 30 6 1 4 7.762336 4 35 4 1 3 3.613458
2. patients(患者表)
结构示例如下:
patient compound lab unit metric_1 metric_2 metric_3 metric_4 0 72070 7 lab_a unit_c 1292774 44351 3454 17.219036 1 43025 3 lab_a unit_b 661842 30200 6147 11.882615 2 45878 8 lab_b unit_b 292885 30928 7864 28.959206 3 697 7 lab_a unit_a 1352669 81372 3769 3.728837 4 51402 8 lab_a unit_c 517981 48154 381 45.606934
已完成的三轮优化测试
1. 初始方案:Pandas默认merge实现
直接使用Pandas默认merge实现内连接,基准规模下得到预期的8334万行结果,运行耗时1分04秒,内存占用7.45GB(通过results.info(memory_usage="deep"统计),核心代码如下:
import pandas as pd results = pd.merge(experiments, patients, on="compound")
结果示例:
trial experiment observation compound value patient lab unit metric_1 metric_2 metric_3 metric_4 0 1 10 1 4 7.578612 11437 lab_a unit_b 1022481 24955 7312 43.395134 1 1 10 1 4 7.578612 60952 lab_a unit_c 873872 98759 1348 5.580664 2 1 10 1 4 7.578612 41207 lab_a unit_b 421455 88188 9705 27.077997 3 1 10 1 4 7.578612 62537 lab_a unit_a 645139 24159 6014 3.610864 4 1 10 1 4 7.578612 59984 lab_a unit_c 1176892 96816 6099 45.588840 ... ... ... ... ... ... ... ... ... ... ... ... ... 83320766 99898 5 1 5 6.272213 33903 lab_a unit_b 894670 98514 6620 9.320307 ...
2. 内存压缩优化
对两个DataFrame的字段做数据类型向下转型,将低基数字段设为category类型、数值字段适配为更小位宽的无符号整型/浮点型,同等数据规模下运行耗时58秒,内存占用降至3.49GB,核心代码如下:
experiments_schema = { "trial": "uint32", "experiment": "category", "observation": "uint8", "compound": "category", "value": "float32" } patients_schema = { "patient": "uint32", "compound": "category", "lab": "category", "unit": "category", "metric_1": "uint64", "metric_2": "uint32", "metric_3": "uint32", "metric_4": "float32", } experiments = experiments.astype(experiments_schema) patients = patients.astype(patients_schema)
3. 计算加速优化
将Pandas merge替换为Dask DataFrame的merge实现,同等数据规模下运行耗时降至26秒,内存占用保持3.49GB,核心代码如下:
import dask.dataframe as dd results = dd.merge(experiments, patients, on="compound", how="inner")
待咨询的优化方向
需要针对150000行experiments规模下的本地OOM问题给出可行方案,包含但不限于以下方向:
- 是否应将流水线后续计算步骤全部迁移至Dask实现,仅在流程最终环节调用
results.compute()? - 是否需要采用Cython或Numba实现merge操作以提升性能?
- 是否应通过磁盘存储中间结果降低内存压力?
- 其他可行的本地优化方案。
问题复现样例脚本
import numpy as np import pandas as pd rng = np.random.default_rng() experiment_count = 15000 experiments = pd.DataFrame( { "trial": np.sort(rng.choice(a=100000, size=experiment_count, replace=False)), "experiment": rng.choice(a=[3, 4, 5, 6, 7, 8, 9, 10], size=experiment_count, replace=True), "observation": rng.choice(a=[1, 2, 3], size=experiment_count, p=[0.9, 0.09, 0.01], replace=True), "compound": rng.choice(a=[1, 2, 3, 4, 5, 6, 7, 8, 9], size=experiment_count, replace=True), "value": 10*rng.random(size=experiment_count, dtype='float32') } ) patient_pool = 50000 patients = pd.DataFrame( { "patient": rng.choice(a=75000, size=patient_pool, replace=False), "compound": rng.choice(a=[1, 2, 3, 4, 5, 6, 7, 8, 9], size=patient_pool, replace=True), "lab": rng.choice(a=["lab_a", "lab_b", "lab_c"], size=patient_pool, p=[0.6, 0.3, 0.1], replace=True), "unit": rng.choice(a=["unit_a", "unit_b", "unit_c"], size=patient_pool, p=[0.3, 0.4, 0.3], replace=True), "metric_1": rng.choice(a=1500000, size=patient_pool, replace=True), "metric_2": rng.choice(a=100000, size=patient_pool, replace=True), "metric_3": rng.choice(a=10000, size=patient_pool, replace=True), "metric_4": 50*rng.random(size=patient_pool, dtype='float32'), } )
内容的提问来源于stack exchange,提问作者FinanceGardener

