【发布时间】:2019-08-16 08:17:27
【问题描述】:
pyspark 的初学者试图理解 UDF:
我有一个 PySpark 数据框 p_b,我通过传递数据框的所有行来调用 UDF。我想从行访问列debit。出于某种原因,这没有发生。请在下面找到 sn-ps。
p_b has 4 columns, id, credit, debit,sum
功能:
def test(row):
return('123'+row['debit'])
转换为 UDF
test_udf=udf(test,IntegerType())
在数据帧 p_b 上调用 UDF
vals=test_udf(struct([p_b[x] for x in p_b.columns]))
print(type(vals))
print(vals)
输出
Column<b'test(named_struct(id, credit,debit,sum))'>
【问题讨论】:
-
您似乎正在尝试将“123”添加到数据框的每一行。不是吗?
-
您必须使用 with 列为您的数据框调用 udf,数据框列值必须作为参数传递。像这样定义你的函数。 def user_func(row): 返回行+123
-
my_func = udf(user_func, IntegerType()) newdf = df.withColumn('new_column',my_func(df.value))
-
感谢 cmets。我试图将“123”添加到“借方”列的所有行
标签: python pyspark pyspark-sql