PySpark:UDF操作后的数据类型判定及两列相加场景实践
在PySpark中用UDF实现两列相加:必须声明返回类型的示例
咱得记住,在PySpark里使用自定义函数(UDF)的时候,一定要在创建UDF时明确声明返回数据类型——要是漏掉这一步,Spark很可能会推断出错误的类型,要么计算出错,要么结果不是你想要的样子。
比如有这么个场景:我们有个包含ID、A、B三列的DataFrame,想要把A和B两列的值相加,生成一个新的Result列。下面是完整的实现代码:
from pyspark.sql.functions import udf, array from pyspark.sql.types import IntegerType # 创建UDF,明确指定返回类型为整数类型 udf_add = udf(lambda x: x[0] + x[1], IntegerType()) # 构造测试DataFrame,调用UDF生成Result列并展示 spark.createDataFrame([(101, 1, 16)], ['ID', 'A', 'B']).withColumn('Result', udf_add(array('A', 'B'))).show()
执行这段代码后,会得到如下输出:
+---+---+---+------+ | ID| A| B|Result| +---+---+---+------+ |101| 1| 16| 17| +---+---+---+------+
这里的关键点在于:我们用array('A', 'B')把A、B两列打包成数组传递给UDF,lambda函数负责取出数组里的两个元素相加,而IntegerType()则明确告诉Spark这个UDF的返回值是整数类型,这样Spark就能正确处理数据并输出预期的结果了。
内容的提问来源于stack exchange,提问作者Clock Slave
相关产品推荐
相关产品推荐

