Flink SQL SESSION窗口中如何不使用UDF获取指定字段的LAST_VALUE值
问题原因
Flink 1.13版本内置的LAST_VALUE聚合函数未实现会话窗口运行所需的merge方法。会话窗口需要在运行时动态合并相邻的会话分片,缺少对应merge实现就会触发你遇到的报错。
可行解决方案(均为原生SQL实现,无需自定义UDF)
方案1:使用MAX_BY聚合函数(最简方案)
MAX_BY(目标字段, 排序字段)聚合函数会返回排序字段取最大值时对应的目标字段值,刚好匹配取窗口内最大时间(即window_end对应时刻)的package值的需求,且该函数1.13版本已支持会话窗口的merge逻辑,直接替换原有LAST_VALUE即可:
SELECT user_account_id, SESSION_START(request_timestamp, INTERVAL '30' MINUTE) AS window_start, SESSION_END(request_timestamp, INTERVAL '30' MINUTE) AS window_end, MAX_BY(package, request_timestamp) AS last_package FROM 你的源表名 GROUP BY SESSION(request_timestamp, INTERVAL '30' MINUTE), user_account_id
方案2:OVER窗口+行号过滤(兼容性更强)
如果有版本兼容顾虑,可以先标记每条数据所属的会话窗口边界,再按窗口分组排序取第一条即可:
WITH window_mark AS ( -- 先给每条数据标记所属的会话窗口边界 SELECT user_account_id, package, request_timestamp, SESSION_START(request_timestamp, INTERVAL '30' MINUTE) AS window_start, SESSION_END(request_timestamp, INTERVAL '30' MINUTE) AS window_end FROM 你的源表名 ), rank_data AS ( -- 按窗口分组,按时间倒序打行号 SELECT *, ROW_NUMBER() OVER ( PARTITION BY user_account_id, window_start, window_end ORDER BY request_timestamp DESC ) AS rn FROM window_mark ) -- 取每个窗口行号为1的记录,即为窗口内最后一条数据 SELECT user_account_id, window_start, window_end, package AS last_package FROM rank_data WHERE rn = 1
注:该方案可以适配所有窗口类型,你之前对OVER窗口不适用的顾虑不存在,此处OVER窗口是基于已经确定的会话窗口边界分区,不会出现数据范围错乱的问题。
内容的提问来源于stack exchange,提问作者michalrudko
相关产品推荐
相关产品推荐

