Apache Airflow环境中如何改写pandas代码降低内存占用?
低内存替代实现方案
原有代码的核心内存消耗来自两个点:一是isin操作会全量加载dadge_df的im_uuid列构建哈希映射表,二是两个DataFrame全量加载到内存时的存储开销,针对这两个点可选择以下适配方案:
方案1:仅加载必要列+分块处理(改造成本最低)
- 优先仅提取过滤需要的
im_uuid列转成Python集合,避免整表加载dadge_df占用额外内存 - 若
arct_df本身数据量极大,可从数据源分块读取过滤,不需要全量加载到内存
# 仅提取过滤用的uuid列转集合,内存占用仅为整表的1/5~1/10 exclude_uuids = set(dadge_df['im_uuid']) # 如arct_df从csv读取,可直接分块处理,chunk_size可根据Airflow worker可用内存调整 chunk_size = 10000 filtered_list = [] for chunk in pd.read_csv("arct_df源文件路径.csv", chunksize=chunk_size): filtered_list.append(chunk[~chunk['im_uuid'].isin(exclude_uuids)]) # 合并过滤后的分块,若结果量极大也可直接逐块写文件,不需要合并到内存 arct_df = pd.concat(filtered_list, ignore_index=True)
方案2:类型优化(适合两个DataFrame已加载到内存的场景)
字符串类型的im_uuid内存占用极高,转成分类类型后可降低70%以上的内存开销:
# 统一两个df的uuid分类基准,避免类型转换时的额外开销 all_uuids = pd.unique(pd.concat([arct_df['im_uuid'], dadge_df['im_uuid']])) arct_df['im_uuid'] = pd.Categorical(arct_df['im_uuid'], categories=all_uuids) dadge_df['im_uuid'] = pd.Categorical(dadge_df['im_uuid'], categories=all_uuids) # 再执行过滤操作,内存占用远低于字符串类型的isin计算 arct_df = arct_df[~arct_df.im_uuid.isin(dadge_df.im_uuid)]
方案3:Dask分布式处理(适合数据量远大于worker内存的场景)
Dask的API和pandas完全兼容,自动做分块计算,不会出现内存溢出问题:
import dask.dataframe as dd # Dask自动分块读取文件,不需要全量加载到内存 dd_arct = dd.read_csv("arct_df源文件路径.csv") dd_dadge = dd.read_csv("dadge_df源文件路径.csv") # 过滤逻辑和pandas完全一致 dd_filtered = dd_arct[~dd_arct.im_uuid.isin(dd_dadge.im_uuid)] # 结果量小可直接拉取到内存,结果量大直接写文件即可 arct_df = dd_filtered.compute() # 大结果直接写文件:dd_filtered.to_csv("过滤后的结果_*.csv", index=False)
额外优化建议
- 过滤前可调用
df.info(memory_usage='deep')查看各列内存占用,优先把大占比的字符串、大数值类型转成更小的存储格式 - 若使用K8sExecutor部署Airflow,也可给当前任务单独配置更高的内存配额,避免代码改造
内容的提问来源于stack exchange,提问作者caasswa
相关产品推荐
相关产品推荐

