寻求类循环的特殊SQL窗口函数实现用户渠道归因
基于会话数据的渠道关联解决方案(支持30天有效期与优先级规则)
业务规则
- 渠道优先级:
paid>organic>direct - 有效期规则:
paid:一旦出现,30天内所有后续会话均关联该渠道,直到新的paid出现或超过30天organic:仅当近30天内无paid记录时生效,生效后30天内覆盖direct会话direct:仅当近30天内无paid和organic记录时生效
目标示例结果
+─────────────+───────+──────────+─────────────+ | date | user | channel | attributed | +─────────────+───────+──────────+─────────────+ | 2022-01-01 | 123 | direct | direct | | 2022-01-14 | 123 | paid | paid | | 2022-02-01 | 123 | direct | paid | | 2022-02-12 | 123 | direct | paid | | 2022-02-13 | 123 | organic | paid | | 2022-03-08 | 123 | direct | direct | | 2022-03-10 | 123 | paid | paid | +─────────────+───────+──────────+─────────────+
现有方案问题
现有方案仅通过lag()函数对比当前行与上一行数据,无法追踪paid的30天有效期,导致2022-02-13的organic会话错误覆盖了未过期的paid关联,最终结果不符合规则。
错误方案代码:
with source as ( -- example data select cast("2022-01-01" as date) as date, 123 as user, "direct" as channel union all select "2022-01-14", 123, "paid" union all select "2022-02-01", 123, "direct" union all select "2022-02-12", 123, "direct" union all select "2022-02-13", 123, "organic" union all select "2022-03-08", 123, "direct" union all select "2022-03-10", 123, "paid" ), flag_new_channel as( -- flag sessions that would override channel informaton ; this only works statically here select *, case when lag(channel) over (partition by user order by date) is null then 1 when date_diff(date,lag(date) over (partition by user order by date),day)>30 then 1 when channel = "paid" then 1 when channel = "organic" and lag(channel) over (partition by user order by date)!='paid' then 1 else 0 end flag from source qualify flag=1 ) select s.*, f.channel attributed_channel, row_number() over (partition by s.user, s.date order by f.date desc) rn -- number of flagged previous sessions from source s left join flag_new_channel f on s.date>=f.date qualify rn=1 --only keep the last flagged session at or before the current session
错误结果:
+─────────────+───────+──────────+─────────────────────+ | date | user | channel | attributed_channel | +─────────────+───────+──────────+─────────────────────+ | 2022-01-01 | 123 | direct | direct | | 2022-01-14 | 123 | paid | paid | | 2022-02-01 | 123 | direct | paid | | 2022-02-12 | 123 | direct | paid | | 2022-02-13 | 123 | organic | organic | | 2022-03-08 | 123 | direct | organic | | 2022-03-10 | 123 | paid | paid | +─────────────+───────+──────────+─────────────────────+
可行方案(BigQuery递归CTE实现)
通过递归CTE模拟迭代计算,每行基于上一次的关联渠道及其有效期判断是否需要更新关联结果,完美适配规则:
WITH source AS ( SELECT CAST("2022-01-01" AS DATE) AS date, 123 AS user, "direct" AS channel UNION ALL SELECT "2022-01-14", 123, "paid" UNION ALL SELECT "2022-02-01", 123, "direct" UNION ALL SELECT "2022-02-12", 123, "direct" UNION ALL SELECT "2022-02-13", 123, "organic" UNION ALL SELECT "2022-03-08", 123, "direct" UNION ALL SELECT "2022-03-10", 123, "paid" ), -- 给每个用户的会话按日期排序,生成行号 ranked_sessions AS ( SELECT *, ROW_NUMBER() OVER (PARTITION BY user ORDER BY date) AS rn FROM source ), -- 递归CTE初始化:每个用户的第一行数据 recursive_attribution AS ( SELECT date, user, channel, channel AS attributed, DATE_ADD(date, INTERVAL 30 DAY) AS valid_until -- 渠道生效截止日期 FROM ranked_sessions WHERE rn = 1 UNION ALL -- 递归迭代后续行 SELECT rs.date, rs.user, rs.channel, -- 判断当前日期是否在之前关联渠道的有效期内,且之前渠道优先级更高 CASE -- 如果之前的关联渠道是paid且未过期,直接沿用 WHEN ra.attributed = 'paid' AND rs.date <= ra.valid_until THEN ra.attributed -- 如果之前是organic未过期,且当前渠道不是paid,沿用organic WHEN ra.attributed = 'organic' AND rs.date <= ra.valid_until AND rs.channel != 'paid' THEN ra.attributed -- 否则,根据当前渠道更新关联结果 ELSE rs.channel END AS attributed, -- 更新生效截止日期:如果沿用之前的渠道,保留原截止日期;否则设为当前日期+30天 CASE WHEN ra.attributed = 'paid' AND rs.date <= ra.valid_until THEN ra.valid_until WHEN ra.attributed = 'organic' AND rs.date <= ra.valid_until AND rs.channel != 'paid' THEN ra.valid_until ELSE DATE_ADD(rs.date, INTERVAL 30 DAY) END AS valid_until FROM recursive_attribution ra JOIN ranked_sessions rs ON ra.user = rs.user AND rs.rn = ra.rn + 1 ) SELECT date, user, channel, attributed FROM recursive_attribution ORDER BY user, date;
验证结果
执行上述SQL后,得到的结果与目标示例完全一致:
+─────────────+───────+──────────+─────────────+ | date | user | channel | attributed | +─────────────+───────+──────────+─────────────+ | 2022-01-01 | 123 | direct | direct | | 2022-01-14 | 123 | paid | paid | | 2022-02-01 | 123 | direct | paid | | 2022-02-12 | 123 | direct | paid | | 2022-02-13 | 123 | organic | paid | | 2022-03-08 | 123 | direct | direct | | 2022-03-10 | 123 | paid | paid | +─────────────+───────+──────────+─────────────+
方案说明
- 递归CTE逐行处理每个用户的会话,每次迭代都基于上一行计算出的
attributed和valid_until判断当前行的关联渠道 - 优先级逻辑通过
CASE语句明确实现,确保paid的最高优先级和30天有效期不会被低优先级渠道打断 - 该方案完全适配BigQuery语法,也可调整后应用于支持递归CTE的其他SQL方言(如PostgreSQL、SQL Server)
内容的提问来源于stack exchange,提问作者Thomas Handorf
相关产品推荐
相关产品推荐

