题干与适用场景
这是数据工程师常见的质量治理设计题。重点不是列出校验规则,而是说明失败记录从发现、隔离、修复到重放的完整生命周期,同时保证正常路径的吞吐和正确性。
面试官考察什么
- 能否区分 schema、业务规则、重复、迟到和毒丸记录。
- 能否设计保留原文、失败原因和来源位置的隔离区。
- 能否用幂等键、版本化规则和可审计状态实现安全重放。
- 能否说明背压、告警、隐私、保留期限与责任边界。
回答前需要澄清的问题
先确认输入是批处理、消息流还是两者并存;是否允许有效记录先提交;事件是否有稳定的 event_id、业务时间和来源 offset;修复后要回放原始事件还是重新从上游拉取;以及合规要求的保留期限和敏感字段处理方式。
30 秒回答框架
我会把流程分成原始不可变层、校验路由层、有效数据路径和 quarantine 路径。每条失败记录都保存原文、来源定位、规则版本、失败原因与状态;有效记录按 event_id 幂等写入。修复通过版本化转换或更正源数据后,在隔离区上生成重放批次,沿同一幂等写入路径执行,并用指标和抽样核对证明没有丢失或重复。
分步骤深入解答
1. 先保留事实,再做路由
入口先写入不可变的原始层或可重放日志,记录 source、partition/offset、接收时间、event_id 和 payload hash。解析与 schema 校验失败时,不要直接丢弃;把记录和错误码送入 quarantine。校验通过的记录进入正常处理。这样坏记录不会阻塞整批,原始事实也可复核。
2. 设计可操作的隔离记录
隔离表至少包含 payload(或受控引用)、失败字段、规则名称与版本、首次失败时间、来源定位、重试次数、修复批次和状态。状态可设为 open、readyforreplay、replayed、rejected。敏感字段按最小权限加密或只保存指针;保留期限由合规和排障需求共同决定。
3. 让重放与正常写入使用同一正确性边界
修复规则必须版本化,不能覆盖原始记录。重放任务读取指定状态和规则版本,先在影子目标或小批量中验证,再调用与实时路径相同的转换和写入逻辑。目标表以 event_id 加业务版本作为幂等键;重复重放应得到同一结果,更新语义要明确是 upsert、忽略还是产生新版本。
4. 处理重复、迟到和毒丸记录
重复消息依靠 event_id、来源 offset 或去重窗口识别,不能把 offset 单独当业务身份。迟到事件按业务时间进入补算或更正流程,并说明水位线影响。反复失败的毒丸记录要限制重试、转人工或永久拒绝,避免占满消费者;单条失败不应让无关分区停止消费。
5. 用可观测性证明系统在工作
分别监控有效通过率、按规则和来源统计的 quarantine 数、未处理年龄、重放成功率、重复写入冲突、端到端延迟和数据新鲜度。质量阈值触发告警或暂停某个来源,而不是盲目暂停全部流水线。定期把原始计数、有效计数、隔离计数和重放计数对账,发现差额时优先检查路由和提交边界。
高质量示范回答
我会先问清输入语义和可接受的一致性,再把原始事件落到不可变存储,随后做 schema、业务规则、重复和时序检查。通过的记录走正常路径,以 event_id 加业务版本幂等写入;失败记录进入隔离区,保存原文引用、来源 offset、失败规则版本、错误详情和状态,绝不静默丢弃。修复时不改原始事件,而是生成一个有审批和批次号的重放任务,在小批量验证后复用正常转换与写入路径。重复重放必须得到同样结果,迟到事件要按业务时间补算,毒丸记录要限流并转人工。最后用通过率、隔离积压年龄、重放成功率、写入冲突和计数对账做告警与审计;敏感字段按最小权限保护,保留期限遵循合规要求。
常见错误
- 校验失败就丢弃或只打印日志,无法找回原文。
- 只保存错误字符串,没有来源 offset、规则版本和 event_id。
- 重放另写一套逻辑,导致实时和补数结果不一致。
- 用“重试直到成功”处理毒丸记录,造成无界积压。
- 把隔离区当成垃圾桶,没有状态、责任人、保留期限和关闭条件。
- 只看整体成功率,不按来源、规则和年龄拆分,无法定位退化。
追问及应对
如果隔离数据包含个人信息,怎么处理?
保留排障所需的最小字段,敏感 payload 加密并限制访问;可用不可逆 hash 关联去重,原文通过受控引用取回。审计访问和到期删除必须可验证。
重放期间新数据仍在流入怎么办?
使用独立重放批次和明确的事件版本,写入端按幂等键合并。若业务要求顺序,按分区或实体设置边界;否则允许并行但记录冲突与最终决策。
什么时候暂停整个流水线?
只有 schema 破坏兼容性、目标不可写或数据损坏风险会污染有效路径时才暂停。单一规则或来源异常优先隔离该分支并保留其他来源。
如何证明没有漏数?
以输入批次或 offset 范围建立账本,对账原始、有效、隔离、重放和拒绝计数;抽样核对 event_id 集合,并把未决记录年龄纳入 SLO。