PySpark按长度优先匹配关联字段出现AnalysisException如何解决
问题原因
你直接在df1的列运算中引用了不属于df1的df2列terms和product_category,两个DataFrame未建立关联的情况下Spark无法解析这两个列的来源,因此抛出属性缺失的分析异常。同时原写法也无法实现「优先匹配长度最长term」的业务规则。
正确实现代码
import pyspark.sql.functions as F from pyspark.sql.window import Window # 1. 给df1添加唯一行ID,避免重复行分组出错 df1 = df1.withColumn("row_id", F.monotonically_increasing_id()) # 2. 预处理df2:按term字符串长度降序排序,生成优先级序号(序号越小优先级越高) df2_ranked = df2.withColumn("term_len", F.length(F.col("terms"))) \ .orderBy(F.col("term_len").desc()) \ .withColumn("priority", F.monotonically_increasing_id()) # 3. 左关联两个DataFrame,关联条件为campaign_name包含terms df_joined = df1.join(df2_ranked, df1.campaign_name.contains(df2_ranked.terms), how="left") # 4. 用窗口函数取每个df1行优先级最高的匹配结果,无匹配的字段填充为other window_spec = Window.partitionBy("row_id").orderBy("priority") df_result = df_joined.withColumn("rn", F.row_number().over(window_spec)) \ .filter(F.col("rn") == 1) \ .fillna("other", subset=["product_category", "product"]) \ .drop("row_id", "term_len", "priority", "rn", "terms")
运行df_result.show()即可得到符合你业务规则的输出结果。
内容的提问来源于stack exchange,提问作者whatsinthename
相关产品推荐
相关产品推荐

