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

如何向PySpark UDF传入两个DataFrame并将结果存入新DataFrame

PySpark中使用FuzzyWuzzy计算字符串相似度的正确方法

问题分析

你的代码存在几个关键错误:

  • UDF定义错误:PySpark UDF是逐行处理单条记录的,并非接收整个DataFrame作为参数;np.vectorize用于批量处理numpy数组,不适合在UDF中使用;返回类型应为IntegerType(fuzz.token_sort_ratio返回0-100的整数),而非StringType。
  • 调用方式错误:无法直接将两个独立DataFrame的列传入UDF,必须先通过关联键合并两个DataFrame,让需要对比的字符串处于同一行。

正确实现步骤

1. 导入依赖库

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, udf
from pyspark.sql.types import IntegerType
from fuzzywuzzy import fuzz

2. 修正UDF定义

@udf(IntegerType())    
def fuzz_ratio(str1, str2):
    # 处理空值避免报错
    if str1 is None or str2 is None:
        return None
    # 直接对单个字符串对计算相似度
    return fuzz.token_sort_ratio(str1, str2)

3. 合并两个输入DataFrame

假设df1和df2有共同关联键(如id),先合并到同一DataFrame:

# 重命名列避免冲突
df1_renamed = df1.withColumnRenamed("VAL", "VAL_FROM_DF1")
df2_renamed = df2.withColumnRenamed("VAL", "VAL_FROM_DF2")

# 按关联键合并(可根据需求改为left/right/full join)
joined_df = df1_renamed.join(df2_renamed, on="id", how="inner")

如果无关联键,需对比所有字符串组合(笛卡尔积):

joined_df = df1_renamed.crossJoin(df2_renamed)

4. 调用UDF生成结果

result_df = joined_df.withColumn("VAL", fuzz_ratio(col("VAL_FROM_DF1"), col("VAL_FROM_DF2")))

5. 可选:性能优化(Pandas UDF)

数据量大时,用Pandas UDF批量处理提升效率:

from pyspark.sql.functions import pandas_udf
import pandas as pd

@pandas_udf(IntegerType())
def fuzzy_ratio_batch(str1_series: pd.Series, str2_series: pd.Series) -> pd.Series:
    return pd.Series([fuzz.token_sort_ratio(s1, s2) if pd.notna(s1) and pd.notna(s2) else None for s1, s2 in zip(str1_series, str2_series)])

# 调用方式一致
result_df = joined_df.withColumn("VAL", fuzzy_ratio_batch(col("VAL_FROM_DF1"), col("VAL_FROM_DF2")))

内容的提问来源于stack exchange,提问作者Danial Khilji

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 07:15:36