将自定义函数应用于 PySpark 中数据框选定列的单元格

Ang*_*gie 5 python apache-spark pyspark spark-dataframe

假设我有一个如下所示的数据框:

+---+-----------+-----------+
| id|   address1|   address2|
+---+-----------+-----------+
|  1|address 1.1|address 1.2|
|  2|address 2.1|address 2.2|
+---+-----------+-----------+
Run Code Online (Sandbox Code Playgroud)

我想将自定义函数直接应用于address1和address2列中的字符串,例如:

def example(string1, string2):
    name_1 = string1.lower().split(' ')
    name_2 = string2.lower().split(' ')
    intersection_count = len(set(name_1) & set(name_2))

    return intersection_count
Run Code Online (Sandbox Code Playgroud)

我想将结果存储在一个新列中,以便我的最终数据框如下所示:

+---+-----------+-----------+------+
| id|   address1|   address2|result|
+---+-----------+-----------+------+
|  1|address 1.1|address 1.2|     2|
|  2|address 2.1|address 2.2|     7|
+---+-----------+-----------+------+
Run Code Online (Sandbox Code Playgroud)

我尝试以一种曾经将内置函数应用于整个列的方式来执行它,但出现错误:

>>> df.withColumn('result', example(df.address1, df.address2))
Traceback (most recent call last):
  File "<stdin>", line 1, in <module>
  File "<stdin>", line 2, in example
TypeError: 'Column' object is not callable
Run Code Online (Sandbox Code Playgroud)

我做错了什么以及如何将自定义函数应用于选定列中的字符串?

dum*_*tru 6

您必须在 spark 中使用 udf(用户定义函数)

from pyspark.sql.functions import udf
example_udf = udf(example, LongType())
df.withColumn('result', example_udf(df.address1, df.address2))
Run Code Online (Sandbox Code Playgroud)

  • 是的,它应该是给定函数的返回类型 (2认同)