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

PySpark跨DataFrame匹配文本子串新增映射列实现方法

PySpark 文本子串匹配映射实现方案

问题根因

之前用Python字典配合map函数无法得到正确结果的核心原因:map仅支持列值与字典key的全量精确匹配,只有列的完整内容和key完全一致时才会返回映射值,无法识别文本中间的子串包含关系。

实现思路

针对城市映射表(B表)数据量极小的场景,优先用笛卡尔积关联后过滤匹配项的方式实现,无需写UDF,性能稳定;如果A表数据量达到亿级,可选择广播变量+UDF的方式减少shuffle开销。两种方案都支持城市名出现在TITLE任意位置的子串匹配。

可直接运行的代码示例

1. 构造测试数据集

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, expr, array_join

# 初始化Spark会话
spark = SparkSession.builder.appName("city_shortcut_match").getOrCreate()

# 构造A表测试数据(覆盖5个城市在文本不同位置的场景)
data_a = [
    (1, "001", "New York is the most populous city in the United States"),
    (2, "002", "The automobile industry was once the core pillar of Detroit"),
    (3, "003", "Miami is famous for its long tropical coastline and beach scenery"),
    (4, "004", "The global film center Hollywood is located in Los Angeles"),
    (5, "005", "Top universities like MIT and Harvard are situated in Boston"),
    (6, "006", "Many business travelers commute between New York and Boston weekly")
]
df_a = spark.createDataFrame(data_a, schema=["ID", "SOME_CODE", "TITLE"])

# 构造B表城市映射数据
data_b = [
    ("New York", "NY"),
    ("Detroit", "DET"),
    ("Miami", "MIA"),
    ("Los Angeles", "LA"),
    ("Boston", "BOS")
]
df_b = spark.createDataFrame(data_b, schema=["City", "Shortcut"])

2. 方案一:Cross Join + 子串过滤(推荐,无UDF,适配小维度表场景)

该方案无需自定义函数,利用Spark内置函数实现,稳定性高,完全适配当前5个城市映射的业务场景:

# 笛卡尔积关联所有可能的城市匹配对,过滤出文本包含对应城市名的记录
matched_df = df_a.crossJoin(df_b) \
    .filter(expr("instr(TITLE, City) > 0")) \
    .groupBy("ID", "SOME_CODE", "TITLE") \
    # 若单条文本包含多个城市,用逗号拼接所有匹配的缩写;单匹配场景可替换为first(Shortcut)
    .agg(array_join(expr("collect_list(Shortcut)"), ",").alias("SHORTCUT"))

# 左关联补全未匹配到城市的记录,保证A表原始数据不丢失
result = df_a.join(matched_df, on=["ID", "SOME_CODE", "TITLE"], how="left")

3. 方案二:广播变量+UDF(适配A表超大数据量场景)

如果A表数据量极大,可将小维度表B广播到所有计算节点,避免笛卡尔积带来的shuffle开销:

from pyspark.sql.functions import udf
from pyspark.sql.types import StringType

# 收集B表数据为字典并广播到所有Executor
city_mapping = {row["City"]: row["Shortcut"] for row in df_b.collect()}
bc_city_mapping = spark.sparkContext.broadcast(city_mapping)

# 定义子串匹配UDF
def get_city_shortcut(title_text):
    matched_sc = []
    for city_name, sc in bc_city_mapping.value.items():
        if city_name in title_text:
            matched_sc.append(sc)
    return ",".join(matched_sc) if matched_sc else None

match_udf = udf(get_city_shortcut, StringType())
result = df_a.withColumn("SHORTCUT", match_udf(col("TITLE")))

结果校验

执行result.show(truncate=False)即可得到符合要求的输出,返回字段为ID、SOME_CODE、TITLE、SHORTCUT,无论城市名出现在文本的开头、中间还是结尾都能正确匹配。


内容的提问来源于stack exchange,提问作者Krzysztof Owadowski

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 21:27:22