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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 22:24:05