在之前分享的文章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,窗口被触发。

举个🌰:
image.png


三、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());

上面的forBoundedOutOfOrdernessforMonotonousTimestamps均为周期性生成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

Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐