You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.25 14:24:05