Flink SQL中LAG函数配窗口框架报错,如何获取当前与前一行数据?
问题解决:Flink SQL中LAG函数报错的修复方案
错误原因
你遇到的错误是因为LAG函数属于偏移量窗口函数,这类函数的逻辑已通过自身参数控制行偏移,无需额外指定ROWS/RANGE窗口范围子句。Calcite(Flink SQL底层引擎)会限制在LAG/LEAD这类函数中使用ROWS/RANGE,因此导致了校验报错。
正确写法
要获取当前行的前一行pk值,直接使用LAG函数的默认行为即可,无需添加ROWS BETWEEN:
SELECT LAG(pk) OVER ( PARTITION BY id ORDER BY procTime ) AS prev_data_hash FROM eoj_table;
如果需要明确指定偏移行数(此处为前1行,是LAG的默认偏移量),也可以显式传入第二个参数:
SELECT LAG(pk, 1) -- 1表示偏移1行,即前一行,该参数可省略 OVER ( PARTITION BY id ORDER BY procTime ) AS prev_data_hash FROM eoj_table;
补充说明
- LAG函数的完整语法为
LAG(column_name, offset, default_value):offset:指定向前偏移的行数,默认值为1;default_value:可选参数,当没有前一行时返回的默认值,默认返回NULL。
ROWS/RANGE子句适用于聚合类窗口函数(如SUM、COUNT)的范围控制,不适用于LAG/LEAD这类偏移量函数。
内容的提问来源于stack exchange,提问作者hitesh
相关产品推荐
相关产品推荐

