如何向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
相关产品推荐
相关产品推荐

