使用Pandas UDF在PySpark中分组过滤DataFrame遇整数越界错误
先给个最简单的替代方案(强烈推荐)
看你的需求是过滤hour在指定范围内的数据,其实完全不需要用Pandas UDF做分组处理——Spark原生的filter操作就能直接搞定,既高效又能避免序列化相关的问题:
df2 = df1.filter((df1.hour > MIN_NIGHT_HOUR) & (df1.hour < MAX_NIGHT_HOUR)).show()
这种方式直接在Spark分布式引擎上执行,比Pandas UDF的序列化/反序列化开销小得多,也不会触发类型转换的错误。
如果必须用Pandas UDF(比如有更复杂的分组后逻辑)
你的错误根源是PyArrow在数据格式转换时的整数类型不匹配:Spark的IntegerType对应32位整数,但Pandas默认用64位整数(int64)存储数值。当PyArrow尝试把Pandas的int64数据转换成Spark的Int32格式时,哪怕实际数值没超出范围,也可能触发边界检查错误。
解决办法是强制将返回的Pandas DataFrame的整数列转换成int32,严格匹配原Spark Schema的类型:
import pandas as pd from pyspark.sql.types import PandasUDFType # 确保边界值是整数类型 MIN_NIGHT_HOUR = 23 MAX_NIGHT_HOUR = 6 @pandas_udf(df1.schema, PandasUDFType.GROUPED_MAP) def filter_data(pdf): # 执行过滤逻辑(这里假设夜间时段是23点到次日6点,用OR逻辑) filtered_pdf = pdf.query('hour > @MIN_NIGHT_HOUR OR hour < @MAX_NIGHT_HOUR') # 强制转换整数列到int32,匹配Spark的IntegerType integer_cols = ['hour', 'day', 'epochtimestamp'] filtered_pdf[integer_cols] = filtered_pdf[integer_cols].astype('int32') return filtered_pdf # 执行分组过滤 df2 = df1.groupBy('idvalue').apply(filter_data).show()
另外要注意:如果某个分组过滤后没有数据,返回的空DataFrame也要保持正确的列结构和类型,避免Spark自动推断类型时出错。
错误原因详解
Spark的Pandas UDF依赖PyArrow完成Pandas DataFrame和Spark数据格式之间的序列化。你的原Schema中hour、day、epochtimestamp都是integer(对应Spark的Int32类型),但Pandas处理后默认会用int64存储整数。当PyArrow尝试把int64数据转换成Spark的Int32格式时,会严格做边界检查,哪怕实际数值在范围内,也可能因为类型不匹配触发"Integer value out of bounds"错误。
内容的提问来源于stack exchange,提问作者Bitswazsky

