PySpark配置Arrow后,Pandas UDF内为何未使用Arrow格式数据?
Spark Pandas UDF中未看到Arrow-backed数据的问题解析
为什么UDF里是object类型而非Arrow格式?
Spark开启Arrow优化后,只是在Spark和Python进程之间用Arrow传输数据,数据进入Pandas UDF后,Pandas会自动把Arrow-backed的数据转成自身默认的类型——比如字符串列变成object dtype,结构体列变成嵌套object。这不是bug,是当前版本的正常逻辑:Arrow是传输层的优化,不是让UDF里的Pandas DataFrame保持Arrow原生格式。
先确认Arrow配置真的生效了
初始化SparkSession时强制开启Arrow并禁用降级:
spark = SparkSession.builder \ .appName("ArrowCheck") \ .config("spark.sql.execution.arrow.pyspark.enabled", "true") \ .config("spark.sql.execution.arrow.pyspark.fallback.enabled", "false") \ .getOrCreate()禁用降级后,如果Arrow没生效会直接报错,能快速验证配置是否起作用。
查看Spark日志:运行UDF时,日志里会出现
Using Arrow for pandas UDF的字样,说明传输阶段确实用了Arrow。
要在UDF里直接用Arrow数据?这么做
如果想全程操作Arrow格式的数据,别用普通的Pandas UDF,改用直接处理Arrow RecordBatch的写法:
from pyspark.sql.functions import pandas_udf import pyarrow as pa @pandas_udf("string") def arrow_native_udf(arrow_batch: pa.RecordBatch) -> pa.Array: # 直接操作Arrow的RecordBatch,不用转成Pandas str_column = arrow_batch.column("your_string_col") return str_column.map(lambda s: s.upper())
这种写法下,数据全程以Arrow格式流转,不会被Pandas转成object类型。
版本兼容提醒
pyspark 3.5.1和pyarrow 15.0.1是兼容的,但要注意:
- 集群所有节点的pyarrow版本必须一致,不然会出现序列化错误
- 禁用降级后,如果遇到Arrow不支持的复杂类型,任务会失败,这时要么调整UDF的类型定义,要么临时开启降级开关
内容的提问来源于stack exchange,提问作者mdurant
相关产品推荐
相关产品推荐

