1. 题目
一个订单事件流按 customerId 聚合五分钟内的金额和订单数。设备时钟可能漂移,网络重试会造成乱序,某些 Kafka 分区还会暂时没有新消息。结果需要尽快展示,同时允许迟到数据在有限时间内更正。请设计事件时间处理方案。
2. 约束与澄清
- 先确认窗口按事件发生时间、写入时间还是处理时间计算;本题选择事件时间。
- 设定最大乱序范围,例如 30 秒,并说明超过范围后的处理方式。
- 明确结果是追加事件还是可更新快照;下游必须能识别修订版本。
- 讨论空闲分区、坏时间戳、重复事件和重放任务,避免把“没有新消息”误判成输入结束。
3. 核心概念
事件时间来自记录本身,水位线表示系统认为事件时间已经推进到某个位置。窗口通常在水位线越过窗口末端后触发;水位线之后到达、但时间戳仍落在窗口内的记录属于迟到事件。并行输入合并时,算子通常要等待各输入的最小水位线,因此一个不活跃分区可能拖慢全局进度,必须支持 idle 标记或单独的分区超时策略。
4. 参考流程
onRecord(event):
ts = extractEventTimestamp(event)
key = canonicalKey(event.customerId)
updateWatermarkGenerator(ts)
window = floorToFiveMinutes(ts)
if ts <= currentWatermark + allowedLateness:
state[window, key] = aggregate(state[window, key], event)
emitUpsert(window, key, state[window, key], revision + 1)
else:
routeToLateData(window, key, event)
onWatermark(wm):
finalizeWindowsBefore(wm)
expireStateAfterRetention()源端提取事件时间并生成有界乱序水位线;输入分区长时间无数据时标记 idle,避免它把合并水位线卡住。窗口状态先保留到允许迟到期限之后,迟到记录通过 upsert 或更正事件更新结果;超过期限的数据进入旁路流,供人工核查或离线回补。
5. 准确性与延迟取舍
允许迟到时间越长,结果越接近完整事件时间视图,但状态和更正次数也越多。只发送一次最终结果延迟较低,却无法表达迟到修订;发送带窗口键、版本号和变更原因的更新事件,要求下游具备幂等合并能力。对于金额等不可丢失指标,可以保留原始事件并安排离线重算;对于实时排行榜,则可接受超过期限后只在下一批刷新。
6. 验证与观测
- 构造有序、乱序、迟到 30 秒内和超过期限的事件,逐窗口核对结果与离线基准。
- 注入空闲分区、时钟跳变、重复事件和任务重启,检查水位线是否停滞或倒退。
- 记录当前水位线、处理时间与事件时间差、窗口状态大小、迟到率、旁路数量和更正次数。
- 对每次更新记录窗口键、版本和输入事件 ID,重放同一日志时比较最终快照是否一致。
7. 常见误区
- 用处理时间窗口代替事件时间窗口,却声称结果能抵抗网络乱序。
- 把水位线当成绝对完成标记;实际水位线可能是基于延迟假设的启发式进度。
- 忽略空闲分区,导致一个没有消息的分区阻塞所有窗口关闭。
- 迟到数据直接丢弃且没有旁路、版本或离线回补路径。
8. 面试评分点
能区分三种时间
应明确事件时间、摄取时间和处理时间的语义,并说明窗口选择如何影响准确性与延迟。
能生成并合并水位线
应解释有界乱序、并行分区最小水位线和 idle 处理,指出水位线是进度估计而非绝对承诺。
能设计迟到与更正路径
应定义允许迟到期限、旁路流、更新版本和下游幂等合并方式,覆盖期限内与期限外两类数据。
能用重放验证结论
应通过乱序、空闲、重启和重复事件测试,与离线基准及实时指标对账,而不是只观察任务是否存活。