前面几篇讲了 Watermark 的原理、乱序问题处理、周期性和断点式 Watermark 生成方式以及企业级 Watermark 案例分析。Watermark 解决了什么时候触发窗口的问题但有一个场景 Watermark 本身解决不了窗口触发后迟到的数据怎么办默认情况下窗口触发后状态立即清除之后到达的迟到数据直接被丢弃。但生产环境中迟到数据是常态——网络延迟、上游重试、跨地域传输都会导致数据迟到。如果直接丢弃统计结果就会偏小业务方不接受。Flink 提供了allowedLateness()机制来解决这个问题窗口触发后状态不立即清除而是再保留一段时间allowedLateness 时长在这段时间内到达的迟到数据会重新触发窗口计算更新结果。这篇从本质定义、时间线可视化、三阶段对比、内部执行流程、状态生命周期、与 Watermark 的关系、完整代码实现、六个常见坑、八条最佳实践把 AllowedLateness 一次性讲透。一、AllowedLateness 本质定义AllowedLateness 的本质一句话概括窗口触发后状态再保留多久允许迟到数据重新触发窗口。下面这张图把 AllowedLateness 的本质定义、时间线可视化、三阶段对比放在一起展示。1.1 默认行为 vs 设置 AllowedLateness 后的行为默认情况下allowedLateness0当 Watermark 超过窗口结束时间时窗口触发计算输出结果立即清除窗口状态之后到达的迟到数据直接丢弃设置 allowedLateness(Time) 后窗口触发计算输出结果不清除状态继续保留在 allowedLateness 时长内到达的迟到数据重新触发窗口计算输出更新后的结果当 Watermark 超过窗口结束时间 allowedLateness时状态才真正清除之后到达的迟到数据设置了 sideOutputLateData 则侧输出否则丢弃1.2 核心公式和三个关键时间点核心公式状态清理时间 窗口结束时间 allowedLateness三个关键时间点窗口结束时间Watermark 超过这个时间点窗口首次触发计算窗口结束时间 allowedLatenessWatermark 超过这个时间点状态真正清除之后迟到数据侧输出sideOutputLateData或直接丢弃1.3 三大核心特点第一延迟容忍。允许数据在窗口触发后延迟到达重新触发计算更新结果。这是 AllowedLateness 最核心的价值——处理迟到数据保证结果准确性。第二状态保留。窗口触发后状态不立即清除保留 allowedLateness 时长。这是代价——状态保留期间持续占用内存Lateness 越大内存占用越久。第三重新触发。迟到数据到达时重新触发窗口函数输出更新后的结果。注意是重新输出不是更新之前的输出——下游会收到多条结果需要幂等处理或去重。二、AllowedLateness 时间线可视化以 5 分钟滚动窗口 allowedLateness(10分钟) 为例看完整的时间线10:00:00 窗口开始数据正常进入窗口累积状态 10:03:00 正常数据到达追加到窗口状态 10:05:00 窗口结束WM超过结束时间首次触发计算输出结果状态保留不清除 10:08:00 迟到数据到达事件时间10:04追加到状态重新触发计算输出更新后的结果 10:12:00 又一条迟到数据到达再次重新触发再次输出更新结果 10:15:00 WM超过 10:05 10分钟 10:15状态真正清除 10:20:00 超期数据到达事件时间10:04状态已清除侧输出到late流或丢弃关键观察10:05 首次触发后状态没有清除而是继续保留10:08 和 10:12 的迟到数据都重新触发了窗口输出了更新后的结果下游收到 3 条结果10:15 是状态清除时间点10:05 10分钟之后到达的数据不再重新触发10:20 的超期数据被侧输出或丢弃三、三阶段数据处理方式对比AllowedLateness 将数据处理分为三个阶段每个阶段的处理方式完全不同阶段时间范围处理方式状态变化输出正常数据WM 窗口结束时间数据正常进入窗口累积状态窗口状态持续累积无等待触发迟到数据窗口结束 WM 结束Lateness迟到数据追加到状态重新触发计算窗口状态更新追加迟到数据重新输出更新结果可能重复超期数据WM 结束Lateness状态已清除侧输出或丢弃无状态已清除侧输出流需设置sideOutputLateData否则丢弃重点理解第二阶段迟到数据迟到数据到达时窗口状态还在因为设置了 allowedLateness迟到数据追加到窗口状态中重新调用窗口函数计算输出更新后的结果每条迟到数据都可能触发一次重新计算迟到数据多时性能开销大下游会收到多条结果首次 每次迟到更新需要幂等处理四、AllowedLateness 内部机制理解了本质和时间线下面深入内部机制看 AllowedLateness 到底是怎么工作的。下面这张图把六步内部执行流程、状态生命周期管理、与 Watermark 的关系放在一起展示。4.1 六步内部执行流程从数据到达到状态清除AllowedLateness 的完整执行流程分为六步第一步数据到达。数据到达窗口算子提取事件时间戳判断属于哪个窗口。如果是新窗口第一条数据创建窗口状态。第二步判断窗口状态。检查窗口是否已触发、状态是否已清除。这一步决定了数据的处理方式窗口未触发 → 正常数据走第三步的正常处理窗口已触发但状态未清除在 allowedLateness 内→ 迟到数据走第三步的迟到处理窗口已触发且状态已清除超过 allowedLateness→ 超期数据侧输出或丢弃第三步正常/迟到处理。正常数据追加到窗口状态中累积迟到数据也追加到窗口状态中状态还在但会触发第四步。第四步重新触发计算。迟到数据到达时重新调用窗口函数计算输出更新后的结果。注意是重新输出不是修改之前的输出。第五步WM 推进检查。每次 Watermark 推进时检查是否有窗口超过了结束时间 allowedLateness。如果有准备清除这些窗口的状态。第六步状态清除。Watermark 超过窗口结束时间 allowedLateness 时清除窗口状态、窗口元数据、定时器。之后到达的迟到数据不再重新触发而是侧输出或丢弃。4.2 状态生命周期管理AllowedLateness 下窗口状态的完整生命周期是创建 → 累积 → 首次触发 → 保留 → 重新触发 → 清除。状态创建与累积第一条属于该窗口的数据到达时创建窗口状态。正常数据到达时追加到状态中ListState/ReducingState/AggregatingState。状态内容是窗口内的所有数据或聚合后的中间状态。首次触发与状态保留Watermark 超过窗口结束时间时首次触发窗口计算。默认行为是触发后立即清除状态但设置了 allowedLateness 后不清除状态继续保留等待可能的迟到数据。保留时长从窗口结束时间开始算保留 allowedLateness 时长。迟到数据重新触发Watermark 在(窗口结束, 结束Lateness)区间内迟到数据到达时追加到窗口状态中重新调用窗口函数计算输出更新后的结果。每条迟到数据都可能触发一次重新计算迟到数据多时性能开销大。ProcessWindowFunction 可以用context.currentWatermark() window.getEnd()判断是否是迟到触发。状态清除与超期处理Watermark 超过窗口结束时间 allowedLateness 时窗口状态、窗口元数据、定时器全部清除。之后到达的超期数据设置了 sideOutputLateData 则侧输出到单独流否则直接丢弃。状态清除后Checkpoint 中也不再包含该窗口状态重启后不会恢复。4.3 与 Watermark 的关系AllowedLateness 的所有行为都由Watermark 的推进驱动没有 WM 推进就没有窗口触发、没有状态清除。三个关键时间点都与 WM 相关首次触发时间点Watermark 窗口结束时间状态清理时间点Watermark 窗口结束时间 allowedLateness迟到数据区间窗口结束时间 Watermark 窗口结束时间 allowedLateness第一WM 驱动窗口首次触发。当 WM 超过窗口结束时间时窗口首次触发计算。这是窗口的第一次输出此时状态不清除如果设置了 allowedLateness。第二WM 决定数据是否被视为迟到。WM 本身不触发重新计算是迟到数据触发的。但 WM 决定了数据是否被视为迟到——如果数据的事件时间 当前 WM且 WM 窗口结束时间这条数据就是迟到数据。第三WM 驱动状态清除。当 WM 超过窗口结束时间 allowedLateness 时窗口状态被清除。这是最关键的时间点之后到达的迟到数据不会再重新触发。注意如果 WM 长时间不推进如空闲分区、断点式标记事件丢失、Source 消费异常窗口状态会一直保留内存占用持续增大。这也是为什么必须设置 withIdleness 的原因之一。五、AllowedLateness 完整代码实现下面是一个完整的 AllowedLateness 实现包含 allowedLateness sideOutputLateData 迟到数据侧输出补算 ProcessWindowFunction 标记迟到触发。下面这张图把完整代码实现、六个常见坑、八条最佳实践放在一起展示。importorg.apache.flink.api.common.eventtime.WatermarkStrategy;importorg.apache.flink.api.java.tuple.Tuple2;importorg.apache.flink.streaming.api.datastream.DataStream;importorg.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;importorg.apache.flink.streaming.api.environment.StreamExecutionEnvironment;importorg.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;importorg.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;importorg.apache.flink.streaming.api.windowing.time.Time;importorg.apache.flink.streaming.api.windowing.windows.TimeWindow;importorg.apache.flink.util.Collector;importorg.apache.flink.util.OutputTag;importjava.time.Duration;publicclassAllowedLatenessExample{// 订单事件publicstaticclassOrderEvent{publicStringuserId;publiclongorderId;publicdoubleamount;publiclongeventTime;publicOrderEvent(){}publicStringgetUserId(){returnuserId;}publiclonggetEventTime(){returneventTime;}}// 统计结果publicstaticclassOrderStat{publicStringuserId;publiclongcount;publicdoubletotalAmount;publicbooleanisLateTrigger;// 标记是否迟到触发publicOrderStat(){}publicOrderStat(StringuserId,longcount,doubletotalAmount,booleanisLateTrigger){this.userIduserId;this.countcount;this.totalAmounttotalAmount;this.isLateTriggerisLateTrigger;}}publicstaticvoidmain(String[]args)throwsException{StreamExecutionEnvironmentenvStreamExecutionEnvironment.getExecutionEnvironment();// 1. 定义迟到数据侧输出标签OutputTagOrderEventlateOrderTagnewOutputTagOrderEvent(late-order){};DataStreamOrderEventsourceenv.addSource(newOrderSource()).assignTimestampsAndWatermarks(WatermarkStrategy.OrderEventforBoundedOutOfOrderness(Duration.ofSeconds(30)).withIdleness(Duration.ofMinutes(1)).withTimestampAssigner((event,ts)-event.eventTime));// 2. 窗口聚合设置 allowedLateness 和 sideOutputLateDataSingleOutputStreamOperatorOrderStatresultsource.keyBy(OrderEvent::getUserId).window(TumblingEventTimeWindows.of(Time.minutes(5))).allowedLateness(Time.minutes(10))// 允许10分钟迟到.sideOutputLateData(lateOrderTag)// 超期迟到数据侧输出.process(newOrderProcessWindowFunction());// 3. 获取迟到数据侧输出流单独补算DataStreamOrderEventlateStreamresult.getSideOutput(lateOrderTag);lateStream.addSink(newLateOrderSink());// 写入补数表T1批量补算// 4. 主流结果输出包含首次触发和迟到触发的结果result.addSink(newOrderStatSink());env.execute(AllowedLateness Example);}// ProcessWindowFunction区分首次触发和迟到触发publicstaticclassOrderProcessWindowFunctionextendsProcessWindowFunctionOrderEvent,OrderStat,String,TimeWindow{Overridepublicvoidprocess(Stringkey,Contextctx,IterableOrderEventelements,CollectorOrderStatout){// 判断是否迟到触发WM超过窗口结束时间 迟到触发longcurrentWatermarkctx.currentWatermark();longwindowEndctx.window().getEnd();booleanisLateTriggercurrentWatermarkwindowEnd;// 统计longcount0;doubletotalAmount0.0;for(OrderEventevent:elements){count;totalAmountevent.amount;}// 输出结果标记是否迟到触发out.collect(newOrderStat(key,count,totalAmount,isLateTrigger));}}}代码关键点OutputTag 定义new OutputTagOrderEvent(late-order) {}注意必须用匿名类带{}因为 OutputTag 是泛型类需要类型信息。allowedLateness 设置.allowedLateness(Time.minutes(10))允许 10 分钟迟到。窗口触发后状态保留 10 分钟期间迟到数据重新触发。sideOutputLateData 设置.sideOutputLateData(lateOrderTag)超过 allowedLateness 的迟到数据侧输出到late-order标签流不直接丢弃。获取侧输出流result.getSideOutput(lateOrderTag)获取迟到数据流单独写入补数表T1 批量补算。ProcessWindowFunction 标记迟到触发用ctx.currentWatermark() ctx.window().getEnd()判断是否是迟到触发。WM 超过窗口结束时间 迟到触发WM 未超过 首次触发。下游可以根据isLateTrigger字段区分首次结果和更新结果做幂等处理。六、sideOutputLateData 配合使用sideOutputLateData 是 AllowedLateness 的最佳搭档两者配合使用形成三层迟到数据处理方案第一层有界乱序forBoundedOutOfOrderness。处理大部分正常范围内的乱序数据数据在乱序容忍时间内到达正常进入窗口触发计算。这是主力处理 95% 以上的数据。第二层allowedLateness。处理少量迟到数据窗口触发后状态保留一段时间迟到数据重新触发窗口计算更新结果。这是兜底处理 4~5% 的迟到数据。第三层sideOutputLateData。处理超过 allowedLateness 的超期数据侧输出到单独流写入补数表T1 批量补算。这是最后防线处理 1% 以内的超期数据保证不丢数据。// 三层迟到数据处理方案完整配置WatermarkStrategyOrderEventstrategyWatermarkStrategy.OrderEventforBoundedOutOfOrderness(Duration.ofSeconds(30))// 第一层30秒乱序容忍.withIdleness(Duration.ofMinutes(1));SingleOutputStreamOperatorOrderStatresultsource.assignTimestampsAndWatermarks(strategy).keyBy(OrderEvent::getUserId).window(TumblingEventTimeWindows.of(Time.minutes(5))).allowedLateness(Time.minutes(10))// 第二层10分钟迟到重新触发.sideOutputLateData(lateOrderTag)// 第三层超期数据侧输出补算.process(newOrderProcessWindowFunction());// 第三层超期数据补算result.getSideOutput(lateOrderTag).addSink(newLateOrderSink());生产环境建议三层都配置这样既能保证低延迟有界乱序正常触发又能保证结果准确allowedLateness 重新触发更新还能保证不丢数据sideOutputLateData 补算。七、六个常见坑7.1 坑一allowedLateness 设置过大导致内存溢出现象设置 allowedLateness(1小时)大促期间 TaskManager OOM作业失败。根因窗口触发后状态保留 1 小时同时有大量窗口状态保留内存占用持续增长。高吞吐 大窗口 大 Lateness 内存爆炸。假设 5 分钟窗口、10 万 QPS、每条数据 1KB1 小时 Lateness 意味着同时保留 12 个窗口的状态每个窗口 3000 万条数据状态大小约 30GB未压缩OOM 是必然的。解决方案allowedLateness 不要设置过大建议 1~10 分钟。需要更长时间补数据用 sideOutputLateData 离线补算不要靠 allowedLateness。allowedLateness 是兜底机制处理少量迟到数据不是主力。7.2 坑二迟到数据重复输出未处理现象下游统计结果偏大同一窗口数据被统计多次。根因allowedLateness 期间每条迟到数据都重新触发窗口输出更新后的结果。下游如果直接累加会导致重复统计。例如首次触发输出 count100迟到数据到达后重新触发输出 count105下游如果两条都累加结果变成 205而正确结果是 105。解决方案下游需要幂等处理按窗口 ID 覆盖更新或去重只取最新结果。在 ProcessWindowFunction 中用isLateTrigger标记是否迟到触发下游据此处理。写入支持幂等的存储如 Redis/HBase 按 key 覆盖不要写入只支持追加的存储如普通文件、Kafka 直接消费累加。7.3 坑三未设置 sideOutputLateData 导致超期数据丢失现象监控发现部分数据丢失但没有任何日志或告警。根因超过 allowedLateness 的迟到数据如果没有设置 sideOutputLateDataFlink 直接丢弃没有任何记录。生产环境数据丢失不可接受尤其是金融、计费等场景。解决方案生产环境必须设置 sideOutputLateData将超期迟到数据写入补数表T1 批量补算。监控侧输出数据量异常增加时告警可能是上游补数据或乱序容忍设置过小。7.4 坑四GlobalWindow 不支持 allowedLateness现象使用 GlobalWindow allowedLateness编译报错或运行时异常。根因GlobalWindow 没有结束时间全局窗口永远不结束除非 Trigger 触发allowedLateness 基于窗口结束时间计算状态清理时间GlobalWindow 无法计算。Flink 对 GlobalWindow 的 allowedLateness 支持有限。解决方案GlobalWindow 不要使用 allowedLateness。如果需要迟到处理改用自定义 Trigger 自定义清理逻辑或改用有界窗口滚动/滑动/会话。GlobalWindow 适合全局聚合如累计 GMV不需要事件时间窗口用 ProcessingTime 更合适。7.5 坑五会话窗口 allowedLateness 行为异常现象会话窗口设置 allowedLateness 后窗口合并行为异常迟到数据导致窗口反复合并。根因会话窗口的窗口边界是动态的由 gap 决定迟到数据可能导致窗口合并两个会话窗口之间的 gap 被迟到数据填充合并为一个窗口。而 allowedLateness 期间状态保留合并逻辑更复杂可能导致窗口反复合并、重复触发。解决方案会话窗口谨慎使用 allowedLateness。如果需要建议 gap 不要设置过大allowedLateness 不要超过 gap。测试验证迟到数据场景下的窗口合并行为确认符合预期。会话窗口本身比滚动/滑动窗口复杂加上 allowedLateness 更复杂能不用就不用。7.6 坑六WM 不推进导致状态永远不清除现象作业运行一段时间后内存持续增长Checkpoint 越来越大作业越来越慢。根因WM 长时间不推进空闲分区未设置 withIdleness、断点式标记事件丢失、Source 消费异常窗口状态永远不清除因为状态清除由 WM 推进驱动内存和 Checkpoint 持续增长。这是生产环境最隐蔽的坑——作业看起来正常但内存慢慢涨最后 OOM。解决方案必须设置 withIdleness。监控 WM Lag超过阈值告警。断点式 WM 必须配合周期性兜底。定期检查各子任务 WM 是否一致发现异常及时排查。八、八条最佳实践 Checklist上线前逐条检查allowedLateness 不要过大建议 1~10 分钟超过 30 分钟需要评估内存影响。需要更长补数据用 sideOutputLateData 离线补算不要靠 allowedLateness。必须设置 sideOutputLateData生产环境必须设置超期迟到数据写入补数表T1 批量补算。监控侧输出数据量异常增加时告警。下游幂等处理重复输出allowedLateness 期间迟到数据重新触发会重复输出下游必须幂等处理按窗口 ID 覆盖或去重只取最新。写入支持幂等的存储。配合 withIdleness 使用必须设置 withIdleness避免空闲分区导致 WM 不推进、状态永远不清除、内存持续增长。ProcessWindowFunction 标记迟到触发用context.currentWatermark() window.getEnd()判断是否迟到触发下游据此区分首次结果和更新结果做不同处理。监控状态大小和迟到数据量监控窗口状态大小Checkpoint 大小、迟到数据量、重新触发频率异常增加时告警排查。GlobalWindow 和会话窗口谨慎使用GlobalWindow 不支持 allowedLateness会话窗口迟到合并行为复杂需要充分测试验证。能用滚动/滑动窗口就不用这两个。优先用有界乱序处理大部分乱序allowedLateness 是兜底机制处理少量迟到数据。大部分乱序应该用 forBoundedOutOfOrderness 处理不要靠 allowedLateness 处理大量乱序。三层方案有界乱序主力→ allowedLateness兜底→ sideOutputLateData最后防线。九、总结与下一篇预告AllowedLateness 深度剖析要点回顾第一本质窗口触发后状态再保留 allowedLateness 时长允许迟到数据重新触发窗口计算更新结果。核心公式状态清理时间 窗口结束时间 allowedLateness。三大特点延迟容忍、状态保留、重新触发。第二时间线以 5 分钟窗口 10 分钟 Lateness 为例10:05 首次触发状态保留→ 10:08/10:12 迟到数据重新触发重复输出→ 10:15 状态清除10:0510分钟→ 10:20 超期数据侧输出或丢弃。第三三阶段对比正常数据WM 窗口结束累积状态无输出→ 迟到数据结束 WM 结束Lateness重新触发重复输出→ 超期数据WM 结束Lateness侧输出或丢弃。第四内部六步流程数据到达 → 判断窗口状态 → 正常/迟到处理 → 重新触发计算 → WM 推进检查 → 状态清除。第五状态生命周期创建 → 累积 → 首次触发 → 保留 → 重新触发 → 清除。设置 allowedLateness 后首次触发不清除状态保留到 WM 超过结束Lateness 才清除。第六与 WM 的关系所有行为由 WM 推进驱动。首次触发WM 结束、迟到区间结束 WM 结束Lateness、状态清除WM 结束Lateness。WM 不推进则状态永远不清除。第七完整代码实现OutputTag 定义 allowedLateness sideOutputLateData getSideOutput 补算 ProcessWindowFunction 标记迟到触发currentWatermark() window.getEnd()。第八三层迟到处理方案有界乱序主力处理 95%→ allowedLateness兜底处理 4~5%→ sideOutputLateData最后防线处理 1%T1 补算。第九六个常见坑Lateness 过大 OOM、重复输出未处理、未设置侧输出丢数据、GlobalWindow 不支持、会话窗口合并异常、WM 不推进状态不清除。第十八条最佳实践Lateness 不要过大、必须设置侧输出、下游幂等处理、配合 withIdleness、标记迟到触发、监控状态大小、GlobalWindow/会话窗口谨慎用、优先有界乱序处理大部分乱序。AllowedLateness 是 Flink 迟到数据处理的核心机制但它是兜底机制不是主力。生产环境的正确做法是优先用有界乱序处理大部分乱序allowedLateness 处理少量迟到sideOutputLateData 处理超期数据补算三层配合既保证低延迟又保证结果准确还保证不丢数据。
