
示例工程大数据【免费下载链接】flink-learningflink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API SQL 等内容的学习案例还有 Flink 落地应用的大型项目案例PVUV、日志存储、百亿数据实时去重、监控告警分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》项目地址https://gitcode.com/gh_mirrors/fl/flink-learning点击查看免费下载导读本文是《Flink 实战与性能优化》3.5 节的核心内容聚焦流式计算中最棘手的两个问题——事件乱序与事件延迟。文章以事件时间Event Time为背景系统讲解 Watermark水印的机制原理、Flink 中两种水印分配方式Punctuated 与 Periodic的代码实现与适用场景并结合本仓库源码给出完整可运行的实战示例最后覆盖 Window 触发后处理迟到数据的三种手段默认丢弃、allowedLateness、sideOutputLateData。读完本文你将掌握在水印机制下正确设计基于事件时间的窗口聚合任务并针对不同业务容忍度选择合适的迟到数据处理方案。为什么需要 Watermark事件乱序与事件延迟在 3.1 节 中我们讲过 Flink 中的三种时间Event Time、Ingestion Time、Processing Time及其使用场景在 3.2 节中又深入讲解了窗口机制与 Flink 自带 Window 的实现原理。这里需要明确一个关键前提如果窗口基于Processing TimeFlink 消费数据时完全不需要关心数据本身的时间因为 Processing Time 代表的是数据在 Flink 中被处理的时间这个时间是顺序递增的不存在乱序问题如果窗口基于Event Time就必须直面两个问题事件乱序与事件延迟。理想情况下Event Time 与 Process Time 相等即数据发生的时间与数据处理的时间之间没有延迟。但现实往往骨感网络的抖动、设备的故障、应用的异常等原因会导致 Process Time 总是滞后于 Event Time 一段距离。所谓乱序是指 Flink 接收到的事件的先后顺序并不是严格按照事件的 Event Time 排列的——先产生的数据可能晚到后产生的数据反而先到。依赖事件时间的典型场景有些场景特别依赖事件发生时间而非处理时间例如错误日志分析错误日志的时间戳代表着错误发生的具体时间开发者只有知道这个时间戳才能还原那个时间点系统到底发生了什么问题或者根据该时间戳去关联其他事件找出触发问题的根源设备监控设备传感器或监控系统按时间点实时上传设备周围监控情况通过监控大屏实时查看不错漏重要或可疑的事件。这类场景下最有意义的是事件发生的顺序而不是事件到达 Flink 后被处理的顺序。庆幸的是Flink 支持用户以事件时间来定义窗口也支持处理时间而为了解决乱序与延迟问题Flink 引入了Watermark 机制。3.5.1 Watermark 简介工作原理与触发机制先看一个业务例子统计 8:00 ~ 9:00 时间段打开淘宝 App 的用户数量。Flink 可以开一个窗口做聚合但由于网络抖动、应用采集发送延迟等原因无法保证在窗口结束那一刻窗口中已经收集齐 8:00 ~ 9:00 内用户打开 App 的所有事件——但又不能无限期等下去。当基于事件时间的数据流进行窗口计算时最困难的一点就是如何确定对应当前窗口的事件已经全部到达。实际上我们并不能百分百准确判断因此业界常用的做法是基于已经收集到的消息来估算是否还有消息未到达这就是 Watermark 的思想。Watermark 的定义Watermark 是一种衡量 Event Time 进展的机制它是数据本身的一个隐藏属性数据本身携带着对应的 Watermark。Watermark 本质上就是一个时间戳代表着比这个时间戳更早的事件已经全部到达窗口即假设不会再有比这时间戳更小的事件到达。这个假设是触发窗口计算的基础只有 Watermark 大于窗口对应的结束时间窗口才会关闭并进行计算。按照这个标准处理数据如果后面还有比该时间戳更小的数据到达则被视为迟到的数据——对于这部分迟到数据Flink 也有相应的机制去处理下文 3.5.7 节详述。Watermark 如何工作一个 4s 窗口的逐步推演以 Flink 从消息队列消费数据为例数据上的数字代表数据本身的 timestampW(4)、W(9)代表水印窗口是基于事件时间定义的 4s 时间窗口数据流乱序到达Flink 消费后数据1、3、2进入第一个窗口[0, 4)数据7进入第二个窗口[4, 8)随后乱序的3依旧进入第一个窗口接着水印W(4)到达水印的 timestamp 与第一个窗口结束时间一致代表后面不会再有比 4 更小的数据到达于是第一个窗口被触发计算后续数据5、6进入第二个窗口数据9进入第三个窗口当水印W(9)到达时水印比第二个窗口的结束时间8还大第二个窗口也随之触发计算以此类推。整个流程印证了窗口计算的触发条件Watermark 大于等于窗口 endTime 时窗口关闭并触发计算。3.5.2 Flink 中 Watermark 的设置方式在 Flink 中数据处理时需要调用 DataStream 的assignTimestampsAndWatermarks方法来分配时间戳与水位线。该方法有两种重载分别接收AssignerWithPeriodicWatermarks和AssignerWithPunctuatedWatermarks其底层实现如下public SingleOutputStreamOperatorT assignTimestampsAndWatermarks(AssignerWithPeriodicWatermarksT timestampAndWatermarkAssigner) { final int inputParallelism getTransformation().getParallelism(); final AssignerWithPeriodicWatermarksT cleanedAssigner clean(timestampAndWatermarkAssigner); TimestampsAndPeriodicWatermarksOperatorT operator new TimestampsAndPeriodicWatermarksOperator(cleanedAssigner); return transform(Timestamps/Watermarks, getTransformation().getOutputType(), operator).setParallelism(inputParallelism); } public SingleOutputStreamOperatorT assignTimestampsAndWatermarks(AssignerWithPunctuatedWatermarksT timestampAndWatermarkAssigner) { final int inputParallelism getTransformation().getParallelism(); final AssignerWithPunctuatedWatermarksT cleanedAssigner clean(timestampAndWatermarkAssigner); TimestampsAndPunctuatedWatermarksOperatorT operator new TimestampsAndPunctuatedWatermarksOperator(cleanedAssigner); return transform(Timestamps/Watermarks, getTransformation().getOutputType(), operator).setParallelism(inputParallelism); }由此设置 Watermark 有两条路线方式核心特征适用场景AssignerWithPunctuatedWatermarks数据流中每一个递增的 EventTime 都可能产生一个 Watermark实时性要求非常高的场景TPS 高时会大量产水印可能给下游算子带来压力需谨慎AssignerWithPeriodicWatermarks周期性一定时间间隔或达到一定记录条数产生一个 Watermark生产环境中的主流选择但必须结合时间与累积条数两个维度否则极端情况下会有很大延时生产环境通常使用周期性生成方式但 Watermark 的生成方式需要根据业务场景的不同进行不同的选择。3.5.3 Punctuated Watermark逐事件判定生成AssignerWithPunctuatedWatermarks接口包含checkAndGetNextWatermark方法该方法会在每次extractTimestamp()被调用后调用由它决定是否生成新的水印。返回的水印只有在不为 null 且时间戳大于先前返回的水印时间戳时才会被发送出去如果返回 null 或时间戳比之前小则不生成新的水印。仓库中提供了完整的自定义实现 WordPunctuatedWatermark.javapublic class WordPunctuatedWatermark implements AssignerWithPunctuatedWatermarksWord { Nullable Override public Watermark checkAndGetNextWatermark(Word lastElement, long extractedTimestamp) { return extractedTimestamp % 3 0 ? new Watermark(extractedTimestamp) : null; } Override public long extractTimestamp(Word element, long previousElementTimestamp) { return element.getTimestamp(); } }该实现以extractedTimestamp % 3 0作为生成水印的判定条件时间戳能被 3 整除的事件到达时生成一个水印其余事件不生成。其配套数据模型 Word.java 只包含word、count、timestamp三个字段时间戳直接从事件中提取。注意这种方式理论上可以为每个事件都生成一个水印但水印要参与下游计算水印过多会导致整体计算性能下降因此只适合对实时性要求极高的场景。配套的可运行示例见 Main.java其中通过env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime)声明使用事件时间并通过env.socketTextStream(localhost, 9001)读取 socket 数据Map 解析为 Word 后调用assignTimestampsAndWatermarks(new WordPunctuatedWatermark())。3.5.4 Periodic Watermark周期性生成生产环境中使用AssignerWithPeriodicWatermarks定期分配时间戳并生成水印更为普遍。仓库中的 WordPeriodicWatermark.java 给出了完整实现比书中示例多加了日志输出便于观察水印推进过程Slf4j public class WordPeriodicWatermark implements AssignerWithPeriodicWatermarksWord { private long currentTimestamp Long.MIN_VALUE; Override public long extractTimestamp(Word word, long previousElementTimestamp) { long timestamp word.getTimestamp(); currentTimestamp Math.max(timestamp, currentTimestamp); log.info(event timestamp {}, {}, CurrentWatermark {}, {}, word.getTimestamp(), DateUtil.format(word.getTimestamp(), YYYY_MM_DD_HH_MM_SS), getCurrentWatermark().getTimestamp(), DateUtil.format(getCurrentWatermark().getTimestamp(), YYYY_MM_DD_HH_MM_SS)); return word.getTimestamp(); } Nullable Override public Watermark getCurrentWatermark() { long maxTimeLag 5000; return new Watermark(currentTimestamp Long.MIN_VALUE ? Long.MIN_VALUE : currentTimestamp - maxTimeLag); } }该类实现两个方法extractTimestamp()从数据中提取 Event Time并将当前时间戳与事件时间比较取最大值后赋给currentTimestamp最后返回事件时间getCurrentWatermark()通过currentTimestamp - maxTimeLag得到水印值。其中maxTimeLag代表数据允许延迟的时间示例中long maxTimeLag 5000;表示最大允许数据延迟 5 秒。超过 5 秒之后如果还来了更早的数据Flink 会将其丢弃——因为窗口中的数据需要被触发不可能一直等待迟到的数据例如因网络问题迟迟未上传的数据而不结束计算。合理设置允许延迟时间是一门细活需要观察生产环境从数据采集、进入消息队列再到 Flink 的整个流程是否出现延迟统计平均延迟的大致波动范围。这也说明一个事实Flink 设计 Watermark 的根本目的是解决部分数据乱序或延迟问题但不能真正做到彻底解决——不过在流处理框架中这已经是非常实用的特性了。Periodic 的四个内置实现类AssignerWithPeriodicWatermarks接口有四个实现类功能与使用方式如下1. BoundedOutOfOrdernessTimestampExtractor用于发出滞后于数据时间的水印作用与上面自定义的类类似只需传入一个时间参数代表允许数据延迟到来的时间。使用方式仓库示例 Main2.java// Time.seconds(10) 代表允许延迟的时间大小 data.assignTimestampsAndWatermarks(new BoundedOutOfOrdernessTimestampExtractorWord(Time.seconds(10)) { // 重写 extractTimestamp() 抽象方法 Override public long extractTimestamp(Word element) { return element.getTimestamp(); } });2. CustomWatermarkExtractor仓库中自定义的周期性生成水印类MetricWatermark.java用于对MetricEvent等数据流生成水印配合 flink-learning-common 模块的公共模型使用。3. AscendingTimestampExtractor用于时间戳单调递增的数据流。如果数据流的时间戳不是单调递增会有专门的处理方法其核心逻辑为public final long extractTimestamp(T element, long elementPrevTimestamp) { final long newTimestamp extractAscendingTimestamp(element); if (newTimestamp this.currentTimestamp) { this.currentTimestamp newTimestamp; return newTimestamp; } else { violationHandler.handleViolation(newTimestamp, this.currentTimestamp); return newTimestamp; } }即新时间戳不小于当前时间戳时正常推进否则触发violationHandler.handleViolation(...)处理乱序违规。4. IngestionTimeExtractor依赖机器系统时间在extractTimestamp和getCurrentWatermark方法中基于System.currentTimeMillis()获取时间而不是基于事件的时间。如果该分配器在数据进入 Flink 后立即分配这个时间就与 Ingestion Time 一致因此得名 IngestionTimeExtractor。两个重要的使用注意点使用周期性方式生成水印时可以通过env.getConfig().setAutoWatermarkInterval(...)设置水印生成间隔每隔 n 毫秒。仓库示例 Main1.java 中设置了env.getConfig().setAutoWatermarkInterval(5000);表示每 5 秒生成一次水印。通常建议在数据源source之后就生成水印或者先做 filter/map/flatMap 等简单操作之后再生成水印——越早生成水印效果越好甚至可以直接在数据源头生成。例如在 source 的run()方法中Override public void run(SourceContextMyType ctx) throws Exception { while (/* condition */) { MyType next getNext(); ctx.collectWithTimestamp(next, next.getEventTimestamp()); if (next.hasWatermarkTime()) { ctx.emitWatermark(new Watermark(next.getWatermarkTime())); } } }通过ctx.collectWithTimestamp携带事件时间输出通过ctx.emitWatermark主动发射水印。3.5.5 每个 Kafka 分区的时间戳当使用 Kafka Connector 作为数据源时Flink 的 Kafka Consumer 会按分区读取数据并为每个 Kafka 分区分别分配时间戳与水印。Kafka 单分区内的消息通常是有序的因此在分区级别生成的水印质量更高当多个分区的水印汇聚到下游算子时Flink 取所有输入分区水印的最小值作为该算子的当前水印以确保不会漏掉任何分区中可能迟到的数据。这一机制可以在使用 Kafka 作为数据源的作业中通过assignTimestampsAndWatermarks覆盖 Kafka Consumer 默认的时间戳提取与水印生成逻辑参考仓库中 Kafka 相关连接器示例。3.5.6 将 Watermark 与 Window 结合起来处理延迟数据将水印与窗口结合是处理乱序、延迟数据的标准做法窗口基于 Event Time 定义窗口的触发由 Watermark 驱动。仓库中的窗口模块 flink-learning-window 也大量采用这一组合。仓库示例 Main3.java 展示了水印 时间窗口 迟到容忍的完整链路env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime); env.setParallelism(1); SingleOutputStreamOperatorWord data env.socketTextStream(localhost, 9001) .map(new MapFunctionString, Word() { Override public Word map(String value) throws Exception { String[] split value.split(,); return new Word(split[0], Integer.valueOf(split[1]), Long.valueOf(split[2])); } }).assignTimestampsAndWatermarks(new WordPeriodicWatermark()); data.keyBy(0) .timeWindow(Time.seconds(10)) .allowedLateness(Time.milliseconds(2)) .sum(1) .print(); env.execute(watermark demo);整个流程为socket 读入形如word,count,timestamp的文本 → 解析为 Word 并分配时间戳/水印允许 5s 延迟→ 按第 0 个字段分组 → 10 秒滚动窗口聚合 count → 窗口触发后额外等待 2ms 的迟到数据。完整的运行环境依赖这些示例需要本地启动nc -lk 9001之类的 socket 服务端发送数据工程位于 flink-learning-examples其依赖关系可查看模块 pom.xml。3.5.7 处理延迟数据的三种方法当窗口已经触发计算后仍可能有迟到数据到达Flink 提供三种处理手段丢弃默认默认行为。窗口触发计算后再到达的迟到数据Watermark 已越过窗口 endTime、且超出 allowedLateness 容忍范围的数据直接丢弃不参与计算也不可恢复。这种方式最简单、开销最低适用于对结果精度要求不高的场景。allowedLateness再次指定允许数据延迟的时间在窗口上调用allowedLateness(Time)可以额外指定一段允许数据延迟的时间。在窗口已经因 Watermark 触发后只要迟到的数据仍未超过窗口结束时间 allowedLateness该数据依然会被重新纳入对应窗口并触发窗口的再次计算输出更新后的聚合结果。仓库示例 Main3.java 中data.keyBy(0) .timeWindow(Time.seconds(10)) .allowedLateness(Time.milliseconds(2)) .sum(1) .print();即 10 秒窗口触发后再等待 2ms 内的迟到数据仍会参与计算。使用注意allowedLateness与 Watermark 中的maxTimeLag是两个不同的延迟容忍前者针对窗口触发之后的迟到数据后者针对水印推进本身的滞后过大的allowedLateness会让窗口迟迟无法真正清理状态、增加状态存储压力并导致结果反复更新需结合实际业务权衡。sideOutputLateData收集迟到的数据对于超出允许范围、即将被丢弃的迟到数据可以通过sideOutputLateData(OutputTag)将其路由到旁路输出流中方便后续单独存储、修复或重新计算做到迟到数据不丢失、可追溯。仓库示例 Main4.java 完整展示了用法OutputTagWord lateDataTag new OutputTagWord(late) { }; SingleOutputStreamOperatorWord data env.socketTextStream(localhost, 9001) .map(new MapFunctionString, Word() { Override public Word map(String value) throws Exception { String[] split value.split(,); return new Word(split[0], Integer.valueOf(split[1]), Long.valueOf(split[2])); } }).assignTimestampsAndWatermarks(new WordPeriodicWatermark()); SingleOutputStreamOperatorWord sum data.keyBy(0) .timeWindow(Time.seconds(10)) // .allowedLateness(Time.milliseconds(2)) .sideOutputLateData(lateDataTag) .sum(1); sum.print(); sum.getSideOutput(lateDataTag) .print(); env.execute(watermark demo);要点拆解先定义一个匿名的OutputTagWordnew OutputTagWord(late) {}泛型与数据类型一致在窗口算子链上调用.sideOutputLateData(lateDataTag)声明迟到数据输出到该 Tag主结果流通过sum.print()打印正常聚合结果迟到数据通过sum.getSideOutput(lateDataTag).print()从旁路输出流中取出打印实现主结果与迟到数据分离消费。这三种方式可以组合使用例如既设置allowedLateness容忍一小段窗口触发后的数据又用sideOutputLateData把更晚的数据保存下来做离线补救。3.5.8 小结与反思回顾本节核心脉络可以概括为问题来源选择 Event Time 就必须面对事件乱序与事件延迟而 Processing Time 天然顺序递增不需要 Watermark。Watermark 本质一种衡量 Event Time 进展的机制是数据携带的隐藏时间戳属性代表此时间戳之前的事件均已到达的假设是触发窗口计算的依据。两种分配方式Punctuated逐事件判定、实时性高、水印量大与 Periodic周期性生成、生产主流均可通过assignTimestampsAndWatermarks接入且需配合setAutoWatermarkInterval控制生成频率越早生成水印效果越好甚至可在 source 源头通过emitWatermark直接发射。内置工具BoundedOutOfOrdernessTimestampExtractor允许固定延迟、AscendingTimestampExtractor单调递增流、IngestionTimeExtractor依赖机器时钟覆盖了常见场景Kafka 场景下可按分区分配时间戳汇聚时取最小水印。迟到数据三板斧默认丢弃、allowedLateness二次容忍、sideOutputLateData旁路收集三者结合可灵活适配从粗粒度近似到精确可回溯的不同需求层次。需要反思的是Watermark 只能缓解、无法根治乱序与延迟。maxTimeLag、allowedLateness等参数都依赖对生产链路采集 → 消息队列 → Flink延迟分布的持续观测与调优脱离业务数据特征盲目设参要么造成结果不准要么造成状态膨胀。建议在实际项目中先基于监控数据统计延迟分布再据此设定水印滞后与迟到容忍参数并在上线后持续观察指标进行调整。本节完整配套源码位于仓库 flink-learning-examples 的streaming/watermark包下Word 数据模型、两种水印生成器WordPeriodicWatermark、WordPunctuatedWatermark以及 Main ~ Main4 五个可运行示例均可直接编译运行验证本文所述机制。赞分享示例工程大数据【免费下载链接】flink-learningflink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API SQL 等内容的学习案例还有 Flink 落地应用的大型项目案例PVUV、日志存储、百亿数据实时去重、监控告警分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》项目地址https://gitcode.com/gh_mirrors/fl/flink-learning点击查看免费下载相关推荐Flink Watermark机制终极指南如何轻松处理乱序数据与延迟数据Flink Watermark机制终极指南如何轻松处理乱序数据与延迟数据 在实时数据处理领域Flink作为顶级流处理框架其Watermark机制是解决乱序示例工程大数据Flink DataStream 窗口Windows完全指南Window Assigner、Trigger、Evictor 与迟到数据处理Flink DataStream 窗口Windows完全指南Window Assigner、Trigger、Evictor 与迟到数据处理 窗口Wind后端大数据流处理批处理Flink Hive Dialect 窗口函数Window Functions完全指南语法、Window Frame 与实战示例Flink Hive Dialect 窗口函数Window Functions完全指南语法、Window Frame 与实战示例 导读 本文聚焦 Flin后端大数据流处理批处理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考