开篇铺平了运行时地基;但流处理真正反直觉、也最容易在生产上出错的地方,其实是时间——一条广告点击到底算在”哪一分钟”?曝光比点击晚到了怎么办?作业挂掉重放一段数据,跑出来的 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)时,其 event time 时钟 = 所有输入 watermark 的最小值——最慢的输入决定全局进度。
- 窗口四类:Tumbling(滚动、不重叠)、Sliding(滑动、重叠)、Session(会话、按活跃间隔切)、Global(全局、靠自定义 Trigger);切分由
WindowAssigner、触发由Trigger、可选剔除由Evictor负责。 - 迟到数据两道兜底:
allowedLateness让窗口”晚点再销毁、迟到数据到了就再触发一次增量更新”;sideOutputLateData把超出容忍的迟到数据侧输出,不丢、留作后期修正。 - 聚合函数两条路线:
ReduceFunction/AggregateFunction是增量聚合(来一条算一条、只存中间结果,省内存);ProcessWindowFunction是全量(能拿窗口元信息但要缓存整窗数据);生产上常把两者组合,兼顾省内存与拿元信息。 - AdTech 落地:曝光/点击按 event-time 窗口算实时 CTR;按
requestId做曝光点击归因用 interval join(点击 join 前 N 分钟的曝光)或 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 端本地缓冲批量上报——到达 Flink 的顺序和事件真实发生顺序几乎必然不一致。
- 要能回放/重跑且结果一致:作业升级、故障恢复、口径修正都需要重放历史数据。如果用 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。此时下游算子的 event time 时钟怎么定?
规则:算子的当前 watermark = 所有输入通道 watermark 的最小值(min)。 直觉是——只要有一个输入还停在 t0,你就不能声称”< t0 的数据都到齐了”,因为那个慢输入随时可能吐出一条 < t0 的数据。所以最慢的输入决定全局进度。
多输入取 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 的窗口机制由三个可插拔组件构成:
WindowAssigner(窗口分配器):决定每条记录属于哪个(些)窗口——这是窗口类型的本质区别。Trigger(触发器):决定窗口什么时候触发计算并输出(event-time 窗口默认在 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 提供两类根本不同的函数,理解它们的差别直接关系到内存占用和能拿到什么信息。
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
把前面所有零件拼起来,看广告实时链路里两个最典型的需求。
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。
参考
- 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 模型)。