You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.15 01:46:00