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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 19:17:33