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 端本地缓冲后批量上报,任何一个环节都足以让到达顺序偏离真实发生顺序,三者叠加之后几乎必然对不上。
  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 默认周期性发射:WatermarkGenerator 的 onPeriodicEmit 每 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。这时下游算子的事件时间时钟该怎么定?

规则很简单:算子的当前 watermark = 所有输入通道 watermark 的最小值(min)。 直觉也很直白——只要还有一个输入停在 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,整个下游的事件时间时钟被这个空闲通道死死拖住,窗口永远等不到 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. 窗口:把无界流切成有界的聚合单元

Watermark 回答的是「什么时候能出结果」,窗口回答的则是「结果按什么粒度聚合」。无界流没法直接求和求平均,必须先切成有限的窗口再聚合;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 窗口的默认 Trigger 是 EventTimeTrigger,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 给了两类根本不同的函数,差别不只在写法,而是直接决定内存占用和你能拿到哪些信息;一旦按 §5 配了 allowedLateness,同一个窗口还要被反复触发,这个选择就更不能随手做了。

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

时间语义、watermark、窗口、迟到兜底、聚合函数,零件到这里就齐了。把它们拼起来,看广告实时链路里两个最典型的需求。

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 之后 0 ~ 30 分钟内的点击,也就是业务上说的「这次点击归因于 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 深挖(开篇)· 流处理模型与运行时架构