使用IForest检测异常值时遇ArrowInvalid类型推断错误求助
问题背景:使用IForest在PySpark pandas分组中进行异常检测触发ArrowInvalid错误
代码实现
import pyspark.pandas as pd import numpy as np from alibi_detect.od import IForest # **************** Modelo IForest ****************************************** # IForest rta - Outlier ---> 1, Not-Outlier ----> 0 od = IForest( threshold=0., n_estimators=5 ) def mode(lm): freqs = groupby(Counter(lm).most_common(), lambda x:x[1]) m=[val for val,count in next(freqs)[1]] if len(m)>1: m=np.median(lm) else: m=float(m[0]) return m def disper(x): x_pred = x[['precio_local', 'precio_contenido']] insumo_std = x_pred.std().to_frame().T mod = mode(x_pred['precio_local']) x_send2 = pd.DataFrame( index=x_pred.index, columns=['Std_precio','Std_prec_cont','cant_muestras','Moda_precio_local','IsFo'] ) x_send2.loc[:,'Std_precio'] = insumo_std.loc[0,'precio_local'] x_send2.loc[:,'Std_prec_cont'] = insumo_std.loc[0,'precio_local'] x_send2.loc[:,'Moda_precio_local'] = mod mod_cont = mode(x_pred['precio_contenido']) x_send2.loc[:,'Moda_precio_contenido_std'] = mod_cont ctn = x_pred.shape[0] x_send2.loc[:,'cant_muestras'] = ctn if x_pred.shape[0]>3: od.fit(x_pred) preds = od.predict( x_pred, return_instance_score=True ) x_preds = preds['data']['is_outlier'] pd.set_option('compute.ops_on_diff_frames', True) x_send2.loc[:,'IsFo']= pd.Series(x_preds, index=x_pred.index) else: x_send2.loc[:,'IsFo'] = 0 print(type(x_send2)) print(x_send2) return x_send2 insumo_all_pd = insumo_all.to_pandas_on_spark()
执行分组操作时触发错误:
df_result = insumo_all_pd.groupby(by=['categoria','marca','submarca','barcode','contenido_std','unidad_std']).apply(disper)
错误信息
ArrowInvalid Traceback (most recent call last) <command-1939548125702628> in <module> ----> 1 df_result = insumo_all_pd.groupby(by=['categoria','marca','submarca','barcode','contenido_std','unidad_std']).apply(disper) 2 display(df_result) /databricks/spark/python/pyspark/pandas/usage_logging/__init__.py in wrapper(*args, **kwargs) 192 start = time.perf_counter() 193 try: ---> 194 res = func(*args, **kwargs) 195 logger.log_success( 196 class_name, function_name, time.perf_counter() - start, signature /databricks/spark/python/pyspark/pandas/groupby.py in apply(self, func, *args, **kwargs) 1200 else: 1201 pser_or_pdf = grouped.apply(pandas_apply, *args, **kwargs) -> 1202 psser_or_psdf = ps.from_pandas(pser_or_pdf) 1203 1204 if len(pdf) <= limit: /databricks/spark/python/pyspark/pandas/usage_logging/__init__.py in wrapper(*args, **kwargs) 187 if hasattr(_local, "logging") and _local.logging: 188 # no need to log since this should be internal call. ---> 189 return func(*args, **kwargs) 190 _local.logging = True 191 try: /databricks/spark/python/pyspark/pandas/namespace.py in from_pandas(pobj) 143 """ 144 if isinstance(pobj, pd.Series): ---> 145 return Series(pobj) 146 elif isinstance(pobj, pd.DataFrame): 147 return DataFrame(pobj) /databricks/spark/python/pyspark/pandas/usage_logging/__init__.py in wrapper(*args, **kwargs) 187 if hasattr(_local, "logging") and _local.logging: 188 # no need to log since this should be internal call. ---> 189 return func(*args, **kwargs) 190 _local.logging = True 191 try: /databricks/spark/python/pyspark/pandas/series.py in __init__(self, data, index, dtype, name, copy, fastpath) 424 data=data, index=index, dtype=dtype, name=name, copy=copy, fastpath=fastpath 425 ) ---> 426 internal = InternalFrame.from_pandas(pd.DataFrame(s)) 427 if s.name is None: 428 internal = internal.copy(column_labels=[None]) /databricks/spark/python/pyspark/pandas/internal.py in from_pandas(pdf) 1458 data_columns, 1459 data_fields, -> 1460 ) = InternalFrame.prepare_pandas_frame(pdf) 1461 1462 schema = StructType([field.struct_field for field in index_fields + data_fields]) /databricks/spark/python/pyspark/pandas/internal.py in prepare_pandas_frame(pdf, retain_index) 1531 1532 for col, dtype in zip(reset_index.columns, reset_index.dtypes): -> 1533 spark_type = infer_pd_series_spark_type(reset_index[col], dtype) 1534 reset_index[col] = DataTypeOps(dtype, spark_type).prepare(reset_index[col]) 1535 /databricks/spark/python/pyspark/pandas/typedef/typehints.py in infer_pd_series_spark_type(pser, dtype) 327 return pser.iloc[0].__UDT__ 328 else: -> 329 return from_arrow_type(pa.Array.from_pandas(pser).type) 330 elif isinstance(dtype, CategoricalDtype): 331 if isinstance(pser.dtype, CategoricalDtype): /databricks/python/lib/python3.8/site-packages/pyarrow/array.pxi in pyarrow.lib.Array.from_pandas() /databricks/python/lib/python3.8/site-packages/pyarrow/array.pxi in pyarrow.lib.array() /databricks/python/lib/python3.8/site-packages/pyarrow/array.pxi in pyarrow.lib._ndarray_to_array() /databricks/python/lib/python3.8/site-packages/pyarrow/error.pxi in pyarrow.lib.check_status() ArrowInvalid: Could not convert Std_precio Std_prec_cont cant_muestras Moda_precio_local IsFo Moda_precio_contenido_std 107 0.0 0.0 3 1.0 0 1.666667 252 0.0 0.0 3 1.0 0 1.666667 396 0.0 0.0 3 1.0 0 1.666667 with type DataFrame: did not recognize Python value type when inferring an Arrow data type
数据集Schema
fecha_ola datetime64[ns] pais object categoria object marca object submarca object contenido_std float64 unidad_std object barcode object precio_local float64 cantidad float64 descripcion object id_ticket object id_item object id_pdv object fecha_transaccion datetime64[ns] id_ref float64 precio_contenido float64 dtype: object
问题根源分析
- 返回结构不一致:初始化
x_send2时定义的列是['Std_precio','Std_prec_cont','cant_muestras','Moda_precio_local','IsFo'],但后续动态添加了Moda_precio_contenido_std列,导致不同分组返回的DataFrame列数/列顺序不统一,Arrow无法推断全局一致的序列化类型。 - 全局模型状态污染:IForest实例
od是全局变量,PySpark分组处理是分布式并行执行的,多个任务共享同一模型实例会导致状态混乱。 - 数据类型推断歧义:未显式指定返回列的数据类型,Arrow在序列化时无法准确推断部分列的类型。
- 赋值错误:
Std_prec_cont列错误赋值为precio_local的标准差,而非precio_contenido的标准差。
解决方案
1. 统一返回DataFrame的列结构
初始化时就包含所有输出列,避免动态添加:
x_send2 = pd.DataFrame( index=x_pred.index, columns=['Std_precio','Std_prec_cont','cant_muestras','Moda_precio_local','Moda_precio_contenido_std','IsFo'] )
2. 每个分组内独立初始化IForest
避免全局模型的状态污染,在disper函数内部创建模型:
def disper(x): x_pred = x[['precio_local', 'precio_contenido']] # 每个分组独立初始化模型 od = IForest(threshold=0., n_estimators=5) # ... 后续代码保持不变
3. 显式指定列数据类型
创建x_send2时明确设置各列的dtype,消除Arrow的类型推断歧义:
x_send2 = pd.DataFrame( index=x_pred.index, data={ 'Std_precio': np.nan, 'Std_prec_cont': np.nan, 'cant_muestras': 0, 'Moda_precio_local': np.nan, 'Moda_precio_contenido_std': np.nan, 'IsFo': 0 }, dtype={ 'Std_precio': 'float64', 'Std_prec_cont': 'float64', 'cant_muestras': 'int64', 'Moda_precio_local': 'float64', 'Moda_precio_contenido_std': 'float64', 'IsFo': 'int64' } )
4. 修正标准差赋值错误
将Std_prec_cont的赋值改为precio_contenido的标准差:
x_send2.loc[:,'Std_prec_cont'] = insumo_std.loc[0,'precio_contenido']
5. 补全缺失的导入
mode函数用到的Counter和groupby需要导入:
from collections import Counter from itertools import groupby
完整修复后的代码
import pyspark.pandas as pd import numpy as np from alibi_detect.od import IForest from collections import Counter from itertools import groupby # **************** Modelo IForest ****************************************** # IForest rta - Outlier ---> 1, Not-Outlier ----> 0 def mode(lm): freqs = groupby(Counter(lm).most_common(), lambda x:x[1]) m=[val for val,count in next(freqs)[1]] if len(m)>1: m=np.median(lm) else: m=float(m[0]) return m def disper(x): x_pred = x[['precio_local', 'precio_contenido']] # 每个分组独立初始化模型 od = IForest(threshold=0., n_estimators=5) insumo_std = x_pred.std().to_frame().T mod = mode(x_pred['precio_local']) # 显式定义所有列和数据类型 x_send2 = pd.DataFrame( index=x_pred.index, data={ 'Std_precio': np.nan, 'Std_prec_cont': np.nan, 'cant_muestras': 0, 'Moda_precio_local': np.nan, 'Moda_precio_contenido_std': np.nan, 'IsFo': 0 }, dtype={ 'Std_precio': 'float64', 'Std_prec_cont': 'float64', 'cant_muestras': 'int64', 'Moda_precio_local': 'float64', 'Moda_precio_contenido_std': 'float64', 'IsFo': 'int64' } ) x_send2.loc[:,'Std_precio'] = insumo_std.loc[0,'precio_local'] # 修正标准差赋值错误 x_send2.loc[:,'Std_prec_cont'] = insumo_std.loc[0,'precio_contenido'] x_send2.loc[:,'Moda_precio_local'] = mod mod_cont = mode(x_pred['precio_contenido']) x_send2.loc[:,'Moda_precio_contenido_std'] = mod_cont ctn = x_pred.shape[0] x_send2.loc[:,'cant_muestras'] = ctn if x_pred.shape[0]>3: od.fit(x_pred) preds = od.predict( x_pred, return_instance_score=True ) x_preds = preds['data']['is_outlier'] pd.set_option('compute.ops_on_diff_frames', True) x_send2.loc[:,'IsFo']= pd.Series(x_preds, index=x_pred.index) else: x_send2.loc[:,'IsFo'] = 0 return x_send2 insumo_all_pd = insumo_all.to_pandas_on_spark()
内容的提问来源于stack exchange,提问作者Daniel Vera
相关产品推荐
相关产品推荐

