题干与适用场景
这是数据平台、实时分析和高级数据工程岗位常见的系统设计题。输入是峰值每秒 50,000 个事件,运营看板要求 60 秒内可见,财务报表要求前一天数据在 7:00 前完成。题目要判断共享执行图何时会把严格 SLO 传递给所有消费者,以及如何用独立消费进度和资源池降低耦合。
面试官考察点
- 能否先把“实时”和“按时完成”写成端到端 SLO,而不是直接报工具名称。
- 能否区分 Beam 图内分支、同一消息源的独立订阅和完全独立的计算资源。
- 能否解释重复读取、成本、背压、重放、迟到事件和故障域之间的取舍。
- 能否给出指标、灰度、回滚和数据对账,证明拆分真的改善了用户结果。
回答前需要澄清的问题
先确认 60 秒是事件到看板的端到端新鲜度,还是处理器延迟;财务报表是否允许迟到分区在次日补算;两个消费者能否共享原始存储;事件是否需要按租户或优先级公平处理;重复、丢失和顺序的容忍度;以及预算是更看重成本还是看板尾延迟。若报表只是 T+1,而看板是严格低延迟,两个 SLO 已经暗示资源隔离值得优先评估。
30 秒回答框架
我会先把两个消费者的端到端 SLO、数据正确性和恢复目标写清楚。若一条共享管道必须同时满足最严格的 60 秒目标,我会先建立单管道基线,再用独立订阅把实时和批处理的消费进度隔开;计算代价高的公共解析可以共用,资源紧张或故障传播明显时再拆成独立作业。每条路径都记录延迟、积压、迟到、重复和报表完成时间,并用回放和故障注入验证拆分是否值得。
分步骤深入解答
1. 从 SLO 推导结构
把实时路径定义为“事件时间到可查询时间的 p99 不超过 60 秒”,把批处理定义为“前一天有效事件在 7:00 前完成,迟到数据进入受控补算窗口”。Google Dataflow 文档强调,混合 SLO 的单管道必须承受更严格的目标;这会让低优先级工作共享实时资源。若两个 SLO 的预算、告警和扩容策略不同,独立路径更容易解释和运维。
2. 比较三种拓扑
单管道分支适合公共解码和轻量路由,单次处理事件后输出多个集合,避免重复执行昂贵的解析。Beam 文档说明,同一 PCollection 可被多个变换读取,但每个变换都会再次处理输入;单个多输出变换可让每个元素只经过一次公共计算。
同一主题的独立订阅适合让实时和批处理拥有各自的确认、积压和重放进度。Google Dataflow 文档给出的多管道方案正是利用独立订阅,让不同作业独立拉取并确认消息。完全独立的作业再进一步隔离 CPU、内存、发布节奏和故障域,代价是重复读取、重复序列化与更高运维成本。
3. 设计推荐路径
我会保留一个不可变原始事件层,并从同一源建立实时订阅和批处理订阅。实时作业只做轻量聚合并写入低延迟查询层;批处理作业按事件日期读取保留数据,写入分区表。公共解析若占总 CPU 的 20% 以上,可在入口做一次规范化并把版本化事件写入原始层;不把实时作业的处理进度当成批处理的提交证据。
raw-events
-> realtime-subscription -> stream-aggregate -> serving-store
-> batch-subscription or retained-raw -> daily-transform -> partitioned-lake4. 处理背压、优先级与成本
实时路径设置独立并发上限和积压告警;批处理在实时资源紧张时降低并发,但不能共享同一无界队列。若成本不允许两套完整计算,先共享解码与落盘,再隔离下游计算。若实时积压超过 60 秒,暂停扩大批处理并优先恢复实时 SLO。每条路径按输入字节、处理 CPU、积压年龄和单位输出成本计费,不能只比较作业数量。
5. 迟到、重放和故障恢复
实时窗口使用 watermark 和有限 allowed lateness;超窗事件进入迟到队列或原始层,由批处理补算并以版本号覆盖受影响分区。每条路径都用事件 ID 或业务主键实现幂等写入。重放时从保存的源位置建立新订阅,禁止把生产消费者的确认位倒退到旧位置。若实时作业失败,批处理仍应能从原始层恢复;若原始层不可用,两条路径都要明确降级和告警。
6. 用验收实验决定是否拆分
先运行共享基线,再对一小部分租户启用独立资源。比较看板 p50/p95/p99 新鲜度、批报表完成时间、消费积压、重复率、重放耗时、CPU、存储和每百万事件成本。注入批处理突发、实时处理器重启、消息重复、分区迟到和订阅暂停,验证实时 SLO 是否仍满足。若拆分只降低作业延迟却增加重复数据或成本超过预算,应保留共享方案并继续优化公共步骤。
高质量示范回答
我会先把两个目标写成端到端 SLO:看板从事件时间到可查询时间 p99 不超过 60 秒,财务报表在次日 7:00 前完成,迟到数据进入有限补算窗口。先做一条共享管道基线,但不会让批处理和实时路径共用确认位。我的默认设计是同一不可变原始事件层、两个独立订阅和两个下游作业;实时作业做轻量聚合,批作业按事件日期读取保留数据。公共解析可在入口版本化一次,只有在资源或故障域需要时才拆出完整作业。两条路径分别监控新鲜度、积压、迟到、重复、报表完成时间和单位成本,重放使用新订阅与幂等键。通过批处理突发、重启和迟到事件的故障注入验证;如果拆分没有改善用户 SLO,就回退到共享计算并优化公共阶段。
常见错误
- 错误表现: 因为有两个消费者就复制两套完整管道。失败原因: 公共解析和落盘成本被无谓翻倍。修正方法: 先共享不可变原始层,再按 SLO 隔离下游计算。
- 错误表现: 用一条全局消费位驱动实时与批处理。失败原因: 慢消费者会拖住快消费者,重放也无法独立。修正方法: 使用独立订阅或可证明独立的消费进度。
- 错误表现: 只看处理器延迟,不看事件到用户可见的延迟。失败原因: 存储、查询和下游刷新仍可能超时。修正方法: 记录端到端新鲜度和分位数。
- 错误表现: 迟到事件直接写回实时结果。失败原因: 可能重复计数或破坏已发布报表。修正方法: 设定 watermark、补算窗口、版本和幂等键。
- 错误表现: 用重启生产消费者的方式回放历史。失败原因: 会扰乱在线确认位并放大流量。修正方法: 从保留源建立独立回放订阅并限速。
追问及应对
如果实时路径和批处理路径都需要同一份昂贵特征计算怎么办?
把特征计算拆成版本化、可重放的中间层,先落盘再由两条路径读取;只有特征状态必须在线且不可共享时,才接受重复计算。用 CPU、延迟和一致性实验比较共享中间层与复制计算。
独立订阅会不会让输入成本翻倍?
会增加读取和确认开销,但可通过共享原始落盘、压缩、保留窗口和按需回放控制。把每百万事件成本与实时 SLO、故障隔离收益一起评估,不能只看存储账单。
批处理落后时,能否临时借用实时资源?
可以设置有界、可抢占的低优先级容量,但实时路径拥有硬上限。借用必须有租约、自动回收和独立积压指标,避免批处理把实时 p99 推过 60 秒。
如何证明两条路径最终结果一致?
用相同事件版本、业务主键和时间边界生成对账集合,比较计数、金额、缺失、重复和迟到修订。允许实时近似时,必须定义最终一致窗口和可解释差异,不能用一次总数相等代替长期对账。