Flink知识点(四)|Watermark(水位线)
在之前分享的文章Flink知识点(二)|Flink中是怎么处理乱序数据的中,提到处理乱序数据中使用Watermark(水位线),有些兄弟私聊Watermark的知识,我准备用两期时间讲下Watermark和窗口。今天我们先聊下Watermark。
一、核心概念
说到Watermark,我们首先要明白Flink的三种时间语义:
- Event Time:事件实际发生的时间(数据本身携带)
- Ingestion Time:数据进入 Flink 的时间
- Processing Time:算子处理数据的时间
比如我们中午点外卖,我选了半天,在11:20:01点了个酸菜鱼,下单。由于中午大家都在使用网络,信号不好,等到11:30:05商家才接到你的单子,11:30:10开始给你做酸菜鱼。
你点击"下单"的时间 11:20:01 → Event Time(事件时间)
商家系统收到订单的时间 11:30:05 → Ingestion Time(摄入时间)
厨师开始处理这个订单的时间 11:30:10 → Processing Time(处理时间)
Watermark只在Event Time模式下有意义。
同样上面的点外卖场景,你同事看见你点的酸菜鱼,看着不错,他也点一份,在11:29:01下了单,但他使用的某某品牌手机,比较给力,11:29:59商家就接到他的单子,11:30:04开始给他做酸菜鱼。
上面的例子可以看出,由于网络等一些问题,导致先发生Event Time的事件却偏后到达处理,这就是Event Time上的乱序问题。
这种乱序的问题,就会导致我们在使用Flink的窗口统计某一指标时,可能统计的指标不准确,那么我们怎么能准确呢?那就在等等,等到齐了在处理。但我们又不知道什么时候到齐,也不能一直等下去,我们就设置个阈值,默认达到这个阈值时,这个窗口的数据都到了,这个阈值就是Watermark。
注意:在生产环境中,只单一使用Watermark,数据极有可能丢失。所以要结合Flink知识点(二)|Flink中是怎么处理乱序数据的,这篇文章一起看。
数据流: [t=1] [t=5] [t=3] [t=8] [t=2] [t=10] ...
↑
乱序、延迟到达
Watermark: W(0) W(3) W(3) W(6) W(6) W(8) ...
当 Watermark(T) 到达某个算子时,意味着所有 event_time <= T 的数据都已到达,触发对应窗口的计算。
二、Watermark的工作原理
- Watermark是一种衡量Event Time进展的机制;
- 对于乱序数据,只使用Watermark是远远不够的,结合Flink知识点(二)|Flink中是怎么处理乱序数据的中的策略使用;
- 数据流中的Watermark(T)用于表示event_time <= T 的数据都已经到了,因此window的执行也是由Watermark触发的;
- Watermark可以理解成为一个延迟机制,我们可以设置Watermark的延时时长t。每次系统会检验已到达数据的最大Event Time,然后如果event_time <= maxEventTime - t,系统认定event_time 之前的数据都来了;当event_time = maxEventTime - t,窗口被触发。
举个🌰:
三、Watermark的使用
Watermark有两种生成方式
1️⃣ SourceFunction产生
如果用的是自定义 Source,可以在源头直接控制 Watermark,精度最高:
public class OrderSource implements SourceFunction<Order> {
private SourceContext<Order> ctx;
@Override
public void run(SourceContext<Order> ctx) throws Exception {
while (true) {
Order order = fetchFromKafka();
synchronized (ctx.getCheckpointLock()) {
ctx.collectWithTimestamp(order, order.getOrderTime());
// 直接在 source 里发 Watermark
ctx.emitWatermark(new Watermark(order.getOrderTime() - 10_000));
}
}
}
}
2️⃣ 使用assignTimestampsAndWatermarks挂载
DataStream<Order> stream = env.addSource(new OrderSource());
stream.assignTimestampsAndWatermarks(
WatermarkStrategy
.<Order>forBoundedOutOfOrderness(Duration.ofSeconds(10))
.withTimestampAssigner((order, ts) -> order.getOrderTime())
);
上面示例用的forBoundedOutOfOrderness ,是针对有界乱序的情况,最常用。对应的还有单调递增。
// 有界乱序,最常用
WatermarkStrategy
.<Order>forBoundedOutOfOrderness(Duration.ofSeconds(10))
.withTimestampAssigner((order, ts) -> order.getOrderTime());
// 单调递增(数据严格有序时用)
WatermarkStrategy
.<Order>forMonotonousTimestamps()
.withTimestampAssigner((order, ts) -> order.getOrderTime());
上面的forBoundedOutOfOrderness和forMonotonousTimestamps均为周期性生成Watermark,默认是200ms,针对于周期性生成的,我们可以扩展下,自定义一个周期性生成器:
public class PeriodicWatermarkGenerator
implements WatermarkGenerator<Order> {
private long maxTimestamp = Long.MIN_VALUE;
private final long outOfOrderness = 10_000L; // 允许10s乱序
@Override
public void onEvent(Order order, long eventTimestamp,
WatermarkOutput output) {
// 每条数据进来时,只更新最大时间戳,不发 Watermark
maxTimestamp = Math.max(maxTimestamp, eventTimestamp);
}
@Override
public void onPeriodicEmit(WatermarkOutput output) {
// 由 Flink 框架周期性调用,这里才发出 Watermark
output.emitWatermark(new Watermark(maxTimestamp - outOfOrderness - 1));
}
}
// 使用
WatermarkStrategy.forGenerator(ctx -> new PeriodicWatermarkGenerator())
.withTimestampAssigner((order, ts) -> order.getOrderTime());
周期间隔可以调整:
env.getConfig().setAutoWatermarkInterval(500); // 改成500ms
周期性的生成Watermark,由于是一定时间才生成一个Watermark,可能不满足有些特殊需求,需要逐条产生:
public class PunctuatedWatermarkGenerator
implements WatermarkGenerator<Order> {
@Override
public void onEvent(Order order, long eventTimestamp,
WatermarkOutput output) {
// 比如订单状态是"已完成"才触发 Watermark
if (order.getStatus().equals("COMPLETED")) {
output.emitWatermark(new Watermark(eventTimestamp - 1));
}
}
@Override
public void onPeriodicEmit(WatermarkOutput output) {
// 逐条模式下这里什么都不做
}
}
3️⃣ Flink SQL中产生
Flink SQL 中 Watermark 直接在建表的 DDL 里定义,不需要写代码。
CREATE TABLE orders (
order_id STRING,
user_id STRING,
amount DECIMAL(10, 2),
order_time TIMESTAMP(3), -- 事件时间字段
WATERMARK FOR order_time AS order_time - INTERVAL '10' SECOND -- 定义 Watermark
) WITH (
'connector' = 'kafka',
'topic' = 'orders',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json'
);
WATERMARK FOR order_time AS order_time - INTERVAL '10' SECOND 对应的就是 DataStream 里的 forBoundedOutOfOrderness(Duration.ofSeconds(10))。
几种常见写法
有界乱序(最常用):
WATERMARK FOR order_time AS order_time - INTERVAL '10' SECOND
单调递增(数据有序):
WATERMARK FOR order_time AS order_time - INTERVAL '0.001' SECOND
不容忍任何乱序:
WATERMARK FOR order_time AS order_time
更多推荐
所有评论(0)