PySpark DataFrame用自定义函数返回值新增列如何实现每行独立取值
问题原因
你当前的写法会导致所有行值相同,是因为getRandomString()是在Driver端仅执行1次,生成的固定值被lit()包装为常量列,全局所有行都复用这同一个值,不会在处理每行时重新执行函数。
解决方法
方法1:注册自定义UDF(通用自定义函数场景)
把你的Python函数注册为PySpark UDF,UDF会在Executor端处理每一行时单独调用,就能生成每行独立的结果:
from pyspark.sql.functions import udf from pyspark.sql.types import StringType import random # 定义UDF,指定返回值类型为字符串 @udf(returnType=StringType()) def getRandomString(): return "woteva" + str(random.randint(0,100)) # 直接调用UDF新增列即可 df = df.withColumn("MyNewColumn", getRandomString())
如果你的函数逻辑可以用Spark内置函数实现,优先用内置函数,性能比UDF高很多,可避免Python和JVM之间的序列化开销。
方法2:用Spark内置函数实现当前需求(性能更优)
你的需求是拼接固定前缀和0-100的随机整数,直接用Spark内置的randint函数(Spark 3.0+支持)和concat函数就能实现,不需要自定义UDF:
from pyspark.sql.functions import concat, lit, randint df = df.withColumn("MyNewColumn", concat(lit("woteva"), randint(0, 100)))
内容的提问来源于stack exchange,提问作者user16985119
相关产品推荐
相关产品推荐

