PySpark调用Jellyfish函数报str argument expected错误的解决方案咨询
PySpark集成Jellyfish计算姓名相似度报错解决方案
错误根因
你当前代码的核心问题是直接将PySpark的Column列对象传入了仅支持Python原生字符串类型的jellyfish.jaro_winkler方法:你调用self.jaro_winkler_score(df[column_left], df[column_right])时传入的参数是Spark逻辑列对象,不是每行的实际字符串值,因此即使加了空值判断也会触发类型错误。
可行解决方案
你需要将字符串匹配逻辑包装为PySpark UDF(用户自定义函数),让Spark可以将这个函数分发到各节点逐行执行计算,适配Spark 2.3版本的代码如下:
第一步:导入依赖并注册UDF
import jellyfish import pyspark.sql.functions as F from pyspark.sql.types import DoubleType from pyspark.sql import SparkSession, DataFrame from pyspark import SparkContext # 注册UDF,指定返回值为浮点型 @F.udf(returnType=DoubleType()) def jaro_winkler_score_udf(s1, s2): if s1 is None or s2 is None: return 0.0 # 强制转换为字符串,避免底层数据类型不匹配 return jellyfish.jaro_winkler(str(s1), str(s2))
第二步:修改列计算逻辑
def jaro_winkler_func(self, df, column_left, column_right): # 使用注册好的UDF计算相似度 df = df.withColumn('test', jaro_winkler_score_udf(F.col(column_left), F.col(column_right))) return df # 调用逻辑保持不变 df = self.jaro_winkler_func(df, 'df1.first_name', 'df2.first_name')
注意事项
- 如果你使用分布式集群运行任务,需要保证所有Worker节点都已安装jellyfish库,否则会报模块不存在错误
- 执行计算前可通过
df.printSchema()确认两个姓名列的类型为StringType,若不是可在传入UDF前用F.col(column_left).cast("string")强制转换 - 若数据量极大,可以升级Spark到2.4+版本后使用Pandas UDF进一步提升计算性能
内容的提问来源于stack exchange,提问作者MAMS
相关产品推荐
相关产品推荐

