函数并行应用于含自定义对象的pandas Series返回NaN?求更快向量化方案
解决方案:向量化优化与并行返回NaN问题排查
一、并行返回NaN的常见原因
- 对象序列化失败:多进程并行时,自定义
Animal类对象可能无法被默认的pickle序列化,导致进程间传递对象出错,触发except块返回np.nan。 - 实例字段/键缺失:部分
Animal实例本身缺少a字段、b['y']键或c['x']['x1']嵌套键,单进程apply时可能只在个别行触发异常,但并行处理时这类问题会集中暴露。 - 并行环境依赖问题:比如子进程中未正确导入
numpy,导致np.nan无法正常生成,最终返回异常值。
二、比apply更快的向量化实现
要彻底提升性能,核心是把自定义对象的字段批量提取为pandas原生的Series,用底层C实现的向量化运算替代Python逐行循环。
方法1:列表推导式批量提取(性能最优)
列表推导式的速度远快于apply,先把所需字段逐一提取并处理异常:
import pandas as pd import numpy as np # 提取a字段,处理对象无a属性的情况 a_vals = [x.a if hasattr(x, 'a') else np.nan for x in df['animal']] # 提取b['y'],处理无b属性或y键不存在的情况 b_y_vals = [] for animal in df['animal']: try: b_y_vals.append(animal.b['y']) except (AttributeError, KeyError): b_y_vals.append(np.nan) # 提取c['x']['x1'],处理多层嵌套的缺失情况 c_x_x1_vals = [] for animal in df['animal']: try: c_x_x1_vals.append(animal.c['x']['x1']) except (AttributeError, KeyError): c_x_x1_vals.append(np.nan) # 转为Series后进行向量化计算 df['value'] = pd.Series(a_vals) + pd.Series(b_y_vals) + 2 * pd.Series(c_x_x1_vals)
方法2:map函数简洁写法(兼顾性能与代码简洁)
如果不想写多个循环,用Series.map配合lambda表达式也能达到接近列表推导式的性能:
s_a = df['animal'].map(lambda x: x.a if hasattr(x, 'a') else np.nan) s_by = df['animal'].map(lambda x: x.b['y'] if hasattr(x, 'b') and 'y' in x.b else np.nan) s_cxx1 = df['animal'].map(lambda x: x.c['x']['x1'] if hasattr(x, 'c') and 'x' in x.c and 'x1' in x.c['x'] else np.nan) df['value'] = s_a + s_by + 2 * s_cxx1
这两种方式的性能都远高于apply,因为pandas的向量化运算跳过了Python层面的逐行循环,直接调用底层优化的C代码。
三、并行优化的正确姿势
如果数据量极大必须用并行,先解决序列化问题:
- 给
Animal类添加__reduce__方法,让它能被pickle正确序列化;或者用dill库替代默认序列化器(需要修改并行工具的序列化配置)。 - 提前把所有
Animal对象的字段提取为普通Python类型(如整数、字典),再对这些普通数据进行并行处理,避免传递自定义对象。
比如用swifter库自动选择最优执行方式(需先解决序列化问题):
import swifter def calculate_value(x): try: return x.a + x.b['y'] + 2 * x.c['x']['x1'] except (AttributeError, KeyError): return np.nan df['value'] = df['animal'].swifter.apply(calculate_value)
内容的提问来源于stack exchange,提问作者Crivi
相关产品推荐
相关产品推荐

