本节摘要:窗口是把无限流切成有限计算的手术刀。本节拆解滚动、滑动、会话三类窗口的分配与触发规则,讲清"允许迟到 + 侧输出流"的迟到数据安置体系,并给出窗口数与状态规模的估算方法。读完你应能为"每分钟统计最近五分钟"这类需求写出参数完全正确的实现。
水位线解决了"何时结算",窗口解决"结算什么"。工程上最容易翻车的地方其实不在概念,而在需求翻译:业务说"每五分钟看一下最近一小时的成交",这句话里藏着滑动窗口(步长五分钟、长度一小时);业务说"按分钟给成交分桶",那是滚动窗口;业务说"把用户连续操作的会话切开统计",那是会话窗口。翻译错一个词,指标口径全歪。本节先教会你翻译,再教你安置迟到,最后教你算清状态账。
滚动窗口(Tumbling):时间轴被切成首尾相接、不重不叠的格子,每条数据恰好落进一个格子。特点是每个格子独立结算、互不重叠,适合分桶报表——"每分钟的成交额"就是它。触发时点:水位线越过格子终点。
滑动窗口(Sliding):由长度(窗口多大)与步长(多久滑动一次)两个参数定义,一条数据会同时落进多个窗口。它天然带来两笔开销:一是输出频率高(每个步长都出一版结果),二是状态翻倍(一条数据被多个窗口记账)。"每分钟输出一次最近一小时的统计"是它的标准用法,也是大屏滚动指标的标准实现。
会话窗口(Session):按活动间隙切分——同一个键的数据之间间隔超过会话超时就切开新会话。它没有固定起止,每个键的会话边界由数据自己长出来,实现上是"先按每条数据开窗口、再把相邻窗口合并"。用户行为分析(一次登录内的操作序列)是它的主场。它的代价是合并逻辑让状态管理更重,且会话长度不可预知。
// 滑动窗口:长度 1 小时,步长 1 分钟——"每分钟给出一版最近一小时的成交" payments.keyBy(Payment::getShopId) .window(SlidingEventTimeWindows.of(Time.hours(1), Time.minutes(1))) .aggregate(new AmountAgg()); // 会话窗口:同一用户操作间隔超过 30 分钟即切开会话 payments.keyBy(Payment::getUserId) .window(EventTimeSessionWindows.withGap(Time.minutes(30))) .aggregate(new SessionAgg());

水位线越窗而过时窗口触发,但门不必立即锁死。第一级安置叫允许迟到(allowed lateness):触发之后窗口再保留一段时间,期间晚到的数据照常补进窗口,并再次输出更新后的结果——这就是"结果修正"的来源,大屏上的数字偶尔跳一下变准,正是它在工作。第二级是侧输出流(side output):连允许迟到的时限都过了的数据,被从主流摘出来送进侧通道,单独落盘或进补偿队列,绝不一丢了之。
// 定义侧输出标签:接住彻底迟到的数据 OutputTag<Payment> lateTag = new OutputTag<Payment>("late-payments"){}; SingleOutputStreamOperator<Stat> stats = payments .keyBy(Payment::getShopId) .window(TumblingEventTimeWindows.of(Time.minutes(1))) .allowedLateness(Time.minutes(10)) // 触发后窗口再保留十分钟收补录 .sideOutputLateData(lateTag) // 十分钟后还迟到的进侧输出 .aggregate(new AmountAgg()); // 侧输出流单独处理:落盘供对账,或回灌补偿 DataStream<Payment> lateStream = stats.getSideOutput(lateTag); lateStream.addSink(new LateDataArchiveSink());
三级安置链的分工口诀:正常触发出快报,允许迟到出修正,侧输出兜底保完整。三段的时长预算(容忍乱序多少、允许迟到多少)要写进口径文档,跟业务方当面确认——它们直接决定大屏的"变准速度"。
⚠️ 常见坑:允许迟到的窗口在时限内不会被销毁,窗口太多、允许太久会显著推高状态规模。给"长允许迟到"配小格子,或者用聚合函数替代全量保存,是常见的减负组合。
窗口类工单繁多,但万变不离三步。第一步:核对触发依据。窗口没出数先看水位线推进——水位线没过窗尾,问题在时间语义(上游停更、空闲源拖住、容忍设超大),与窗口本身无关。第二步:核对数据归属。出了数但数值不对,按"这条数据该进哪个窗口"逐条回溯——时间戳分配错(取了采集时间而非业务时间)、时区偏移(业务口径是东八区零点,窗口按 UTC 切分导致"每日窗口"整八小时的错位)、键字段取错(聚合维度与需求不符)。第三步:核对迟到安置。数值小幅度偏低且随时间跳变修正,查允许迟到的时长是否覆盖了实际乱序分布;彻底偏低且无修正,查侧输出流的量——那批数据很可能已经被判定"不可救药"。三步走完仍无头绪的窗口工单极少;大多数情况下,卡住我们的不是窗口机制,而是"业务口径与窗口参数没有逐词对齐"的翻译问题。
窗口参数纸上谈兵终觉浅——下一节把同一批乱序数据放进不同参数组合里跑一遍,让你亲眼看着窗口触发、迟到补记与侧输出分流的全过程。