Spark中基于x与cum_inb规则生成inb_date_assigned字段的实现求助
问题分析
你的代码存在以下关键问题:
- 列名不匹配:原始DataFrame列名为
cum_inb和inb_date,但代码中使用了cum_inb_qty和inbound_date,会导致列不存在或取值错误。 - 窗口函数逻辑错误:按
x排序的窗口中使用last函数,无法实现区间匹配逻辑,完全不符合需求。 - 未处理规则2:没有判断
x是否大于cum_inb的最大值,无法触发返回null的逻辑。
解决方案
以下提供两种高效实现方式,均能满足需求:
方式一:区间Join实现(推荐,适合大数据量)
利用Spark的SQL Join能力匹配区间,无需UDF,性能更优:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 1. 提取有效参照数据:非空的cum_inb和对应的inb_date,生成区间边界 ref_df = df.filter(F.col("cum_inb").isNotNull() & F.col("inb_date").isNotNull()) \ .select("inb_date", "cum_inb") \ .orderBy("cum_inb") \ .withColumn("prev_cum_inb", F.lag("cum_inb", 1, 0).over(Window.orderBy("cum_inb"))) # 2. 获取cum_inb的最大值,用于判断规则2 max_cum_inb = ref_df.agg(F.max("cum_inb")).first()[0] # 3. 区间Join匹配inb_date,再应用规则生成目标字段 df_result = df.join(ref_df, (F.col("x") > ref_df.prev_cum_inb) & (F.col("x") <= ref_df.cum_inb), "left") \ .withColumn("inb_date_assigned", F.when(F.col("x") == 0, F.current_date()) .when(F.col("x") > max_cum_inb, None) .otherwise(F.col("inb_date"))) \ .drop(ref_df.inb_date, "prev_cum_inb", ref_df.cum_inb) # 查看结果 df_result.show()
方式二:UDF实现(适合小数据量)
如果数据量较小,也可以用UDF直接匹配区间:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 1. 提取有效参照数据并广播,减少重复计算 ref_df = df.filter(F.col("cum_inb").isNotNull() & F.col("inb_date").isNotNull()) \ .select("inb_date", "cum_inb") \ .orderBy("cum_inb") \ .withColumn("prev_cum_inb", F.lag("cum_inb", 1, 0).over(Window.orderBy("cum_inb"))) max_cum_inb = ref_df.agg(F.max("cum_inb")).first()[0] ref_broadcast = spark.sparkContext.broadcast(ref_df.collect()) # 2. 定义UDF匹配区间 def get_assigned_date(x): if x == 0: return None # 后续用current_date处理 if x > max_cum_inb: return None for row in ref_broadcast.value: if row.prev_cum_inb < x <= row.cum_inb: return row.inb_date return None get_assigned_date_udf = F.udf(get_assigned_date) # 3. 生成目标字段 df_result = df.withColumn("inb_date_assigned", F.when(F.col("x") == 0, F.current_date()) .otherwise(get_assigned_date_udf(F.col("x")))) # 查看结果 df_result.show()
内容的提问来源于stack exchange,提问作者user23436382
相关产品推荐
相关产品推荐

