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

