PySpark中Pandas UDF性能不及Python UDF?求原因排查
Pandas UDF性能反而不如Python UDF?
我理解Pandas UDF借助Arrow减少序列化开销,还支持向量计算,性能应该比Python UDF好,但下面的测试代码结果却相反,这是什么原因?还是我的测试方式有问题?
from time import perf_counter import pandas as pd from pyspark.sql import SparkSession from pyspark.sql.functions import * from pyspark.sql.types import * spark = SparkSession.builder.appName("TEST").getOrCreate() sdf = spark.range(0, 1000000).withColumn( 'id', col('id') ).withColumn('v', rand()) @pandas_udf(DoubleType()) def pandas_plus_one(pdf): return pdf + 1 @udf(DoubleType()) def plus_one(num): return num + 1 # Pandas UDF res_pdf = sdf.select(pandas_plus_one(col("v"))) st = perf_counter() for _ in range(10): res_pdf.show() print(f"Pandas UDF Time: {(perf_counter() - st) * 1000} ms") # Python UDF res = sdf.select(plus_one(col("v"))) st = perf_counter() for _ in range(10): res.show() print(f"Python UDF Time: {(perf_counter() - st) * 1000} ms")
核心原因分析
1. 测试方法完全不准确
show()默认只返回前20行,触发的是部分计算,而且包含大量额外开销:数据从集群传输到Driver的IO、终端渲染输出等。这些开销会完全掩盖UDF本身的性能差异,根本无法反映真实的计算性能。
2. 任务太简单,Pandas UDF的优势无法体现
你的UDF逻辑只是简单的+1,属于极轻量计算。Python UDF的逐行序列化开销在这种场景下占比极低,而Pandas UDF存在批量数据转换、Arrow序列化的固定启动开销,小任务下反而会显得更慢。只有当计算逻辑复杂(比如批量字符串处理、滑动窗口统计、复杂数值运算)、数据量足够大时,Pandas UDF的向量计算优势才会显现。
3. Spark的隐性优化抵消了Pandas UDF的优势
对于这种极简的+1操作,Spark可能会对Python UDF做隐性优化(甚至有可能将其转换为JVM端的内置表达式),而Pandas UDF必须走完整的Arrow序列化/反序列化流程,在简单计算下这个流程的开销超过了它的优势。
修正后的测试代码
用count()触发全量计算,避免show()的额外开销,同时放大数据量来体现Pandas UDF的优势:
from time import perf_counter import pandas as pd from pyspark.sql import SparkSession from pyspark.sql.functions import * from pyspark.sql.types import * spark = SparkSession.builder.appName("TEST").getOrCreate() # 放大数据量到1000万行 sdf = spark.range(0, 10000000).withColumn('v', rand()) @pandas_udf(DoubleType()) def pandas_plus_one(pdf): return pdf + 1 @udf(DoubleType()) def plus_one(num): return num + 1 # Pandas UDF 测试:触发全量计算 res_pdf = sdf.select(pandas_plus_one(col("v"))) st = perf_counter() res_pdf.count() print(f"Pandas UDF Time: {(perf_counter() - st) * 1000} ms") # Python UDF 测试 res = sdf.select(plus_one(col("v"))) st = perf_counter() res.count() print(f"Python UDF Time: {(perf_counter() - st) * 1000} ms")
当数据量足够大、计算逻辑复杂时,Pandas UDF的批量处理和向量计算优势会显著超过逐行处理的Python UDF。
内容的提问来源于stack exchange,提问作者jiashenC
相关产品推荐
相关产品推荐

