Python多进程场景下内存泄漏问题排查求助
内存泄漏问题分析与解决方案
核心问题定位
你的代码内存耗尽并非传统意义的"泄漏",而是不必要的大内存对象生成和多进程下内存未及时回收,主要集中在以下几点:
- 全量笛卡尔积生成:
itertools.product(geo, geo)和pd.merge(how='cross')会生成N²条记录(N为df_emd行数),内存呈平方级增长,数据量稍大就会撑爆内存。 - 多进程子进程内存残留:
Pool的子进程可能持有大对象引用,加上Python GC延迟回收,导致内存持续累积。 - 低效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
相关产品推荐
相关产品推荐

