如何使用whoswho库通过UDF逐行比较两个Spark DataFrame的姓名数据
实现方案
方案1:普通Python UDF实现逐行对比
如果你的数据量较小,直接自定义Python UDF即可满足需求,步骤如下:
- 首先保证所有Spark executor节点都已安装
whoswho依赖,可通过pip install whoswho安装。 - 定义匹配UDF并调用:
from pyspark.sql import functions as F from pyspark.sql.types import BooleanType from whoswho import who # 定义姓名匹配UDF,返回布尔值表示是否匹配 match_name_udf = F.udf(lambda name1, name2: who.match(name1, name2), BooleanType()) # 场景1:按两个DataFrame的行位置一一对应比对 # 给两个DF分别加行号、重名字段避免冲突 df1_with_id = df1.withColumnRenamed("name", "name_df1").withColumn("row_id", F.monotonically_increasing_id()) df2_with_id = df2.withColumnRenamed("name", "name_df2").withColumn("row_id", F.monotonically_increasing_id()) # 按行号关联后执行匹配 row_match_result = df1_with_id.join(df2_with_id, on="row_id", how="inner") \ .withColumn("is_match", match_name_udf(F.col("name_df1"), F.col("name_df2"))) \ .select("name_df1", "name_df2", "is_match") # 场景2:找出两个DF中所有互相匹配的姓名对 # 仅适合数据量较小的场景,数据量大时建议先按姓氏/首字母分桶过滤再关联 all_match_result = df1.crossJoin(df2.select(F.col("name").alias("name_df2"))) \ .withColumn("is_match", match_name_udf(F.col("name"), F.col("name_df2"))) \ .filter(F.col("is_match") == True)
方案2:更优方案 - Pandas向量化UDF
普通Python UDF存在JVM和Python进程之间频繁序列化/反序列化的开销,数据量大时性能较差,推荐使用Pandas UDF做向量化处理,性能可提升3~10倍:
import pandas as pd from pyspark.sql import functions as F from pyspark.sql.types import BooleanType from whoswho import who # 定义向量化匹配UDF,批量处理数据 @F.pandas_udf(BooleanType()) def match_name_pandas_udf(name1: pd.Series, name2: pd.Series) -> pd.Series: return pd.Series([who.match(n1, n2) for n1, n2 in zip(name1, name2)]) # 调用方式和普通UDF完全一致 result = df1_with_id.join(df2_with_id, on="row_id", how="inner") \ .withColumn("is_match", match_name_pandas_udf(F.col("name_df1"), F.col("name_df2")))
注意事项
- 数据量超过百万级时,不要直接使用交叉连接找全量匹配对,可以先对姓名做预处理,提取姓氏、首字母作为关联键,先过滤掉完全不可能匹配的组合,再调用匹配函数,可大幅减少计算量。
- 提交Spark任务时如果依赖没有预装在节点上,可以通过
--py-files参数打包whoswho依赖分发到所有executor。
内容的提问来源于stack exchange,提问作者John Doe
相关产品推荐
相关产品推荐

