SparkSQL动态时间窗口计算:统计参与者事件后3天内其他事件数
需求:统计参与者事件后未来3天内的其他事件数量
数据样本
| 参与者(performer) | 事件(event) | 事件时间(event_time) |
|---|---|---|
| A | event_a | 2022-07-21 |
| C | event_b | 2022-07-20 |
| B | event_c | 2022-07-18 |
| C | event_d | 2022-07-11 |
| A | event_e | 2022-07-12 |
| B | event_f | 2022-07-11 |
| B | event_g | 2022-07-15 |
| B | event_h | 2022-07-14 |
| B | event_i | 2022-07-13 |
期望结果
| 参与者(performer) | 事件(event) | 事件时间(event_time) | 未来3天内其他事件数(count_other_event_in_NEXT_3_days) |
|---|---|---|---|
| A | event_a | 2022-07-21 | 0 |
| C | event_b | 2022-07-20 | 0 |
| B | event_c | 2022-07-18 | 0 |
| C | event_d | 2022-07-11 | 0 |
| A | event_e | 2022-07-12 | 0 |
| B | event_f | 2022-07-11 | 2 |
| B | event_g | 2022-07-15 | 1 |
| B | event_h | 2022-07-14 | 1 |
| B | event_i | 2022-07-13 | 2 |
现有代码(统计过去3天事件数)
以下SparkSQL代码可实现过去3天的统计需求:
-- 这段SparkSQL用于统计过去3天的事件数量 select performer, event, event_time, count(event) over(partition by performer order by cast(event_time as timestamp) range between interval 3 days preceding and current row ) as cnt from table
解决方案(统计未来3天内的其他事件数)
要实现统计每位参与者在某一事件后的未来3天内其他事件数量,可以调整窗口函数的范围为当前行到未来3天,并减去当前事件本身的计数:
select performer, event, event_time, -- 统计当前事件及未来3天内的所有事件数,减去1得到其他事件的数量 (count(event) over( partition by performer order by cast(event_time as timestamp) range between current row and interval 3 days following ) - 1) as count_other_event_in_NEXT_3_days from table
逻辑说明
- 按
performer分区,确保只统计同一参与者的事件 - 按
event_time排序,确定事件的时间顺序 - 窗口范围设为
current row到interval 3 days following,覆盖当前事件及未来3天内的所有事件 - 用总数减去1,排除当前事件本身,得到未来3天内其他事件的数量
内容的提问来源于stack exchange,提问作者Applewald
相关产品推荐
相关产品推荐

