Skip to content
Charles Shao
Go back

Flink 深挖 · 时间语义、窗口与 Watermark

views

开篇铺平了运行时地基;但流处理真正反直觉、也最容易在生产上出错的地方,其实是时间——一条广告点击到底算在”哪一分钟”?曝光比点击晚到了怎么办?作业挂掉重放一段数据,跑出来的 CTR 还和上次一样吗?这些问题的答案,全落在时间语义、Watermark 与窗口这三件事上。

这一篇只钻时间这一维:Flink 怎么定义时间、怎么用 Watermark 在乱序流里判断”数据到齐了没”、怎么把无界流切成有界的窗口去聚合,最后落到广告实时链路——按 event time 算 CTR、按 requestId 做曝光点击归因 join。

一句话定位:时间语义决定”结果对不对”,Watermark 决定”什么时候能出结果”,窗口决定”结果按什么粒度聚合”。 三者配齐,无界乱序流才能算出可复现、可对账的实时指标。

TL;DR

Table of contents

Open Table of contents

1. 三种时间语义:processing / ingestion / event time

Flink 处理的是无界流,每条记录其实有好几个”时间”可以谈。搞清楚它们的区别,是后面一切窗口/watermark 讨论的地基。

Event Time 与 Processing Time 对比示意图。图上方是 event time 时间轴(事件真实发生时刻,从早到晚),下方是 processing time 时间轴(数据到达 Flink 的先后顺序)。同一批六条广告事件分别标注真实发生秒 t=1s 到 t=6s。上轴按真实时间顺序排列,下轴按到达顺序排列,两轴之间用连线连接同一条事件:发生在 t=2s 的事件却排在 t=3s 之后到达、发生在 t=4s 的事件排在 t=6s 之后到达,这两条连线交叉并用红色高亮标注"乱序到达"。底部说明:processing time 只看 Flink 何时看到,受网络抖动、Kafka 分区、算子并行度、回放影响,天然乱序且不可复现;event time 是事件自带时间戳,不随处理时机变化,结果可复现、可回放、跨重跑一致。

同一批事件:上轴按真实发生时间、下轴按到达顺序,连线一交叉就是乱序。processing time 图省事但结果随处理时机漂移;event time 才是可对账的那一个。

在新版 Flink(1.12+)里,event time 已经是默认processing time 需要显式用 processing-time 窗口/timer 才启用;ingestion time 作为独立概念被淡化了。下面除非特别说明,我们讨论的都是 event time。

2. 为什么实时广告统计必须用 event time

广告实时链路的数据路径很长:用户设备上的 SDK 产生一次曝光/点击 → 上报到边缘网关 → 写入 Kafka → Flink 消费。这条路上有三个”不可控”,直接决定了不能用 processing time:

  1. 网络延迟不均:移动端弱网、海外链路、CTV(联网电视)设备回传,延迟从几十毫秒到几十秒甚至分钟级都有。用 processing time,晚到的点击会被算进错误的分钟,CTR 曲线整个错位。
  2. 天然乱序:Kafka 多分区并行、多个上报网关、SDK 端本地缓冲批量上报——到达 Flink 的顺序和事件真实发生顺序几乎必然不一致
  3. 要能回放/重跑且结果一致:作业升级、故障恢复、口径修正都需要重放历史数据。如果用 processing time,重放时”墙上时钟”完全变了,同一批数据算出的结果和线上对不上——对账直接崩。而 event time 下,重放 3 天前的数据,落进的还是 3 天前那些窗口,结果逐窗口可复现

一句话:广告是要拿去结算、对账、做归因的业务,“结果可复现”不是锦上添花而是硬需求。 只要指标要和别的系统(计费、报表、DSP 回传)对得上,就必须用 event time——把”什么时候发生”和”什么时候被处理”彻底解耦。这也是为什么下游所有窗口、join、pacing 逻辑都建立在 event time 之上。

给数据流指定 event time 时间戳与 watermark,标准写法是在 source 之后(或直接在 fromSource 里)挂一个 WatermarkStrategy

DataStream<AdEvent> events = env
    .fromSource(kafkaSource,
        WatermarkStrategy
            .<AdEvent>forBoundedOutOfOrderness(Duration.ofSeconds(30)) // 乱序容忍 30s
            .withTimestampAssigner((event, ts) -> event.getEventTimeMillis()) // 提取 event time
            .withIdleness(Duration.ofMinutes(1)),                     // 空闲分区处理,见 §3.4
        "ad-events-source");

withTimestampAssigner 告诉 Flink 从每条记录的哪个字段取 event time;forBoundedOutOfOrderness 定义能容忍多大的乱序——这两个参数,直接决定后面所有窗口的正确性。

3. Watermark:在乱序流里判断”数据到齐了没”

有了 event time,新问题立刻来了:既然数据是乱序到达的,Flink 什么时候才能确信”截止到 12:00:00 的曝光都到齐了,可以计算 11:59 那一分钟的窗口并输出结果”?无界流永远有”下一条”,你不可能等到”所有数据”,所以需要一个近似的完整性信号——这就是 Watermark。

3.1 Watermark 的本质

Watermark 是一条混在数据流里、随流下推的特殊记录,携带一个时间戳 W(t)。它的语义是一个声明

event time 已经推进到 t 了,我认为所有 event time < t 的数据都已经到达(或即将到达的可以忽略)。

一旦某个窗口的结束时间 ≤ 当前 watermark,Flink 就认为这个窗口”到齐了”,触发计算并输出。所以 Watermark 本质是一个赌注:赌”比 t 更早的数据不会再来了”。这个赌注押得:

Watermark 的调参,本质就是在”完整性”和”延迟”之间找平衡点。 广告 CTR 这种既要准(分母不能缺点击)又要快的场景,这个平衡尤其难。

3.2 生成策略与生成时机

最常用的内置策略是 forBoundedOutOfOrderness(Duration)——有界乱序假设:认为乱序不会超过某个上限。它的 watermark 公式是:

W = 当前已见的最大 event time − outOfOrderness − 1ms

减 1ms 是为了保证”恰好等于窗口边界”的边界语义正确。这个策略假设”最多晚 outOfOrderness 这么久”,超过的就是迟到数据。另一个内置策略 forMonotonousTimestamps() 用于保证有序的场景(outOfOrderness = 0),罕见。

生成时机上,Flink 的 watermark 是周期性生成的(WatermarkGeneratoronPeriodicEmit,默认每 200ms 触发一次,可用 env.getConfig().setAutoWatermarkInterval(...) 调整)。也支持标点式(punctuated,onEvent 里根据特殊事件立即发射),但周期性是绝对主流。

需要自定义时(比如按 p99 延迟动态调整乱序容忍),实现 WatermarkGenerator

public class BoundedLatenessGenerator implements WatermarkGenerator<AdEvent> {
    private final long maxOutOfOrderness = 30_000L; // 30s
    private long maxTimestamp = Long.MIN_VALUE + maxOutOfOrderness + 1;

    @Override
    public void onEvent(AdEvent event, long eventTimestamp, WatermarkOutput output) {
        maxTimestamp = Math.max(maxTimestamp, eventTimestamp); // 只更新已见最大 event time
    }

    @Override
    public void onPeriodicEmit(WatermarkOutput output) {
        // 周期性发射:W = maxEventTime - outOfOrderness - 1
        output.emitWatermark(new Watermark(maxTimestamp - maxOutOfOrderness - 1));
    }
}

注意 watermark 严格单调不减maxTimestamp 只增不减,即便来了一条超晚的旧数据,也不会把 watermark 往回拉。

3.3 多输入算子:取所有输入的最小 watermark

单条流内的 watermark 好理解,但真实作业里算子几乎总是多输入的:上游有多个并行子任务、或者做了 union/connect。此时下游算子的 event time 时钟怎么定?

规则:算子的当前 watermark = 所有输入通道 watermark 的最小值(min)。 直觉是——只要有一个输入还停在 t0,你就不能声称”< t0 的数据都到齐了”,因为那个慢输入随时可能吐出一条 < t0 的数据。所以最慢的输入决定全局进度

Watermark 生成与多输入取 min 传播示意图。左侧有两个上游 source 子任务:子任务 A 已见最大 event time = 00:16,其 watermark W_A = 16 − 3s = 00:13;子任务 B 是冷清分区,长时间无数据或 CTV 回传慢,watermark W_B = 00:10 推不动。中间是下游 KeyedProcess/窗口算子,对多个输入 watermark 做对齐取最小:min(W_A, W_B) = min(00:13, 00:10) = 00:10,因此事件时间时钟只推进到 00:10。右侧是窗口触发条件:当 Watermark ≥ 窗口末端时才输出该窗口结果。底部说明:生成策略 forBoundedOutOfOrderness 周期性(默认 200ms)发射 W = maxEventTime − outOfOrderness − 1ms,Watermark 只增不减、严格单调;多输入取 min 的代价是任何一个上游停滞(空闲分区、回传延迟大)都会拖住全局事件时间时钟、导致窗口迟迟不触发,需要 withIdleness 把空闲通道标记为 idle 排除在 min 之外。

多输入取 min:即便 A 已经推进到 00:13,只要冷清的 B 停在 00:10,下游时钟就被 B 拖在 00:10。这也是”空闲源”会卡住窗口的根因。

3.4 idle source:空闲输入把 watermark 卡死

多输入取 min 有个直接的副作用,也是生产上最常见的坑之一:如果某个输入通道长时间不来数据(某个 Kafka 分区在低峰期没消息、某个 source 子任务分到的分区是空的),它的 watermark 就一直停在原地不推进。取 min 之后,整个下游的 event time 时钟被这个空闲通道死死拖住——窗口永远等不到 watermark 越过它的末端,结果迟迟不输出

解法是 withIdleness(Duration):如果一个通道在指定时长内没有数据,就把它标记为 idle临时排除在 min 计算之外,让其他活跃通道的 watermark 能正常推进;等它再来数据时自动恢复 active。

WatermarkStrategy
    .<AdEvent>forBoundedOutOfOrderness(Duration.ofSeconds(30))
    .withTimestampAssigner((e, ts) -> e.getEventTimeMillis())
    .withIdleness(Duration.ofMinutes(1)); // 1 分钟无数据 → 标记 idle,不再拖住全局 watermark

凡是分区数 > source 并行度、或存在冷热极不均的分区/渠道,就该配 withIdleness 否则一个半夜没流量的分区,就能让你整个 CTR 大盘”卡住不更新”——而监控上看又”没报错”,极其难查。

4. 窗口:把无界流切成有界的聚合单元

无界流没法直接”求和/求平均”,必须先切成有限的窗口再聚合。Flink 的窗口机制由三个可插拔组件构成:

4.1 四种窗口类型

![三种窗口切分方式对比图。顶部是一条共享的事件时间轴(0s 到 16s,带刻度),上面用圆点标出八条曝光/点击事件,分别落在 1s、2s、3s、6s、7s、9s、13s、14s,并有竖向虚线向下贯穿三个窗口带。第一行①滚动 Tumbling,size=4s,固定不重叠、每条事件只属于一个窗口,画出四个相邻不重叠的矩形 W1[0,4)、W2[4,8)、W3[8,12)、W4[12,16)。第二行②滑动 Sliding,size=4s slide=2s,重叠、一条事件可落入多个窗口,用交错两行画出多个重叠矩形 [0,4)、[2,6)、[4,8)、[6,10)、[8,12)、[10,14)、12,16) 相互重叠。第三行③会话 Session,gap=2s,无固定边界,相邻事件间隔大于 gap 即切分出新会话,画出会话1(覆盖 1-3s 的密集事件)、会话2(覆盖 6-9s)、会话3(覆盖 13-14s)三个宽度不等的矩形。底部说明:Tumbling 适合按固定周期出报表(每分钟 CTR),Sliding 适合滑动窗口平滑趋势(近 5 分钟每 1 分钟刷新),Session 适合按用户活跃段聚合(一次浏览会话内的曝光点击);切分由 WindowAssigner 决定,触发由 Trigger 决定,Evictor 可选做窗口内元素剔除。

同一条事件流的三种切法:滚动不重叠、滑动重叠、会话按活跃间隔断开。选哪种,取决于你要回答”每个固定周期""滑动趋势”还是”每段活跃会话”的问题。

一个滚动窗口算实时 CTR 的骨架:

events
    .keyBy(AdEvent::getCampaignId)                          // 按 campaign 分组
    .window(TumblingEventTimeWindows.of(Time.minutes(1)))   // 1 分钟滚动窗口
    .aggregate(new CtrAggregate());                         // 增量聚合,见 §6

4.2 Trigger 与 Evictor

event-time 窗口的默认 TriggerEventTimeTrigger当 watermark 越过窗口末端时触发一次。你可以自定义提前触发(early firing,窗口没结束就先出个中间结果)或延迟触发。Evictor 用得少,典型场景是”窗口触发前先剔除掉超过某个数量/时间的旧元素”——但用了 Evictor 就无法增量聚合(必须缓存全量元素才能剔除),代价大,慎用。

5. 迟到数据处理:allowedLateness + 侧输出兜底

Watermark 是个赌注,赌错了——即真的有数据比 watermark 晚到——这条数据就是迟到数据(late data)。默认行为是:迟到数据被直接丢弃。对广告统计这是不可接受的(丢掉的可能是真实点击,直接让 CTR 分母缺数)。Flink 给了两道兜底。

5.1 allowedLateness:让窗口”晚点再销毁”

正常情况下 watermark 一越过窗口末端,窗口触发后状态立即清理allowedLateness(Duration) 让窗口在触发后先不销毁,再多留一段时间:在这段”额外容忍期”内,每来一条迟到数据就再触发一次窗口计算(增量更新之前已输出的结果),直到 watermark > 窗口末端 + allowedLateness 才真正销毁清理状态。

代价是窗口状态要多保留一段时间(内存/状态成本),且下游要能接受”同一个窗口结果被更新多次”(幂等写入或 upsert)。

5.2 sideOutputLateData:超出容忍也不丢

超过 allowedLateness 的更晚数据,用 sideOutputLateData(OutputTag) 把它们侧输出到一条旁路流,而不是默默丢弃。这条侧流可以:落到存储做离线修正计数报警(迟到量突然飙升往往意味着某个渠道回传出问题)、或走一条单独的补偿链路。

OutputTag<AdEvent> lateTag = new OutputTag<>("late-clicks") {};

SingleOutputStreamOperator<CtrResult> result = events
    .keyBy(AdEvent::getCampaignId)
    .window(TumblingEventTimeWindows.of(Time.minutes(1)))
    .allowedLateness(Time.minutes(5))       // 窗口末端后再容忍 5 分钟迟到,期间每来一条就再触发
    .sideOutputLateData(lateTag)            // 超过 5 分钟的更晚数据侧输出,不丢
    .aggregate(new CtrAggregate(), new CtrWindowFunction());

// 迟到侧流:计数 / 报警 / 落库做后期修正
DataStream<AdEvent> lateClicks = result.getSideOutput(lateTag);
lateClicks.map(e -> 1).keyBy(x -> 0).sum(0).addSink(new LateCountAlertSink());

迟到侧输出的计数,是排查”数据准不准”的第一手指标。 迟到量的绝对值和占比,直接告诉你 outOfOrderness 设得合不合理:迟到占比常年偏高,说明乱序容忍设小了。“迟到数据”这层信息本身,比”丢掉的那几条数据”更值钱。

6. 窗口聚合函数:增量 vs 全量

窗口”怎么算”,Flink 提供两类根本不同的函数,理解它们的差别直接关系到内存占用能拿到什么信息

6.1 ReduceFunction / AggregateFunction:增量聚合

增量聚合的关键是:窗口里不缓存原始数据,只维护一个聚合中间结果。每来一条记录,立刻和已有的中间结果合并,然后丢弃原始记录。窗口触发时直接吐出中间结果。内存占用是 O(1)(每窗口只存一个累加器),这是生产的默认选择。

public class CtrAggregate
        implements AggregateFunction<AdEvent, long[], Double> {
    // 累加器:[曝光数, 点击数]
    @Override public long[] createAccumulator() { return new long[]{0L, 0L}; }

    @Override public long[] add(AdEvent e, long[] acc) {
        if (e.getType() == EventType.IMPRESSION) acc[0]++;
        else if (e.getType() == EventType.CLICK) acc[1]++;
        return acc; // 只存两个 long,不缓存原始事件 —— O(1) 内存
    }

    @Override public Double getResult(long[] acc) {
        return acc[0] == 0 ? 0.0 : (double) acc[1] / acc[0]; // CTR = 点击 / 曝光
    }

    @Override public long[] merge(long[] a, long[] b) {
        return new long[]{a[0] + b[0], a[1] + b[1]}; // session 窗口合并时用
    }
}

6.2 ProcessWindowFunction:全量聚合

ProcessWindowFunction 拿到的是窗口内全部元素的迭代器,以及一个 Context(能拿到窗口元信息:窗口起止时间、当前 watermark、访问窗口态/全局态、侧输出)。它的能力最强——能做需要看到全量数据的计算(如中位数、精确去重、topN)——但代价是窗口内所有原始元素都要缓存在状态里,内存 O(n),大窗口 + 高吞吐容易撑爆状态。

6.3 二者组合:既省内存又拿元信息

生产上最优雅的写法是把两者组合:用 AggregateFunction增量聚合(省内存),再接一个 ProcessWindowFunction窗口元信息做最终包装。此时 ProcessWindowFunction 收到的不是全量元素,而是增量聚合后的那一个结果——两全其美:

events
    .keyBy(AdEvent::getCampaignId)
    .window(TumblingEventTimeWindows.of(Time.minutes(1)))
    .aggregate(
        new CtrAggregate(),            // 增量:只存累加器,O(1)
        new ProcessWindowFunction<Double, CtrResult, Long, TimeWindow>() {
            @Override public void process(Long campaignId, Context ctx,
                    Iterable<Double> aggResult, Collector<CtrResult> out) {
                Double ctr = aggResult.iterator().next(); // 增量聚合的唯一结果
                long windowEnd = ctx.window().getEnd();    // 拿到窗口元信息
                out.collect(new CtrResult(campaignId, windowEnd, ctr));
            }
        });

默认就用”AggregateFunction + ProcessWindowFunction”这个组合:增量聚合扛住内存与吞吐,process function 负责把窗口时间戳打进结果(下游按 windowEnd 对账、写时序库都要用)。只有当计算本身必须看全量数据(中位数、精确 distinct)时,才退回纯 ProcessWindowFunction

7. 落到 AdTech:实时 CTR 与曝光点击归因 join

把前面所有零件拼起来,看广告实时链路里两个最典型的需求。

7.1 实时 CTR:event-time 滚动窗口

CTR = 点击 / 曝光。曝光流和点击流按 event time 打时间戳、按 campaignId(或更细的 campaignId + slotId)keyBy、进 1 分钟滚动窗口、用 §6.3 的组合聚合——就是前面的代码。关键是曝光和点击都按各自的 event time 落窗,而不是”处理时刻”,这样某个渠道回传慢导致点击晚到,只要在 outOfOrderness + allowedLateness 容忍内,仍能算进正确的那一分钟,CTR 不失真。

7.2 曝光点击归因:interval join

更难的是归因:一次点击要关联到它对应的那次曝光(同一个 requestId/impressionId),才能算”这次点击是哪个创意、哪个位置带来的”。这是两条流按 key 的 join,但流 join 不能无限等——必须限定时间范围

Flink 的 interval join 正是为此设计:让点击流去 join 它 event time 前后一段区间内的曝光流。语义是”点击 join 发生在它之前 0~30 分钟的曝光”:

impressions
    .keyBy(Impression::getRequestId)
    .intervalJoin(clicks.keyBy(Click::getRequestId))
    // 点击的 event time 落在 [曝光时间, 曝光时间 + 30min] 内才算归因命中
    .between(Time.minutes(0), Time.minutes(30))
    .process(new ProcessJoinFunction<Impression, Click, Attribution>() {
        @Override public void processElement(Impression imp, Click click,
                Context ctx, Collector<Attribution> out) {
            out.collect(new Attribution(imp.getRequestId(),
                    imp.getCampaignId(), click.getClickTimeMillis()));
        }
    });

interval join 底层按 event time 缓存两条流在时间区间内的数据到状态里,watermark 推进后自动清理过期状态——所以区间越大,状态越大。区间的选取是业务归因窗口(如”点击必须在曝光后 30 分钟内”)和状态成本的权衡。

另一种是 window join:把两条流放进同一个窗口再 join(impressions.join(clicks).where(...).equalTo(...).window(...)),适合”同一时间窗内的配对”;interval join 更适合”点击相对曝光的滑动区间”这种归因语义,也是广告场景更常用的。

归因 join 对 watermark 尤其敏感:曝光和点击是两条独立的流,进入 join 算子后走的是 §3.3 的多输入取 min——任何一条流的 watermark 滞后,都会拖慢 join 的状态清理和输出。曝光通常量大、点击相对稀疏,点击流很容易出现 §3.4 的 idle 分区把 join 卡住。这条链路的 exactly-once 语义与状态一致性,见端到端 Exactly-Once;状态清理与 TTL 细节见状态管理与 Checkpoint

参考


views
Share this post on:

Previous Post
Flink 深挖 · 端到端 Exactly-Once 与两阶段提交
Next Post
Flink 深挖(开篇)· 流处理模型与运行时架构