Flink(二):时间语义、Watermark 与窗口机制
Flink(二):时间语义、Watermark 与窗口机制
导语:实时计算最容易答错的一章。核心就三个问题——「用谁的时间」「怎么判断数据齐了」「数据迟到了怎么办」。本篇围绕 Event Time、Watermark 的生成/传播/对齐、allowedLateness 与侧输出三件套、窗口类型与触发器、窗口生命周期展开,共 14 题。
一、时间语义
1. Flink 有哪几种时间语义?为什么只推荐 Event Time?
答: 三种时间语义:
| 时间语义 | 含义 | 由谁产生 | 现状 |
|---|---|---|---|
| Event Time(事件时间) | 事件真实发生的时间,通常写在数据里(如日志打的 ts) | 业务/数据生产方 | 推荐,生产必用 |
| Ingestion Time(接入时间) | 事件进入 Flink Source 的时间 | Flink Source | 1.12 起已废弃(官方建议直接用 Processing Time 或 Event Time) |
| Processing Time(处理时间) | 事件被算子处理时的机器时间(System.currentTimeMillis()) | 运行机器 | 简单但不准,仅用于不关心顺序的场景 |
为什么 Event Time 是唯一正确选择:
- 结果可复现:同一天的数据重跑(补数据、回溯)结果一致;Processing Time 重跑结果会变;
- 不受反压影响:反压导致数据在管道里排队,Processing Time 下「同一批数据」会被分到不同窗口,窗口结果失真;
- 乱序场景的唯一答案:只有 Event Time + Watermark 能正确处理「事件到达顺序 ≠ 发生顺序」。
Processing Time 的三个典型问题(面试要点):
- 反压放大偏差:下游变慢 → 数据在缓冲区堆积 → 同一个事件时间段的记录被推到多个处理时间窗口;
- 无法补数:离线把昨天数据重导一次,Processing Time 会把它们算进「今天」;
- 跨时区/时钟漂移:多机房机器时间不同步会直接撕裂窗口统计。
一句话:Event Time 用「数据的语义时间」做窗口,Processing Time 用「机器的物理时间」做窗口。定时/告警类不关心顺序的场景可以用 Processing Time 省掉 Watermark 的复杂度。
2. Watermark 是什么?它解决了什么问题?
答: Watermark 是「事件时间进度」的标记,它的核心语义是:
当 Flink 观察到 watermark = T 时,意味着「事件时间 ≤ T 的数据理论上已经到齐了」——于是所有
end ≤ T的窗口可以安全触发计算。
为什么必须有它:Event Time 的窗口面临一个死循环——数据乱序 → 你不知道「还有没有更早的数据在路上」→ 就永远不敢关窗。Watermark 就是打破死循环的「进度承诺」(虽然是通过「允许迟到」来换取确定性)。
两个必须记住的结论:
- 窗口触发条件是
watermark >= window.maxTimestamp(),其中maxTimestamp = end - 1(毫秒窗口的右开区间); - Watermark 是「单调整体推进」的(单调不减),某个分区/算子的 watermark 回退会被 Flink 过滤掉——所以并行度越高,最小 watermark 越慢(木桶效应)。
3. Watermark 的生成方式有哪两种?BoundedOutOfOrderness 的延迟怎么设?
答: 两种生成方式:
| 方式 | 触发时机 | 适用 |
|---|---|---|
| Periodic(周期性,默认) | 按 pipeline.auto-watermark-interval(默认 200ms)周期产出 | 绝大多数场景 |
| Punctuated(标点式) | 每来一条数据就判断是否要发(通常靠数据里的特殊标记) | 稀疏数据流、需要精确控制 |
生产常用的两种 WatermarkStrategy:
// ① 有界乱序(最常用):允许最多 5 秒乱序
WatermarkStrategy.<Order>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner((order, ts) -> order.getEventTime())
.withIdleness(Duration.ofMinutes(1)); // 空闲分区兜底,见第 5 题
// ② 单调递增:数据保证不乱序时,直接用水位=最大事件时间,延迟最小
WatermarkStrategy.<Order>forMonotonousTimestamps()
.withTimestampAssigner((order, ts) -> order.getEventTime());BoundedOutOfOrderness 的内部逻辑:维护 maxTimestamp,产出 maxTimestamp - outOfOrdernessMillis - 1。
延迟怎么定(高频实战题):延迟 = 你能容忍的乱序上界,本质是「延迟 vs 完整性的权衡」:
| 延迟设太小 | 延迟设太大 |
|---|---|
| 迟到数据变多 → 窗口结果偏小(丢数据) | 窗口触发变慢 → 端到端延迟变大、状态驻留时间变长(内存压力) |
实操建议:
- 先量后定:统计上游数据的「到达时间 - 事件时间」分布,取 P99/P999 分位作为延迟,不要拍脑袋;
- 接受"少数迟到":把延迟设成 P99,剩下 1% 的极端乱序交给 allowedLateness + 侧输出 兜底,而不是把延迟设成 P9999(那会让所有正常数据都等很久);
- 注意单调性:
forMonotonousTimestamps只能在上游严格有序时用,否则会大量丢数据; - Watermark 是全局最慢进度:多分区(如 Kafka 多 partition)下,只要一个分区没数据,它的 watermark 不推进会卡住整个算子——这就是
withIdleness要解决的问题。
4. Watermark 在并行流中是怎么传播和对齐的?
答: Watermark 的传播规则是「广播 + 取最小」:
三条传播规则:
- 广播:上游一个子任务产生的 watermark,会发给所有下游子任务(而不是按 key 路由);
- 取最小:下游子任务的 watermark = 它所有输入通道 watermark 的最小值;
- 单调不减:过滤掉小于当前值的 watermark,防止回退。
由此推导出的重要结论:
- 并行度越高,窗口触发越慢(要等所有并行实例的进度都对齐);
- 数据倾斜会拖慢开关:某个 key 特别热导致该子任务处理慢,它的 watermark 滞后 → 整条链的窗口都晚触发;
- 多分区 Kafka + 空分区:空分区不推进 watermark,直接卡死全局 —— 用
withIdleness或对空闲分区发WatermarkStatus.IDLE解决; - Union / Connect 多流:合并后的 watermark 是各输入的最小值,因此「一条慢流会拖住整个 join」。
5. 什么是空闲分区(Idleness)问题?怎么解决?
答: 场景:Kafka topic 有 8 个 partition,其中 2 个长时间没有数据。这两个 partition 的 WatermarkGenerator 不会产出 watermark(或停在旧值),导致下游 min() 永远取到旧值 → 整个作业的窗口永不触发,一段时间后完全没有任何输出。
解决方案:withIdleness
WatermarkStrategy.<Order>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner(...)
.withIdleness(Duration.ofMinutes(1)); // 1 分钟没数据则认为该分区"空闲"原理:分区空闲超过阈值后,Flink 会把它标记为 IDLE,在计算下游 watermark 的 min 时忽略它;一旦该分区又有数据,会先发一个当前 watermark 通知下游恢复参与计算(避免「恢复瞬间把 watermark 拉回」)。
其他可选做法:
| 做法 | 说明 | 缺点 |
|---|---|---|
withIdleness | 通用解法,改了就好了 | 阈值要结合数据稀疏度调 |
| 按窗口聚合打散前先过滤空分区/合并分区 | 减少分区数 | 治标不治本 |
| 把热 & 冷数据分两个流处理 | 冷流不用关心实时开关 | 架构更复杂 |
| 改用 Processing Time | 对不关心的场景最简单 | 丢掉 Event Time 的好处 |
面试加分句:「空闲分区不推进 watermark,本质是全局 watermark 取 min 的副作用」——把「为什么」说出来比记住 API 更重要。
6. 迟到数据怎么处理?三种补救手段的区别与选择
答: 数据迟到有三种情况与三层补救:
| 层次 | 手段 | 能做什么 | 代价 |
|---|---|---|---|
| ① 减少迟到 | 增大 outOfOrderness 延迟 | 让更多迟到数据能进窗口 | 增大端到端延迟 |
| ② 窗口内补救 | allowedLateness | 窗口触发后延迟销毁,期间再来的数据重新触发窗口计算 | 每个迟到事件都会触发一次全量重新计算并输出(结果会"更新");窗口状态驻留更久 |
| ③ 窗口外补救 | sideOutputLateData(侧输出) | 超过 allowedLateness 的极致迟到数据从侧输出拿到,单独处理 | 需要自己写补数/对账逻辑 |
OutputTag<Order> lateTag = new OutputTag<Order>("late"){};
SingleOutputStreamOperator<Result> result = orders
.keyBy(Order::getUserId)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.allowedLateness(Time.seconds(30)) // ② 窗口销毁前可再更新
.sideOutputLateData(lateTag) // ③ 兜底
.aggregate(new SumAmount());
DataStream<Order> late = result.getSideOutput(lateTag); // 拿到彻底过期的数据关键结论(面试必答):
allowedLateness会让「同一个窗口」输出多次——下游必须能处理更新(幂等/主键 upsert),否则会重复累加;allowedLateness依赖状态:窗口状态要保留window end + allowedLateness这么久,设太大会吃内存/加大 Checkpoint;- 侧输出不是万能:它只是把迟到数据捞出来,业务正确性还得靠自己的补数/对账流程(落库后离线修正、或触发一次独立的补算任务);
allowedLateness默认值是 0,即窗口一触发就销毁,此时迟到的数据直接被丢弃(如果有侧输出则进侧输出);- 窗口触发 ≠ 窗口销毁:触发是「算一次并输出」,销毁是「状态被清理」,两者分别在
Trigger和WindowLifecycleListener里控制。
记忆图:
7. 为什么 Watermark 不能彻底解决乱序?
答: 因为 Watermark 是一个「取舍」,不是「魔法」:
- Watermark =
maxTimestamp - 允许乱序 - 1,它用「等待时间」换取「确定性」:等待越久越准,但延迟越高; - 总会有超出等待窗口的迟到数据,此时只能靠
allowedLateness/ 侧输出 / 离线补数; - Watermark 只能保证「不遗漏在延迟窗口内到达的数据」,不能保证「结果与离线一致」——要绝对正确必须配合补数。
推论(面试加分):
- 实时计算的"正确"通常定义为「近似正确 + 可对账」,不是「绝对正确」;
- 因此实时链路不能作为最终账,一般架构是「实时出指标给业务看,离线出账给财务用」,两者通过对账任务保持一致;
- 如果业务要求「绝对不能错」,就要换思路:不用窗口触发,而用「按数据驱动的定时补算」「离线回流修正」等方式。
二、窗口机制
8. 窗口有哪些类型?各自适用什么场景?
答: 按「分配规则」和「是否按 key」两个维度看:
① 按是否分区:
| 类型 | API | 特点 |
|---|---|---|
| Keyed Window | stream.keyBy(...).window(...) | 每个 key 一个独立窗口,可并行,最常用 |
| Non-Keyed Window | stream.windowAll(...) | 全局单窗口(并行度强制为 1),是性能瓶颈,慎用 |
② 按分配规则(四种窗口 Assigner):
| 窗口类型 | 分配器 | 语义 | 窗口是否重叠 | 场景 |
|---|---|---|---|---|
| 滚动窗口 Tumbling | TumblingEventTimeWindows.of(size) | 固定长度、首尾相接不重叠 | 否 | 每分钟订单量(最常用) |
| 滑动窗口 Sliding | SlidingEventTimeWindows.of(size, slide) | 固定长度、按 slide 滑动可重叠 | 是 | 每 30s 看最近 5 分钟的移动平均 |
| 会话窗口 Session | EventTimeSessionWindows.withGap(gap) | 按活跃间隔切分,gap 内无数据则关窗 | 否 | 用户会话、点击流会话 |
| 全局窗口 Global | GlobalWindows.create() | 永不自动关闭,必须配 Trigger | — | 需要自定义触发(如凑满 100 条) |
③ 按计量方式:
- 时间窗口(上面四种)与 计数窗口(
countWindow(size)/countWindow(size, slide),基于GlobalWindows + CountTrigger)。
四个高频细节(易错):
- 滑动窗口的
size / slide关系:slide = size就退化成滚动窗口;slide > size会丢数据(期间不覆盖),所以业务上要求slide ≤ size; - 滑动窗口的状态开销与
size/slide成正比:size=1h, slide=1s意味着同一份数据同时存在于 3600 个窗口里,状态与写入量爆炸——这是典型的「滑动窗口写崩 Flink」事故; - 会话窗口为什么不能增量聚合得那么顺:会话窗口的边界由数据决定(
newGap可能合并多个窗口),所以每个迟到数据可能触发窗口合并(merge),比滚动/滑动更重; windowAll的并行度是 1:windowAll之后所有数据汇到一个子任务,只能靠前置算子并行——所以 "全局 TopN" 场景要考虑两阶段聚合(先 keyBy 局部 TopN,再 windowAll 全局合并)。
9. 窗口的生命周期是怎样的?什么时候创建、什么时候销毁?
答: 窗口不是「预先创建好的容器」,而是随数据到达而创建、随 deadline 到期而销毁:
三个关键点:
- 窗口对象按 key 隔离:
keyBy后每个 key 都会有一份自己的窗口集合,numWindows = numKeys × 窗口数(这是「key 基数爆炸 → 状态爆炸」的根因); window.maxTimestamp() = end - 1:因为窗口是左闭右开[start, end),所以end-1毫秒已经属于该窗口;watermark >= end - 1时窗口触发(首次触发条件);- 销毁与「清理」是两个概念:
trigger决定是否输出,purge决定是否清状态。如果自定义 Trigger 只onFire不purge,窗口状态会一直留在 StateBackend 里直到过期(默认 EventTimeTrigger 会 purge,所以allowedLateness场景下需要显式组合ContinuousEventTimeTrigger+purge或改用WindowLifecycleListener逻辑)。
10. 窗口的 Trigger(触发器)有哪些?怎么自定义?
答: Trigger 决定「什么时候对窗口求值并输出」,核心回调有五个:
| 回调 | 触发时机 |
|---|---|
onElement | 每来一条数据 |
onEventTime | Event Time 定时器触发(watermark 推进到窗口边界) |
onProcessingTime | Processing Time 定时器触发 |
onMerge | 会话窗口等需要合并窗口时 |
clear | 窗口销毁时清理定时器与状态 |
内置 Trigger:
| Trigger | 触发条件 | 说明 |
|---|---|---|
EventTimeTrigger | watermark >= window.maxTimestamp() | 默认(Event Time 窗口使用) |
ProcessingTimeTrigger | Processing Time 到达窗口边界 | Processing Time 窗口的默认 |
ContinuousEventTimeTrigger | 每个 interval 触发一次(如每 10s) | 用于「周期性输出中间结果」 |
ContinuousProcessingTimeTrigger | 同上,Processing Time 版 | |
CountTrigger | 元素数达到阈值 | countWindow 用它 |
PurgingTrigger | 包装器,触发后立即 purge | 需要「触发后即清状态」时包一层 |
DeltaTrigger | 元素值变化量达到阈值 | 少用(如轨迹抽稀) |
TriggerResult 四种返回:
CONTINUE:不触发,继续等;FIRE_AND_PURGE:触发计算并清理状态(默认EventTimeTrigger的行为);FIRE:触发计算,保留状态(多次触发场景,如ContinuousEventTimeTrigger);PURGE:不触发,只清状态。
自定义 Trigger 典型场景:「窗口内数据超过 1000 条就立刻输出,同时 5 秒超时也输出」——组合 CountTrigger(onElement 返回 FIRE)+ 事件时间定时器(onEventTime 返回 FIRE),并注意 clear() 里清理定时器防止状态泄漏。
高频追问:「
allowedLateness和CountTrigger能一起用吗?」 可以,但要注意allowedLateness依赖的是定时器的重复注册:allowedLateness > 0时EventTimeTrigger会在窗口最后一次触发时清理;如果自定义 Trigger 没有正确处理,可能造成窗口状态永不过期。
11. 窗口聚合:reduce / aggregate / process 有什么区别?
答: 三者的本质区别是「状态里存什么」:
| 方式 | 状态中保存 | 触发时计算量 | 适用 |
|---|---|---|---|
reduce(ReduceFunction) | 一个累积值(前后类型相同) | O(1) 合并 | 累加、求最大值(类型一致) |
aggregate(AggregateFunction) | 一个累加器(Accumulator) | O(1) 合并 | 输入输出类型可不同(如 Order → Double 求平均/去重),最推荐 |
process(ProcessWindowFunction) | 窗口内所有元素(Iterable) | O(n) 遍历 | 需要整个窗口的上下文(全量排序、TopN) |
reduce/aggregate + ProcessWindowFunction | 累积值 + 窗口元信息 | 增量 + 补充上下文 | 最佳实践:既增量又拿到 window 信息 |
为什么推荐「增量聚合」:
- 状态里只存一个累加器(几十字节),而不是整个窗口的元素;
- 触发时不需要遍历,直接从累加器取结果;
- 对 Checkpoint 和内存都友好——这是「大窗口不 OOM」的关键。
// 增量聚合 + 补充窗口信息(生产推荐形态)
orders.keyBy(Order::getUserId)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.aggregate(
// ① AggregateFunction:增量累加,状态只存一个累加器
new AggregateFunction<Order, SumAcc, Double>() {
public SumAcc createAccumulator() { return new SumAcc(); }
public SumAcc add(Order o, SumAcc acc) { acc.sum += o.getAmount(); acc.cnt++; return acc; }
public Double getResult(SumAcc acc) { return acc.sum; }
public SumAcc merge(SumAcc a, SumAcc b) { ... }
},
// ② ProcessWindowFunction:拿到 window 起止时间等上下文,但输入已是聚合结果
new ProcessWindowFunction<Double, Result, String, TimeWindow>() {
public void process(String key, Context ctx, Iterable<Double> in, Collector<Result> out) {
out.collect(new Result(key, in.iterator().next(),
ctx.window().getStart(), ctx.window().getEnd()));
}
});ProcessWindowFunction 的真实代价(面试要点):
- 它把窗口内所有元素缓存在
ListState中,窗口越大、元素越多,状态越大; - 只有在真的需要「看到全部元素」时才用它(TopN、去重、排序);
- 另一个用途是
ProcessWindowFunction的Context提供windowState()/globalState()/sideOutput(),需要这些能力时也用它。
12. 窗口 Join 和 Interval Join 有什么区别?怎么做双流关联?
答: 双流关联有三个层次,代价依次升高:
| 方式 | 语义 | 状态开销 | 适用 |
|---|---|---|---|
window join(join) | 两条流同一窗口内按 key 关联(Tumbling/Sliding 窗口) | 每边各缓存窗口内数据 | 两边都实时、窗口对齐的场景 |
intervalJoin | 一条流的每条数据与另一条流事件时间区间内的数据关联(如 b.ts ∈ [a.ts-5min, a.ts+2min]) | 只需保留区间内的数据,可清理 | 最常见的双流关联(如「订单 + 支付」在 5 分钟内匹配) |
coGroup | 类似 join,但只要有一侧有数据就会输出(含单边数据) | 同上 | 需要处理「未匹配」的场景 |
// Interval Join:订单发生后 5 分钟内到达的支付记录
orders.keyBy(Order::getOrderId)
.intervalJoin(payments.keyBy(Payment::getOrderId))
.between(Time.minutes(-5), Time.minutes(0)) // 支付时间 ∈ [订单时间-5min, 订单时间]
.process(new ProcessJoinFunction<Order, Payment, Joined>() { ... });关键要点:
intervalJoin的between下界为负、上界可为正,表示「另一条流的时间可以在主流的哪一侧」;intervalJoin能回收状态:超出区间的数据可以清理,而window join必须等窗口关闭——所以能用 intervalJoin 就用它;join只能处理「窗口内」,如果两条流的事件时间偏移超过一个窗口就永远匹配不上;- 两边都要
keyBy同一个 key,否则无法共置(同一 key 落到同一子任务); - 水位线取两者最小值:慢流会拖住快流,所以双流 join 一定要考虑
withIdleness与延迟设置。
与「维表关联」的区别:维表关联是 Temporal Join / Lookup Join(流 × 外部表),属于流维 join,见《Flink(五)》。
13. 为什么窗口结果会「输出多次」?下游该怎么处理?
答: 输出多次的三个来源:
allowedLateness > 0:迟到数据会重新触发窗口 → 同一窗口多次onFire;ContinuousEventTimeTrigger/ContinuousProcessingTimeTrigger:周期性输出中间结果(每一次FIRE都是一次输出);- 故障恢复重放:Checkpoint 回滚后,从上次 Checkpoint 到失败点之间的窗口可能被重算并重新输出。
下游的正确姿势(这是"实时数据对不上"的第一大原因):
| 方案 | 说明 |
|---|---|
| 幂等 + 主键 upsert(推荐) | 输出带业务主键(如 userId + windowStart),下游按主键 upsert,重复输出是"覆盖"而非"累加" |
| 输出"全量结果"而不是"增量" | 窗口每次输出的是该窗口的完整聚合值,而不是 +1 这样的增量,天然可覆盖 |
| 去重表 / 版本号 | 下游按 (windowStart, windowEnd, key) 去重,保留最新版本 |
避免 allowedLateness + 可累加下游 | 如果下游只会 +=,那 allowedLateness 就是灾难 |
结论:窗口的下游天然要能"重放"。这与 Kafka/DB 的 Exactly-Once 能力配合(事务提交)才能形成端到端一致;否则只能依赖幂等。
14. 综合实战:设计一个「每分钟 UV/PV,允许 5 秒乱序,迟到 1 分钟内可修正」的作业
答: 把本篇串起来,标准答案如下:
需求拆解:
· UV/PV → 需要 Global 去重能力(不能用简单 sum)→ 用 ProcessWindowFunction + Set 状态
· 5 秒乱序 → forBoundedOutOfOrderness(5s)
· 1 分钟内可修正 → allowedLateness(1min),输出带窗口主键供下游覆盖
· 极端迟到 → 侧输出,落库做离线补算OutputTag<Log> lateTag = new OutputTag<>("late"){};
SingleOutputStreamOperator<UvPv> res = logs
// ① 时间语义 + 允许乱序
.assignTimestampsAndWatermarks(
WatermarkStrategy.<Log>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner((l, ts) -> l.getTs())
.withIdleness(Duration.ofMinutes(1)))
.keyBy(Log::getPageId)
// ② 窗口:1 分钟滚动
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
// ③ 迟到修正窗口 + 侧输出兜底
.allowedLateness(Time.minutes(1))
.sideOutputLateData(lateTag)
// ④ UV 需要去重 → 只能看到全部元素,用 ProcessWindowFunction
.process(new ProcessWindowFunction<Log, UvPv, String, TimeWindow>() {
private transient ValueState<Set<String>> seen; // 跨窗口复用去重(或按天 TTL)
public void process(String key, Context ctx, Iterable<Log> in, Collector<UvPv> out) {
Set<String> set = new HashSet<>();
long pv = 0;
for (Log l : in) { set.add(l.getUserId()); pv++; }
// 输出"全量值 + 窗口主键",下游 upsert 覆盖(应对重复输出)
out.collect(new UvPv(key, ctx.window().getStart(), ctx.window().getEnd(), pv, set.size()));
}
});面试时要主动说明的取舍(这是区分深度的地方):
| 决策点 | 选择 | 理由 / 代价 |
|---|---|---|
| 时间语义 | Event Time | 补数可重算、抗反压 |
| 乱序延迟 | 5s(P99) | 再大就牺牲端到端延迟;剩余乱序交 allowedLateness |
| 迟到修正 | allowedLateness 1min + 全量输出 + 主键 upsert | 输出可覆盖,下游幂等;状态多留 1 分钟 |
| 极端迟到 | 侧输出 + 离线补算 | 实时无法保证绝对正确 |
| UV 统计 | 窗口内 Set | 内存 = 窗口内 distinct 用户数,key 基数大时需用 Bitmap/HLL 近似(如 RoaringBitmap、HyperLogLog) |
| 大 key 倾斜 | 两阶段或加盐 | 见《Flink(四)》 |
加分收尾:真实生产里 UV 常做两级统计——窗口内先做分桶去重(Bitmap)得到「分钟级基数」,再按天做不可合并估算的修正(或用 HLL 做可合并估算),否则一个页面几百万用户会让窗口状态直接爆掉。
下一篇:《Flink(三)》进入状态管理与容错——Keyed/Operator State 的类型与选择、StateBackend 选型、Checkpoint 的 Barrier 对齐与异步快照原理、非对齐 Checkpoint、端到端 Exactly-Once、重启策略与 State TTL。
