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

如何在AWS Kinesis Data Analytics中实现全量1小时滑动窗口热门统计

解决AWS Kinesis Data Analytics中滑动窗口输出全量历史物品的问题

我完全懂你的痛点——现在的实现只能输出最近一分钟有新事件的物品统计,但你想要的是过去1小时内所有出现过的物品,哪怕它们最近没动静,也要一直输出直到最后一次事件滑出窗口。问题核心出在滑动窗口的触发逻辑上:只有当TUMBLE_STREAM有新数据流入时,KDA才会计算对应item的窗口值,没新数据的item就被“遗忘”了。

咱们换个思路,核心是要强制每分钟触发一次全量计算,把过去1小时内的所有item都拉出来统计。这里给你一套优化后的SQL方案:

步骤1:生成每分钟的时间触发流

首先创建一个每分钟生成一条时间戳记录的流,作为全局触发信号,确保每分钟都会执行一次统计:

CREATE OR REPLACE STREAM "MINUTE_TICK_STREAM" ("tick_time" TIMESTAMP);
CREATE OR REPLACE PUMP "MINUTE_TICK_PUMP" AS
INSERT INTO "MINUTE_TICK_STREAM"
SELECT STREAM FLOOR(CURRENT_TIMESTAMP TO MINUTE) AS tick_time
FROM "SOURCE_SQL_STREAM_001"
GROUP BY FLOOR(CURRENT_TIMESTAMP TO MINUTE);

如果你的源流可能长时间没有数据,可以改用KDA的SYSTEM_STREAM来生成时间戳,保证触发的稳定性。

步骤2:修改翻滚窗口聚合流

给TUMBLE_STREAM加上窗口时间字段,方便后续关联匹配:

CREATE OR REPLACE STREAM "TUMBLE_STREAM" ("item_id_count" bigint, "item_id" bigint, "window_time" TIMESTAMP);
CREATE OR REPLACE PUMP "TUMBLE_PUMP" AS
INSERT INTO "TUMBLE_STREAM"
SELECT STREAM 
    COUNT(*) AS item_id_count, 
    "item_id",
    FLOOR("SOURCE_SQL_STREAM_001".ROWTIME TO MINUTE) AS window_time
FROM "SOURCE_SQL_STREAM_001"
GROUP BY FLOOR("SOURCE_SQL_STREAM_001".ROWTIME TO MINUTE), "item_id";

步骤3:关联触发流与历史item,确保全量覆盖

把时间触发流和翻滚聚合流做关联,同时拉取过去1小时内所有出现过的item,保证没有遗漏:

CREATE OR REPLACE STREAM "JOINED_STREAM" ("tick_time" TIMESTAMP, "item_id" bigint, "current_minute_count" bigint);
CREATE OR REPLACE PUMP "JOIN_PUMP" AS
INSERT INTO "JOINED_STREAM"
SELECT STREAM
    t.tick_time,
    COALESCE(ts.item_id, prev.item_id) AS item_id,
    COALESCE(ts.item_id_count, 0) AS current_minute_count
FROM "MINUTE_TICK_STREAM" t
-- 左关联当前分钟的新数据
LEFT JOIN "TUMBLE_STREAM" ts 
    ON t.tick_time = ts.window_time
-- 全关联过去1小时内的所有历史item,避免遗漏
FULL JOIN (
    SELECT DISTINCT item_id 
    FROM "TUMBLE_STREAM" 
    WHERE window_time >= t.tick_time - INTERVAL '1' HOUR
) prev ON COALESCE(ts.item_id, prev.item_id) = prev.item_id;

步骤4:计算1小时滑动窗口总和并输出

最后基于关联后的数据集,计算每个item过去1小时的累计计数:

CREATE OR REPLACE STREAM "DESTINATION_SQL_STREAM" ("count" bigint, "item_id" bigint);
CREATE OR REPLACE PUMP "SLIDE_PUMP" AS
INSERT INTO "DESTINATION_SQL_STREAM"
SELECT STREAM
    SUM(current_minute_count) OVER (
        PARTITION BY item_id
        ORDER BY tick_time
        RANGE INTERVAL '1' HOUR PRECEDING
    ) AS "count",
    item_id
FROM "JOINED_STREAM"
ORDER BY tick_time, "count" DESC;

关键逻辑说明

  • 强制触发:MINUTE_TICK_STREAM保证每分钟都会触发一次统计,不管源流有没有新数据
  • 全量覆盖:通过FULL JOIN拉取过去1小时的所有item,哪怕它们最近没有新事件,也会被纳入计算
  • 滑动窗口累加:用OVER窗口函数对每个item的分钟计数做1小时范围的累加,得到你想要的实时统计

这样修改后,你就能得到期望的输出:每分钟都会输出过去1小时内所有出现过的item,直到它们的最后一次事件时间滑出1小时窗口为止。比如在17:42的输出里,item1和item2虽然没有新事件,依然会保留它们的计数,直到1小时后从窗口中移除。

注意事项

  • 如果你的item数量极大,FULL JOIN可能会有性能压力,可以考虑用KDA的状态存储来跟踪活跃item列表,优化关联逻辑
  • 可以根据业务需求调整水印(Watermark)设置,避免延迟数据影响统计准确性

内容的提问来源于stack exchange,提问作者WhiteHotLoveTiger

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 07:09:41