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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 13:54:06