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

Python多进程场景下内存泄漏问题排查求助

内存泄漏问题分析与解决方案

核心问题定位

你的代码内存耗尽并非传统意义的"泄漏",而是不必要的大内存对象生成和多进程下内存未及时回收,主要集中在以下几点:

  1. 全量笛卡尔积生成:itertools.product(geo, geo)和pd.merge(how='cross')会生成N²条记录(N为df_emd行数),内存呈平方级增长,数据量稍大就会撑爆内存。
  2. 多进程子进程内存残留:Pool的子进程可能持有大对象引用,加上Python GC延迟回收,导致内存持续累积。
  3. 低效DataFrame操作:大量apply(axis=1)逐行处理、重复merge和临时DataFrame未及时清理,进一步加剧内存占用。

具体修复方案

1. 替换笛卡尔积为空间索引查询(彻底解决内存爆炸)

用scipy.spatial.cKDTree快速定位1km内的邻居,避免生成全量两两组合:

from scipy.spatial import cKDTree

def process_emd(args, tuple_emd):
    sd, sgg, emd = tuple_emd
    
    # ... 前面读取数据、预处理代码不变 ...
    
    df_emd.reset_index(inplace=True)
    # 转换经纬度为弧度(匹配haversine计算逻辑)
    coords = np.radians(df_emd[['lon_x', 'lat_y']].values)
    # 构建KDTree
    tree = cKDTree(coords)
    # 1km对应的弧度:地球半径6367km,1km弧度为1/6367
    radius_rad = 1.0 / 6367
    # 查询每个点周围1km内的所有点索引
    neighbor_indices = tree.query_ball_tree(tree, radius_rad)
    
    # 生成有效邻居对(排除自身)
    pairs = []
    for i, neighbors in enumerate(neighbor_indices):
        for j in neighbors:
            if i != j:
                pairs.append((i, j))
    
    # 构建仅包含有效邻居对的DataFrame,替代原cross join
    df_pairs = pd.DataFrame(pairs, columns=['idx_x', 'idx_y'])
    df_emd_cp = df_pairs.merge(df_emd, left_on='idx_x', right_index=True) \
                        .merge(df_emd, left_on='idx_y', right_index=True, suffixes=('_x', '_y'))
    
    # 计算距离(KDTree已过滤1km内数据,可验证或直接使用)
    lon1, lat1 = df_emd_cp['lon_x_x'], df_emd_cp['lat_y_x']
    lon2, lat2 = df_emd_cp['lon_x_y'], df_emd_cp['lat_y_y']
    df_emd_cp['euclidean_distance'] = haversine_np(lon1, lat1, lon2, lat2)
    
    # ... 后续过滤条件不变 ...

2. 优化多进程内存管理

  • 避免子进程持有全局大对象:如果pair_road_cond、pair_land_usage是大字典,尽量在子进程内部加载,避免跨进程传递大对象。
  • 子进程内手动触发GC:在process_emd末尾强制回收内存:
import gc

def process_emd(args, tuple_emd):
    # ... 所有业务代码 ...
    
    # 清理变量后手动触发GC
    del df_emd_cp, df_land_ai, df_land, df_final, df_emd, coords, tree, neighbor_indices, pairs
    gc.collect()
    
    t.toc(msg=f'Elapsed time for getting similiar properties of {emd} : ', restart=True)

3. 替换低效apply为矢量化计算

把逐行apply改成Pandas矢量化操作,减少内存占用和计算时间:

# 先补上原代码遗漏的total_price_prop字段
df_emd_cp['total_price_prop'] = np.where(
    (df_emd_cp['종류_x'] != '토지') & (df_emd_cp['건물가격_추정_x'] != 0.),
    (df_emd_cp['건물가격_추정_y'] + df_emd_cp['공시가격_2022_y'] * df_emd_cp['Mutiple_Pred_y']) / 
    (df_emd_cp['건물가격_추정_x'] + df_emd_cp['공시가격_2022_x'] * df_emd_cp['Mutiple_Pred_x']),
    np.nan
)

# 矢量化计算check1和check2
lb_total, ub_total = 0.5, 1.5
lb_unit, ub_unit = 0.5, 1.5

df_emd_cp['check1'] = df_emd_cp['total_price_prop'].isna() | \
                      ((df_emd_cp['total_price_prop'] >= lb_total) & (df_emd_cp['total_price_prop'] <= ub_total))
df_emd_cp['check2'] = (df_emd_cp['unit_price_prop'] >= lb_unit) & (df_emd_cp['unit_price_prop'] <= ub_unit)

4. 简化groupby操作,避免嵌套列表残留

原groupby生成的嵌套列表容易残留内存引用,改成扁平化DataFrame处理:

TOP_K = 5
# 替换原df_similiar_min_candidates生成逻辑
df_min = df_emd_cp[df_emd_cp['check1']].copy()
# 按pnu_x分组,取距离最小的TOP_K条
df_similiar_min_candidates = df_min.groupby('pnu_x').apply(
    lambda x: x.sort_values('euclidean_distance').head(min(TOP_K, len(x)))[['pnu_y', 'euclidean_distance']]
).reset_index()
# 整理为需要的嵌套列表格式
df_similiar_min_candidates = df_similiar_min_candidates.groupby('pnu_x').apply(
    lambda x: list(zip(x['pnu_y'], x['euclidean_distance']))
).reset_index(name='similiar_min_candidates')

其他辅助优化

  • 限制子进程数量:NUM_CORES不要超过CPU核心数的一半,避免内存竞争。
  • 读取数据库时限制字段:get_land_ai和get_land只查询业务需要的字段,减少DataFrame内存占用。
  • 指定数据类型:读取数据时给数值字段指定更小的类型(如float32代替float64),降低内存开销。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 15:29:53