基于城市的场馆名称模糊Join实现(PySpark场景)
PySpark实现城市精确匹配+场馆名称模糊匹配关联数据集
需求说明
将新数据集df_new_stadium_data与现有场馆信息数据集df_stadium_information关联,规则如下:
- 城市(
city)必须精确匹配 - 场馆名称(
venue_name)需模糊匹配(非完全相等) - 同时满足上述条件时返回对应
venue_id,无符合条件的匹配则返回null
示例数据集
现有数据集(df_stadium_information)
| venue_name | city | venue_id |
|---|---|---|
| Sree Kanteerava Stadium | Bengaluru | 1 |
| Sree Kanteerava Stadium | Kochi | 2 |
| Eden Gardens | Kolkata | 3 |
| Narendra Modi Stadium | Ahmedabad | 4 |
新数据集(df_new_stadium_data)
| venue_name | city |
|---|---|
| Sri Kanteerava Indoor Stadium | Bengaluru |
| Eden Gardens | Kolkata |
期望输出
| venue_name | city | venue_id |
|---|---|---|
| Sri Kanteerava Indoor Stadium | Bengaluru | 1 |
| Eden Gardens | Kolkata | null |
实现代码
from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, regexp_extract, lower # 初始化SparkSession spark = SparkSession.builder.appName("StadiumVenueMatch").getOrCreate() # 构建示例数据集 # 现有场馆信息 existing_data = [ ("Sree Kanteerava Stadium", "Bengaluru", 1), ("Sree Kanteerava Stadium", "Kochi", 2), ("Eden Gardens", "Kolkata", 3), ("Narendra Modi Stadium", "Ahmedabad", 4) ] df_stadium_information = spark.createDataFrame(existing_data, ["venue_name", "city", "venue_id"]) # 新场馆数据 new_data = [ ("Sri Kanteerava Indoor Stadium", "Bengaluru"), ("Eden Gardens", "Kolkata") ] df_new_stadium_data = spark.createDataFrame(new_data, ["venue_name", "city"]) # 左连接:基于城市精确匹配 joined_df = df_new_stadium_data.alias("new").join( df_stadium_information.alias("existing"), col("new.city") == col("existing.city"), "left" ) # 定义模糊匹配条件: # 1. 场馆名称不完全相等 # 2. 新场馆名称包含现有场馆的核心关键词(提取如"X Stadium/X Gardens"这类核心部分) match_condition = ( col("new.venue_name") != col("existing.venue_name") & lower(col("new.venue_name")).contains(lower(regexp_extract(col("existing.venue_name"), r"(\w+ Stadium|\w+ Gardens)", 1))) ) # 生成结果:满足条件取venue_id,否则留null result_df = joined_df.withColumn( "venue_id", when(match_condition, col("existing.venue_id")) ).select( col("new.venue_name").alias("venue_name"), col("new.city").alias("city"), col("venue_id") ).groupBy("venue_name", "city").agg( col("venue_id").first().alias("venue_id") ) # 查看结果 result_df.show(truncate=False)
代码说明
- 精确匹配:通过
city字段的等值判断完成左连接,确保只关联同城市的场馆数据 - 模糊匹配逻辑:使用正则提取现有场馆名称的核心关键词(如"Kanteerava Stadium"),并检查新场馆名称是否包含该核心词,同时排除完全匹配的情况
- 结果聚合:通过分组聚合确保每个新数据行只返回一个有效
venue_id(若存在多个匹配可根据需求调整聚合规则)
内容的提问来源于stack exchange,提问作者Harshit Mahajan
相关产品推荐
相关产品推荐

