PySpark/SQL实现按日期合并行并计算相邻日ID集合差的新增与流失计数
实现方案
以下分别给出PySpark DataFrame API和Spark SQL两种实现方式,逻辑完全匹配你的需求。
PySpark DataFrame API实现
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession(已有可跳过) spark = SparkSession.builder.appName("cal_ab_field").getOrCreate() # --------------- 构造示例源数据,你的真实数据替换这部分即可 --------------- source_data = [ ("2021-08-18", 12), ("2021-08-18", 15), ("2021-08-18", 10), ("2021-08-19", 15), ("2021-08-19", 10), ("2021-08-20", 15), ("2021-08-20", 10), ("2021-08-20", 14) ] df = spark.createDataFrame(source_data, ["date", "id"]) # ----------------------------------------------------------------------- # 1. 按日期聚合,得到每日全量id集合和id总数 daily_agg = df.groupBy("date") \ .agg( F.collect_set("id").alias("curr_ids"), F.count_distinct("id").alias("curr_id_cnt") ) # 2. 用窗口函数取上一个日期的全量id集合 date_order_win = Window.orderBy("date") daily_agg = daily_agg.withColumn("prev_day_ids", F.lag("curr_ids", 1).over(date_order_win)) # 3. 计算每日的A、B值 daily_ab = daily_agg.withColumn("A", F.when(F.col("prev_day_ids").isNull(), 0) \ .otherwise(F.size(F.array_except("prev_day_ids", "curr_ids"))) ).withColumn("B", F.when(F.col("prev_day_ids").isNull(), F.col("curr_id_cnt")) \ .otherwise(F.size(F.array_except("curr_ids", "prev_day_ids"))) ).select("date", "A", "B") # 4. 关联回原表,得到最终结果 result_df = df.join(daily_ab, on="date", how="left").orderBy("date", "id") result_df.show()
Spark SQL实现
-- 注册源数据为临时视图,你的真实表替换即可 CREATE OR REPLACE TEMP VIEW source_df AS SELECT * FROM values ('2021-08-18', 12), ('2021-08-18', 15), ('2021-08-18', 10), ('2021-08-19', 15), ('2021-08-19', 10), ('2021-08-20', 15), ('2021-08-20', 10), ('2021-08-20', 14) AS t(date, id); WITH daily_agg AS ( SELECT date, collect_set(id) AS curr_ids, count(DISTINCT id) AS curr_id_cnt, lag(collect_set(id), 1) OVER(ORDER BY date) AS prev_day_ids FROM source_df GROUP BY date ), daily_ab AS ( SELECT date, CASE WHEN prev_day_ids IS NULL THEN 0 ELSE size(array_except(prev_day_ids, curr_ids)) END AS A, CASE WHEN prev_day_ids IS NULL THEN curr_id_cnt ELSE size(array_except(curr_ids, prev_day_ids)) END AS B FROM daily_agg ) SELECT s.date, s.id, a.A, a.B FROM source_df s LEFT JOIN daily_ab a ON s.date = a.date ORDER BY s.date, s.id;
补充说明
如果你的业务需要按自然日的前一天计算(而非数据中存在的上一个日期),只需要把date字段转为日期类型,调整窗口的排序逻辑即可,当前实现完全匹配你给出的示例规则。
内容的提问来源于stack exchange,提问作者Cassius Clay
相关产品推荐
相关产品推荐

