Apache Flink SQL使用LAG函数时查询结果缓慢问题求助
我在AWS Kinesis Data Analytics中运行Apache Flink SQL查询读取Kinesis流数据,发现使用LAG函数时性能极差,部分场景耗时长达数分钟;但不使用LAG函数时,查询仅需毫秒级就能返回结果。
伪代码如下:
SELECT COLUMN1, lag(COLUMN1) over(PARTITION BY id ORDER BY event_time ASC) as previous_value FROM TABLE
更新内容(完整表定义):
CREATE TABLE events ( `specversion` STRING, `type` STRING, `id` STRING , event_sid as JSON_VALUE(data, '$.payload.sid'), task_sid as JSON_VALUE(data, '$.payload.task_sid'), reservation_sid as JSON_VALUE(data, '$.payload.reservation_sid'), event_type as JSON_VALUE(data, '$.payload.eventtype'), agent_name as JSON_VALUE(JSON_VALUE(data, '$.payload.worker_attributes'), '$.full_name'), queue as JSON_VALUE(data, '$.payload.task_queue_name'), direction as CASE WHEN JSON_VALUE(JSON_VALUE(data, '$.payload.task_attributes'), '$.direction') is null THEN 'Inbound' else JSON_VALUE(JSON_VALUE(data, '$.payload.task_attributes'), '$.direction') end, reservation_version as JSON_VALUE(data, '$.payload.reservation_version'), -- 业务属性 brand as JSON_VALUE(JSON_VALUE(data, '$.payload.task_attributes'), '$.conversations.case'), country as JSON_VALUE(JSON_VALUE(data, '$.payload.task_attributes'), '$.conversations.campaign'), product as JSON_VALUE(JSON_VALUE(data, '$.payload.task_attributes'), '$.conversations.content'), business_channel as JSON_VALUE(JSON_VALUE(data, '$.payload.task_attributes'), '$.conversations.initiative'), `language` as JSON_VALUE(JSON_VALUE(data, '$.payload.task_attributes'), '$.conversations.language'), handling_department_name as JSON_VALUE(JSON_VALUE(data, '$.payload.task_attributes'), '$.conversations.handling_department_name'), handling_team_name as JSON_VALUE(JSON_VALUE(data, '$.payload.task_attributes'), '$.conversations.handling_team_name'), conversation_id as JSON_VALUE(JSON_VALUE(data, '$.payload.task_attributes'), '$.conversations.conversation_id'), communication_channel as JSON_VALUE(data, '$.payload.task_channel_unique_name'), task_age_in_queue AS JSON_VALUE(data, '$.payload.task_age_in_queue' RETURNING integer), task_age AS JSON_VALUE(data, '$.payload.task_age' RETURNING integer), task_age_in_queue2 AS JSON_VALUE(data, '$.payload.task_attributes.callStartDateTime'), `time` string, `data` STRING, task_attributes as JSON_VALUE(data, '$.payload.task_attributes'), callStartDateTime as JSON_VALUE(JSON_VALUE(data, '$.payload.task_attributes'), '$.callStartDateTime'), call_start_date_time as JSON_VALUE(data, '$.payload.task_attributes.callStartDateTime'), conference as JSON_VALUE(data, '$.payload.task_attributes.conference'), segments as JSON_VALUE(JSON_VALUE(data, '$.payload.task_attributes'), '$.task_attributes.segments'), record_time AS TO_TIMESTAMP(REPLACE(left(`time`, CHAR_LENGTH(`time`) - 1), 'T', ' ')), task_assignment_status as JSON_VALUE(data, '$.payload.task_assignment_status'), proc_time AS PROCTIME(), WATERMARK FOR record_time AS record_time - INTERVAL '1' second, PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'kinesis', 'stream' = '', 'aws.region' = '', 'scan.stream.initpos' = 'AT_TIMESTAMP', 'scan.stream.initpos-timestamp' = '2023-07-29T18:18:00.000-00:00', 'format' = 'json', 'scan.stream.recordpublisher' = 'EFO', 'scan.stream.efo.consumername' = 'App', 'scan.stream.efo.registration' = 'LAZY' );
核心原因分析
1. 分组状态膨胀
LAG函数属于窗口函数,需要基于PARTITION BY id分组并按event_time排序。这意味着Flink要为每个id维护状态,存储该分组内的历史数据,直到水印推进到可清理状态的时间点。如果id基数极大(比如每个事件id唯一),状态会爆炸式增长,导致内存占用过高、状态读写IO开销剧增,直接拖慢查询速度。
2. 水印与乱序数据的延迟
表定义中水印为record_time - INTERVAL '1' second,如果Kinesis流中存在大量乱序事件,Flink会持续等待水印推进,无法及时清理窗口状态,导致状态累积。同时,基于事件时间的窗口需要等待足够时间对齐乱序数据,会让LAG结果输出延迟,表现为性能差。
3. 重复JSON解析的CPU开销
表定义中大量使用嵌套JSON_VALUE解析字段(如agent_name、brand),使用LAG函数时,Flink需要为每个分组的历史记录重复解析这些JSON字段(除非做了投影优化),高CPU开销的重复解析会进一步放大性能问题。
4. EFO初始加载的状态构建压力
使用EFO增强型扇出和LAZY注册模式,初始启动时从指定时间点AT_TIMESTAMP加载大量历史数据,LAG函数需要为这些历史数据构建完整的分组状态,初始化过程耗时极长,导致首次查询结果延迟极高。
5. 主键设置的潜在干扰
表中设置PRIMARY KEY (id) NOT ENFORCED,如果id并非真正唯一键(或流中存在重复id的事件),会导致状态维护逻辑混乱,额外增加状态处理开销。
优化建议
- 调整分组键:若
id基数过大,改用更粗粒度的字段(如conversation_id),减少分组数量,降低状态规模。 - 优化水印策略:业务允许的话,适当增大水印延迟,或改用处理时间(
proc_time)排序,避免乱序数据导致的状态累积。 - 预解析JSON:在上游Kinesis流提前解析并扁平化存储JSON数据,减少Flink侧的重复解析开销;或用物化视图预计算字段。
- 调整EFO初始加载:无需大量历史数据时,将
scan.stream.initpos改为LATEST,减少初始状态构建压力;或分批次加载历史数据。 - 监控状态指标:在Kinesis Data Analytics控制台查看作业的状态大小指标,确认是否存在状态膨胀,必要时开启状态后端的增量checkpoint优化。
内容的提问来源于stack exchange,提问作者Peter Donnelly

