如何用单条PySpark SQL实现窗口聚合计算滚动周期去重访客数
首先将你构建的DataFrame注册为Spark临时视图,执行如下语句:
sparkDF.createOrReplaceTempView("transactions")
之后执行以下单条Spark SQL即可得到你要的统计结果:
WITH daily_unique_cust AS ( -- 先按日期+客户ID去重,排除同个客户当日多次访问的重复记录 SELECT to_date(timestamp) AS trxn_date, customer_id FROM transactions GROUP BY to_date(timestamp), customer_id ), daily_base AS ( -- 计算每日独立客户访问量 SELECT trxn_date, COUNT(customer_id) AS unique_cust_visits FROM daily_unique_cust GROUP BY trxn_date ) SELECT a.trxn_date, a.unique_cust_visits, COUNT(DISTINCT b.customer_id) AS next_7_day_visits, COUNT(DISTINCT c.customer_id) AS next_30_day_visits FROM daily_base a -- 关联当前日期往后7天(含当日)的所有去重客户 LEFT JOIN daily_unique_cust b ON b.trxn_date BETWEEN a.trxn_date AND date_add(a.trxn_date, 6) -- 关联当前日期往后30天(含当日)的所有去重客户 LEFT JOIN daily_unique_cust c ON c.trxn_date BETWEEN a.trxn_date AND date_add(a.trxn_date, 29) GROUP BY a.trxn_date, a.unique_cust_visits ORDER BY a.trxn_date
如果统计出来的窗口范围和你预期的样例不完全匹配,可以自行调整date_add的偏移天数,比如7天窗口如果要包含当日往后共7个自然日,偏移量改成7即可。
内容的提问来源于stack exchange,提问作者Aesop
相关产品推荐
相关产品推荐

