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

使用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

问题根源分析

  1. 返回结构不一致:初始化x_send2时定义的列是['Std_precio','Std_prec_cont','cant_muestras','Moda_precio_local','IsFo'],但后续动态添加了Moda_precio_contenido_std列,导致不同分组返回的DataFrame列数/列顺序不统一,Arrow无法推断全局一致的序列化类型。
  2. 全局模型状态污染:IForest实例od是全局变量,PySpark分组处理是分布式并行执行的,多个任务共享同一模型实例会导致状态混乱。
  3. 数据类型推断歧义:未显式指定返回列的数据类型,Arrow在序列化时无法准确推断部分列的类型。
  4. 赋值错误: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 10:25:24