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

PySpark实现campaign_name包含与精确匹配的性能优化方案咨询

PySpark 数据映射业务性能优化方案咨询

1. 基础数据

业务表包含campaign_name字段,样例如下:

campaign_name
abcloancde
abcsolcdf
abcemicdef
emic_estore
Personalloa-nemic_sol
personalloa_nemic
abc/emic-dg-upi:bol

另有映射字典表,包含terms、product_category、product、function四个字段。

2. 原有逻辑与实现

2.1 逻辑说明

首先筛选字典表中function='contains'的记录,按terms长度倒序排序,匹配campaign_name时优先取第一个命中的terms对应的分类与产品信息,未命中的prod、prod_cat标记为other。

2.2 原有实现代码

#Dataset loaded below
initialData = spark.read.option("header", "true").csv("file://..../sample_data.csv")
initialData.show()
#Dictionary loaded below
df = spark.read.option("header", 
"true").csv("file://..../mapper.csv")

df_contains = df.filter(df.function == 'contains').drop("function")
df_contains = df_contains.orderBy(length(col("terms")).desc())
w = Window.partitionBy(lit('A')).orderBy(length(col("terms")).desc())
df_contains = df_contains.withColumn("rw", row_number().over(w))


df3 = df_contains.na.fill("").groupBy(lit(1)).agg(collect_list(
    concat(col("rw"), lit(":"), col("terms"), lit(":"), col("product_category"), lit(":"), col("product"))).alias(
    "Check")).withColumn("Check", concat_ws(",", col("Check"))).drop("1")


def categoryFunction(name, Check):
    # checkList = Check.lower().split(",")
    out = ""
    match = False
    for Key in Check.lower().split(","):
        keyword = Key.split(":", 2)
        terms = keyword[1]
        tempOut = keyword[2]
        if terms in name.lower():
            out = tempOut
            match = True
        if match:
            break
    return out


def categoryFunction1(name, Check):
    # checkList = Check.lower().split(",")
    out = ""
    match = False
    for Key in Check.lower().split(","):
        keyword = Key.split(":", 2)
        terms = keyword[1]
        tempOut = keyword[2]
        if terms == name.lower():
            out = tempOut
            match = True
        if match:
            break
    return out


categoryUDF = udf(categoryFunction, StringType())
categoryUDF1 = udf(categoryFunction1, StringType())

df4 = initialData.crossJoin(df3)

finalDF = df4.withColumn("out", categoryUDF(col("campaign_name"), col("Check"))).drop("Check").withColumn("out", split(
    col("out"), ":")).withColumn("product_category", col("out")[0]).withColumn("product", col("out")[1]).drop(
    "out").withColumn("prod", when(col("product").isNull(), "other").otherwise(col("product"))).withColumn("prod_cat",
                                                                                                           when(
                                                                                                               col("product_category") == "",
                                                                                                               "other").otherwise(
                                                                                                               col("product_category"))).drop(
    "product", "product_category")

2.3 原有逻辑输出结果

+---------------------+-----+--------+
|campaign_name        |prod |prod_cat|
+---------------------+-----+--------+
|abcloancde           |     |lending |
|abcsolcdf            |sol  |lending |
|abcemicdef           |other|other   |
|emic_estore          |other|other   |
|personalloan-emic_sol|     |lending |
|personalloan_emic    |     |lending |
|abc/emic-dg-upi:bol  |other|other   |
+---------------------+-----+--------+

3. 新增需求与性能问题

需要对prod、prod_cat为other的记录,将campaign_name按"_"拆分后,与字典表中function='match'的记录做精确匹配,优先取第一个命中的结果更新分类与产品信息。目前通过UDF二次处理后union结果的方式实现,面对数十亿条数据性能较差。

4. 特殊匹配规则

  • contains匹配时需先去除所有分隔符再比较
  • 精确匹配时保留原分隔符,仅按"_"拆分单词后比较

5. 优化诉求

请问是否可以通过case...when..then、explode、crossJoin等算子实现更高效的逻辑?预期最终输出样例如下:

+---------------------+-----+--------+
|campaign_name        |prod |prod_cat|
+---------------------+-----+--------+
|abcloancde           |     |lending |
|abcsolcdf            |sol  |lending |
|abcemicdef           |other|other   |
|emic_estore          |emic |cards   |
|personalloan-emic_sol|     |lending |
|personalloan_emic    |     |lending |
|abc/emic-dg-upi:bol  |other|other   |
+---------------------+-----+--------+

欢迎各位提供优化建议,非常感谢。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 11:24:03