开篇把运行时地基铺平了,但流处理最反直觉、也最容易在生产上翻车的地方其实是时间:一条广告点击到底该算进哪一分钟?曝光比点击晚到了怎么办?作业挂掉后重放一段数据,跑出来的 CTR 还和上次一样吗?这些问题的答案,全压在时间语义、Watermark 与窗口这三件事上。
这一篇只钻时间这一维——Flink 怎么定义时间、怎么靠 Watermark 在乱序流里判断「数据到齐了没」、又怎么把无界流切成有界窗口去聚合,最后一路落到广告实时链路:按 event time 算 CTR、按 requestId 做曝光点击归因 join。
一句话定位:时间语义决定「结果对不对」,Watermark 决定「什么时候能出结果」,窗口决定「结果按什么粒度聚合」,三者配齐,无界乱序流才能算出可复现、可对账的实时指标。
TL;DR
- 三种时间语义:processing time(算子处理时刻,最快但不可复现)、ingestion time(进入 Flink 的时刻,折中)、event time(事件自带时间戳,可复现、可回放、可对账)——实时统计几乎只该用 event time。
- 广告统计为什么非 event time 不可:曝光与点击经 SDK 上报、网关、Kafka 多跳,天然乱序且延迟不均;只有按事件真实发生时刻聚合,重放与重跑才能得到一致的结果,否则对账永远对不平。
- Watermark 的本质:一个随流下推的时间戳
W(t),声明「event time 已推进到 t,认为 < t 的数据基本到齐了」——它是完整性与延迟之间的一个赌注。 - 生成策略:
WatermarkStrategy.forBoundedOutOfOrderness(Duration)最常用,周期性(默认 200ms)发射W = maxEventTime − outOfOrderness − 1ms,且 watermark 严格单调不减。 - 多输入取 min:算子有多个输入通道(上游多并行度、
union、connect)时,其事件时间时钟等于所有输入 watermark 的最小值,最慢的输入决定全局进度。 - 窗口四类:Tumbling(滚动、不重叠)、Sliding(滑动、重叠)、Session(会话、按活跃间隔切)、Global(全局、靠自定义 Trigger);切分归
WindowAssigner、触发归Trigger、可选剔除归Evictor。 - 迟到数据两道兜底:
allowedLateness让窗口晚点再销毁、迟到数据一到就再触发一次增量更新;sideOutputLateData把超出容忍的数据侧输出,不丢,留作后期修正。 - 聚合函数两条路线:
ReduceFunction/AggregateFunction是增量聚合,来一条算一条、只存中间结果,省内存;ProcessWindowFunction是全量聚合,能拿窗口元信息但要缓存整窗数据;生产上常把两者组合使用。 - AdTech 落地:曝光与点击按 event-time 窗口算实时 CTR,按
requestId做曝光点击归因,用 interval join(曝光 join 其后 30 分钟内的点击)或 window join。 - idle source 是隐形杀手:某个 Kafka 分区或 source 子任务长期无数据,会让 watermark 推不动、窗口迟迟不触发;用
withIdleness(Duration)把空闲通道排除出 min 计算。
Table of contents
Open Table of contents
1. 三种时间语义:processing / ingestion / event time
Flink 处理的是无界流,而一条记录身上其实挂着好几个「时间」,谈「什么时候」之前得先说清楚谈的是哪一个——这是后面所有窗口与 watermark 讨论的地基。
- Processing Time(处理时间):算子处理这条记录那一刻的机器墙上时钟。它最简单、延迟也最低,不用提取时间戳、不用 watermark,系统时间一到窗口就触发;代价是结果不可复现,同一批数据今天跑和明天跑、快机器和慢机器、正常跑和重放,落进的窗口都可能不同。更麻烦的是它还受处理速度牵连:作业一旦积压,被追着处理的数据会整片挤进同一个处理时间窗口。
- Ingestion Time(摄入时间):记录进入 Flink source 那一刻、由 source 算子赋予的时间戳。它比 processing time 稳定,下游各算子看到的是同一个摄入时间戳,不会各自去取墙上时钟;但它仍然无法反映事件真实发生的时刻,更处理不了数据在进入 Flink 之前就已经乱序这件事,实践中用得很少。
- Event Time(事件时间):事件真实发生时就打进数据里的时间戳(曝光发生的毫秒、点击发生的毫秒)。它是数据的固有属性,不随处理时机改变,无论什么时候、在哪台机器上、重放多少次,同一条事件永远落进同一个窗口。代价是必须正面处理乱序——真实世界的事件不会排好队再到达,而这正是 watermark 存在的理由。
同一批事件:上轴按真实发生时间、下轴按到达顺序,连线一交叉就是乱序。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 的死刑。
- 网络延迟不均:移动端弱网、海外链路、CTV(联网电视)设备回传,延迟从几十毫秒到几十秒乃至分钟级都有。用 processing time,晚到的点击会被算进错误的分钟,CTR 曲线整条错位。
- 天然乱序:Kafka 多分区并行消费、多个上报网关、SDK 端本地缓冲后批量上报,任何一个环节都足以让到达顺序偏离真实发生顺序,三者叠加之后几乎必然对不上。
- 要能回放重跑且结果一致:作业升级、故障恢复、口径修正都要重放历史数据。用 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 更早的数据不会再来」。这个赌注押得:
- 太紧(
outOfOrderness太小)→ watermark 推得快、窗口触发快、延迟低,但晚到的数据会被判为迟到、算不进窗口,结果偏少、不准。 - 太松(
outOfOrderness太大)→ 等得久、几乎不丢数据,但结果延迟高,实时性差。
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 的数据都到齐了」,因为那个慢输入随时可能再吐出一条更早的记录。换句话说,最慢的输入决定全局进度。
多输入取 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 把窗口机制拆成三个可插拔组件,职责分得很清楚:
WindowAssigner(窗口分配器):决定每条记录属于哪个(或哪些)窗口,这是各类窗口的本质区别所在。Trigger(触发器):决定窗口什么时候触发计算并输出,event-time 窗口默认要等 §3 那套 watermark 越过窗口末端才触发。Evictor(剔除器,可选):在触发计算前后从窗口里剔除部分元素,例如只保留最近 N 条。
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 可选做窗口内元素剔除。
同一条事件流的三种切法:滚动不重叠、滑动重叠、会话按活跃间隔断开。选哪种,取决于你要回答的是「每个固定周期」「滑动趋势」还是「每段活跃会话」。
- Tumbling(滚动窗口):固定长度、互不重叠,每条事件只落进一个窗口,也是最常用的一类——「每 1 分钟的曝光数与 CTR」就是它,
TumblingEventTimeWindows.of(Time.minutes(1))。 - Sliding(滑动窗口):固定长度加固定滑动步长,窗口之间重叠,一条事件可能同时落进多个窗口,适合「近 5 分钟 CTR、每 1 分钟刷新一次」这类平滑趋势,
SlidingEventTimeWindows.of(Time.minutes(5), Time.minutes(1))。需要留意的是 size/slide 比值越大,一条数据被复制进的窗口就越多,状态与计算成本成比例放大。 - Session(会话窗口):没有固定边界,由活跃间隔决定——相邻两条事件间隔超过
gap就切出一个新会话,适合「一次浏览会话内的行为聚合」,EventTimeSessionWindows.withGap(Time.minutes(30))。它还会随新数据动态合并:两个原本分开的会话,中间来了一条事件把它们连起来就并成一个,实现上比前两者复杂不少。 - Global(全局窗口):把同一 key 的所有数据分进一个永不结束的窗口,必须配自定义
Trigger才有意义,否则永远不触发;典型用法是「每攒够 100 条触发一次」这类计数触发逻辑(CountTrigger)。
一个用滚动窗口算实时 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),这也是生产上的默认选择。
ReduceFunction:输入、输出与中间结果同类型,适合求和、取最大这类同类型归并。AggregateFunction:更通用,输入类型、累加器类型与输出类型可以各不相同。算 CTR 正好用得上——累加器里同时攒曝光数和点击数,输出时相除:
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。
参考
- Apache Flink. Timely Stream Processing(时间语义与 Watermark 概念):Flink 官方对 event/processing time 与 watermark 的权威定义。
- Apache Flink. Generating Watermarks(WatermarkStrategy / forBoundedOutOfOrderness / withIdleness):watermark 生成策略、周期生成、空闲源处理的 API 文档。
- Apache Flink. Windows(WindowAssigner / Trigger / Evictor / allowedLateness / side output):窗口类型、聚合函数与迟到数据处理的官方指南。
- Tyler Akidau. Streaming 101: The world beyond batch:event time vs processing time、乱序与延迟的经典入门。
- Tyler Akidau. Streaming 102: The world beyond batch:watermark、窗口触发、迟到数据与完整性的进阶论述。
- Akidau et al. The Dataflow Model:event-time 窗口与 watermark 的理论源头(Google Dataflow / Beam 模型)。