题干与适用场景
一个电商平台把订单事件写入可重放的持久日志,同时保留不可变原始数据。稳态吞吐未给出,峰值为每秒 2 万条,促销时可能短时升至 10 倍。库存异常要在事件产生后 10 秒内告警;运营看板要在 5 分钟内更新;财务在次日关账,要求同一业务日可重跑、可对账。约 2% 的事件会迟到,最长 6 小时。
以上吞吐、突发倍数、迟到比例和时限都是题目假设,不代表某个引擎的性能承诺。假设每条事件有稳定的 eventid、orderid、event_time 和 schema 版本;财务口径使用明确时区的业务日。团队有四名数据工程师,现有仓库和调度器能可靠运行批任务,但还没有成熟的有状态流作业值班体系。
问题不要求为三个消费者统一选择一种模式。目标是从行动时限、结果完整性、恢复方式和维护成本推导最小足够方案。本题归入 data,因为核心能力是数据处理语义、时效与数据质量取舍,而不是搭建跨业务域的平台架构。
面试官考察点
第一项信号是候选人是否从业务决定的最后时限倒推,而不是看到消息队列就默认流处理。源是无界事件流,只说明数据持续到达;它不决定每个下游必须逐条处理。财务次日关账可以读取同一日志的有界快照,用批处理得到更容易复算的结果。
第二项是能否区分处理延迟和行动延迟。引擎在 500 毫秒内完成计算,但接收器每 5 分钟刷新一次,看板仍不可能达到秒级。强回答会把预算拆成采集、排队、计算、写入、缓存与告警投递,并为每段设指标。
第三项是正确性边界。流引擎的“正好一次”处理不能保证迟到数据已经完整,也不能自动保护外部副作用。批任务可重跑也不等于安全;覆盖范围、业务键、快照版本和原子发布不清楚时,重跑同样会重复或混合口径。
最后看运维判断。连续流处理需要长期容量、积压恢复、状态、检查点、发布与值班能力。若 5 分钟时限能由微批稳定满足,引入逐条有状态计算只会增加故障面。反过来,10 秒告警如果真会触发补货或熔断,按小时跑批就失去业务价值。
回答前需要澄清的问题
- 时限从哪里开始,到哪里结束? 本题从事件在源系统提交成功开始,到告警送达或查询结果可见为止。若只承诺“进入数仓”,架构会低估接收器和缓存延迟。
- 10 秒告警会触发自动动作吗? 自动扣减、停卖或通知需要去重、幂等和审计;仅供观察的告警可以接受更高误报或重复率。
- 5 分钟看板需要最新估算还是完整数字? 若允许带“截至时间”的可修订估算,微批足够;若每次展示都必须包含最长 6 小时的迟到事件,5 分钟新鲜度与完整性要求互相冲突,必须改产品契约。
- 财务关账何时冻结,之后能否调整? 次日生成初版、迟到后出调整分录,与等待 6 小时再冻结是两种口径。它们决定批任务的截点和修订协议。
- 三个用途能否共享原始层与变换定义? 可以共享规范化事件、业务键和测试样本;若分别复制过滤、时区和金额规则,批流结果迟早漂移。
- 团队能否全天候维护状态流作业? 没有积压告警、检查点恢复演练和安全发布能力时,应把连续流范围限制在确实需要秒级动作的最小路径。
30 秒回答框架
“我先按业务行动时限拆分消费者,不按数据来源一刀切。10 秒库存异常使用连续流处理;5 分钟运营看板用一到两分钟微批,留出写入和缓存预算;次日财务关账使用按业务日快照的批处理,并作为对账真源。三条路径共享不可变原始事件、规范化规则和业务键。
流结果都是可修订视图,不宣称迟到数据已经完整;外部动作按事件与规则版本幂等。上线前以同一历史区间回放批、微批和流结果,影子运行比较数量、金额与尾延迟,再只为秒级路径启用动作。选择标准是满足行动 SLO 的最简单模式,并且能证明恢复后的结果收敛。”
分步骤深入解答
第一步:把三个需求写成结果契约
先为每个输出写清楚键、时限、完整性和修订方式,而不是先选 Spark、Flink 或其他产品。
| 用途 | 结果键 | 可见时限 | 完整性与修订 | 推荐模式 |
|---|---|---|---|---|
| 库存异常 | 商品、规则版本、时间窗 | 10 秒 | 可快速修订;动作必须幂等 | 连续流 |
| 运营看板 | 指标、维度、窗口 | 5 分钟 | 显示截至时间;迟到后覆盖新版本 | 1–2 分钟微批 |
| 财务关账 | 业务日、账户、币种 | 次日 | 按冻结快照重算;差异走调整流程 | 批处理 |
同一事件可服务三份不同契约。消息队列只是输入传输层,不能把表中三行强行变成同一种处理方式。
第二步:从行动时限分配端到端预算
10 秒告警可以先分配一个可测预算:源提交与传输 2 秒,排队 2 秒,计算 2 秒,接收器与规则动作 2 秒,告警投递 2 秒。实际数字要通过压测修正,但总和不能超过业务时限。每段都记录 p95、p99 与最大积压年龄;只看算子耗时会遗漏慢接收器。
5 分钟看板可每 1 或 2 分钟启动微批。假设使用 2 分钟切片,调度等待最多约 2 分钟,计算和写入各预留 1 分钟,缓存刷新预留 1 分钟。若 10 倍突发让计算超过预算,先增加并行度、降低切片或延后非关键维度;不能靠平均耗时宣称达标。
财务的主要约束是可复算和口径冻结,不是最低延迟。批任务读取明确的输入快照或 offset 范围,把结果写入带 run_id 的暂存区,校验通过后原子发布。重跑相同输入和规则版本应得到相同结果。
第三步:让三条路径共享事实,不复制实现债务
原始事件首先进入可重放日志或对象存储,保留原始 payload、schema 版本、摄取时间与来源位置。规范化层统一处理 schema 演进、时区、金额单位、取消状态和业务键。之后再由三种模式读取。
共享事实不要求三个引擎逐行代码相同。更实用的边界是共享数据契约、黄金样本和确定性规则;若批与流必须分别实现窗口逻辑,就用同一输入样本做等价测试。财务还可以保留更严格的独立校验,不能因为复用代码而失去制衡。
第四步:分别设计重复、迟到与外部副作用
库存流按 event_id 去重,按事件时间计算短窗口,输出包含规则版本和结果版本。发送停卖、补货或通知时使用稳定的动作幂等键;检查点成功不代表外部 HTTP 调用只发生一次。失败重试必须能识别已执行动作。
运营微批按固定的半开区间读取来源 offset,每次对目标窗口做版本化覆盖,而不是追加一个新的总数。迟到事件进入后续批次并提升结果版本。看板显示“数据截至某时”,让用户知道五分钟新鲜度不等于六小时尾部已完整。
财务批处理等待约定截点后读取冻结输入。超过截点的事件进入下一次调整运行,保留原运行、输入范围、规则版本和差异审批。流式正好一次可以保证某些持久化处理结果不重复,却不能证明迟到数据完整,也不能为外部副作用自动提供相同语义。
第五步:把容量与成本比较成可复算问题
题设峰值每秒 2 万条;10 倍突发就是每秒 20 万条。若每条规范化事件按 1 KiB 粗估,入口峰值约为每秒 195 MiB。这个数只用于容量推导,真实设计必须用压缩前后尺寸分布复算。
连续流要按峰值与积压恢复速率配置长期容量。若故障 10 分钟期间持续每秒 2 万条,会积压 1200 万条;恢复后只提供等于当前流量的能力,积压永远清不掉。微批则要证明单个切片在下一切片到来前完成。批处理可在低价时段集中资源,但小批过密会让启动、提交和小文件开销占比上升。
成本比较至少包含常驻计算、状态与检查点存储、接收器写入、数据扫描、值班和发布复杂度。不能只比较云账单;四人团队维护三套相近逻辑的时间也是成本。推荐方案把连续流限制在库存告警,避免为财务和看板承担不必要的全天状态成本。
第六步:用影子回放完成迁移,而不是一次切换
先固定一段包含正常流量、10 倍突发、重复、乱序和 6 小时迟到的历史输入。旧批结果作为基线,新微批和流作业在影子模式运行,不触发真实库存动作。比较每个业务键的计数、金额、版本和迟到修订轨迹,并解释每个差异。
随后逐级放量:先只写影子表,再开放内部看板,最后才让流告警触发可回滚动作。故障注入覆盖检查点前后崩溃、接收器超时、分区停滞、schema 变更和积压追赶。验收指标包括端到端 p99、最大积压年龄、微批完成时间、迟到修订率、重复动作数、批流差异金额和恢复时间。
退出条件也要明确。若微批连续超出 5 分钟,检查调度、倾斜、接收器和突发容量后再考虑连续流;若 10 秒告警不再驱动行动,可以降为微批以减少值班成本。处理模式是可验证的业务选择,不是永久身份。
高质量示范回答
“我会把无界输入和处理模式分开看。三个消费者的行动时限不同,所以不会为了技术统一把它们全做成流。库存异常确实要在 10 秒内触发动作,我会用连续流,逐段分配源、排队、计算、写入和通知预算,并按 eventid 与规则版本生成幂等动作。运营看板只要求 5 分钟,我会先用一到两分钟微批,结果按窗口和版本覆盖,页面显示数据截至时间。财务关账则读取冻结的业务日快照做批处理,用 runid 暂存、核对后原子发布,迟到数据走调整运行。
三条路径共享原始可重放事件、规范化数据契约和黄金样本,但财务保留独立校验。题设每秒 2 万条、短时 10 倍,我会验证每秒 20 万条入口以及故障期间积压的恢复倍率,而不只压测稳态。流作业的检查点不能替代外部动作幂等,正好一次也不能证明 6 小时迟到尾部已经完整。
迁移时先回放同一段历史,在影子表比较批、微批和流的键级结果及修订轨迹。连续流先只发观察告警,达到 10 秒 p99、没有重复动作、积压可在目标时间清空后再启用自动动作。最终规则是:选择满足行动时限、完整性和恢复要求的最简单模式;更低延迟只有在改变业务决策时才值得承担额外状态与值班成本。”
常见错误
- 看到 Kafka 就全部改成流处理 → 输入连续不代表每个结果都需逐条更新,财务与看板承担了无收益的状态和运维成本 → 按每个消费者的行动时限分别选择。
- 只说“流处理延迟低” → 接收器、缓存和通知可能吃掉全部预算 → 拆出端到端路径并监控各段尾延迟。
- 把 5 分钟新鲜度叫作最终正确 → 最长 6 小时的迟到事件仍会改变结果 → 明确估算、修订、冻结和调整协议。
- 认为检查点等于外部动作正好一次 → 重试可能重复调用库存或通知服务 → 使用动作幂等键、审计记录和可回放测试。
- 微批不断追加总数 → 同一窗口每批被重复累计 → 按窗口与版本原子覆盖,或定义清楚可撤回增量协议。
- 批与流各自复制业务规则 → 时区、取消单和金额口径逐渐漂移 → 共享契约与黄金样本,并做键级等价比较。
- 只按稳态吞吐扩容 → 10 倍突发和故障积压无法在时限内清空 → 同时验证峰值、积压恢复倍率和接收器容量。
- 上线即触发真实动作 → 语义差异直接影响库存 → 先影子写入和观察告警,再逐级开放可回滚动作。
追问及应对
追问一:5 分钟看板为什么不用连续流?
若一到两分钟微批在峰值和恢复场景下稳定完成,并且最终写入与缓存仍在 5 分钟内,连续流不会改变运营决策。它只增加长期状态、检查点、积压恢复和发布成本。若后来时限缩短到 30 秒,或微批调度与启动开销已经占掉大部分预算,再用同一历史回放比较连续流;升级条件应是 SLO 证据,不是“更实时”这个标签。
追问二:能否用一条流同时产出告警、看板和财务结果?
技术上可以,但要避免一个检查点、schema 变更或错误窗口同时阻断三种用途。告警可读取规范化流并维护短状态,看板从版本化聚合读取,财务仍从不可变历史按冻结快照重算。共享输入和定义,隔离发布与失败域。只有团队证明统一作业的故障影响、回填和审计都可接受时,才值得合并运行单元。
追问三:正好一次处理是否能让财务直接采用流结果?
不能单凭这个标签决定。处理语义可以防止已提交输出因重试重复,却不保证迟到记录已经到齐,也不自动覆盖外部副作用。财务还需要冻结输入范围、规则版本、可重复运行、总账约束与差异审批。流结果可作为早期估算;只有这些审计条件同样成立,并经过长期对账证明后,才可能升级为关账输入。
追问四:10 倍突发时如何证明系统能恢复?
记录突发时长与故障时长,计算积压条数,再测实际净清空速率。若当前输入为每秒 2 万条、消费者处理每秒 3 万条,净清空只有每秒 1 万条;积压 1200 万条需要约 20 分钟。验收同时看端到端告警时延、接收器限额、状态增长和自动扩缩滞后。只展示峰值吞吐而不计算净恢复能力,会掩盖长尾违约。
追问五:何时应该把连续流降回微批?
当业务不再根据秒级结果行动、微批已经能满足新时限,或流作业的值班与状态成本持续高于其避免的损失时,可以降级。先影子运行微批并比较结果与时延,再停用真实动作、保留可重放输入和回滚窗口。降级后仍要监控迟到修订和峰值完成时间,避免用成本优化换来不可见的数据陈旧。