数据工程面试:如何设计可审计的数据质量隔离通道?
题干与适用场景
上游每天发送 2 亿条订单事件,约 0.2% 可能缺少必填字段、类型解析失败或违反业务规则。有效事件要继续进入事实表和下游指标;坏事件不能静默丢弃,修复后还要能从原始输入重放,并保证同一事件不会重复计入。
本文采用批流统一的约束:每条输入有稳定 event_id,原始载荷保存不可变副本;质量规则按 schema 版本发布;隔离记录保留失败原因、规则版本和重试状态。题目核心是数据工程中的分流、可观测性和恢复,而非某个具体 Spark API。
面试官考察点
- 能否把解析失败、字段缺失和业务校验失败分成可行动的错误类型。
- 是否保留原始证据、规则版本和 lineage,支持审计与重放。
- 是否区分
fail-fast、drop、redirect/quarantine的适用场景。 - 是否用幂等键、去重状态和输出版本防止重复计数。
- 是否设计质量指标、告警阈值、人工修复和回放门禁。
- 是否说明毒性数据、PII、保留期限和隔离区访问控制。
回答前需要澄清的问题
- 0.2% 是允许的业务缺陷比例,还是任何超标都必须阻断发布?这决定门禁是阈值还是零容忍。
- 事件是否可补发、乱序或重复?若可重复,必须让
event_id成为幂等边界。 - 规则失败是否可以自动修复,还是必须人工确认?这改变重放队列和审批流程。
- 下游指标能否接受延迟或更正?若不能,隔离修复后需要补偿分区和版本化报表。
- 原始载荷是否包含个人信息?这会改变加密、脱敏、访问和删除策略。
30 秒回答框架
我把管道分成原始不可变层、解析层、规则验证层和有效/隔离两个输出。解析失败与规则失败都生成带 eventid、schema 版本、规则版本、原因码和原始引用的 quarantine 记录;有效事件通过幂等写入进入事实表。隔离区提供修复、审批和重放队列,重放仍使用同一 eventid 并在目标分区去重。监控有效率、各原因码占比、隔离年龄和重放成功率,门禁按业务 SLO 决定阻断或告警。
分步骤深入解答
1. 先保存原始证据
原始对象存储或日志层按批次、来源和接收时间分区,内容不可变并带校验和。解析任务只追加处理状态,不覆盖原文。这样规则升级、解析器修复或供应商争议都能从同一输入重现结果;原始层与隔离层要分离权限,避免支持人员直接修改事实数据。
2. 将失败分层并保留原因
先做字节/格式解析,再做 schema 类型和必填字段检查,最后做跨字段与业务规则检查。每次失败写入结构化 reasoncode,例如 MALFORMEDJSON、MISSINGORDERID、INVALID_CURRENCY。一条记录可以有多个原因,但要保存首次失败阶段和规则版本,避免修复后无法解释历史结果。
quarantine_record = {
event_id, source_batch, raw_uri, payload_hash,
schema_version, rule_version, failed_stage,
reason_codes, first_seen_at, status
}Spark 的文件读取选项可以记录坏文件或忽略损坏文件,但“继续运行”不等于业务记录安全;题目要求把可恢复的记录显式写入隔离通道,而不是只打开忽略开关。
3. 选择 fail、drop 还是 quarantine
基础设施不可读、签名不可信或可能污染全批的数据应 fail,让批次重试并保留告警。可定位到单条记录且不会影响其他事件的错误应 quarantine,让有效记录继续。只有经过业务负责人批准、且记录不可恢复并有审计要求时才 drop;drop 必须计数并可追溯,不能把静默丢失当成成功。
4. 有效写入与幂等边界
以 event_id 加来源版本构造唯一键,事实表写入使用幂等 upsert 或去重日志。隔离记录重放时先检查目标表是否已经接受该事件,再决定跳过、更新或写入补偿版本。对于订单金额这类可更正事实,不要直接覆盖旧值;应写入更正事件并让下游按版本或有效时间重算。
5. 修复、审批和重放
修复工具只能生成新载荷或修复补丁,不能修改原始层。每次修复记录操作者、原因、输入哈希和规则版本,并进入审批队列。重放 worker 从隔离状态读取,重新执行完整验证链;成功后把状态从 QUARANTINED 原子更新为 REPLAYED,失败则增加尝试次数和下次时间。并发重放用租约或数据库锁避免同一事件同时提交。
6. 质量指标与发布门禁
至少监控总接收量、有效率、每个 reason_code 的比率、隔离区年龄分位数、重放成功率、重复事件数和下游更正量。门禁按规则严重度分层:例如签名错误零容忍,缺少可选字段只告警;0.2% 也不能直接当作“正常”,要与历史基线、来源和业务损失一起判断。超过门槛时冻结下游发布或切换到上一版规则,并保留人工放行记录。
7. 保留、隐私与失败恢复
隔离区只保留修复所需的最小原始字段,敏感载荷加密并限制访问;保留期限与删除请求应能映射到 event_id 和原始对象。队列、元数据表和对象存储要分别备份。若下游写入成功但状态更新失败,依靠唯一键重试;若状态已标记重放但写入未确认,则通过提交日志或目标表校验恢复,不能凭 worker 返回值猜测成功。
高质量示范回答
我会先确认缺陷是否允许、事件是否可重复补发,以及订单更正能否延迟到报表。管道保存不可变原始层,解析、schema 和业务规则分阶段执行;单条可恢复错误写入 quarantine,基础设施或安全错误阻断批次。隔离记录包含事件幂等键、原始引用、规则版本、原因码和状态。有效事件幂等写入事实表,修复工具生成新载荷并经过审批,重放再次执行完整验证,成功后原子标记已重放。监控有效率、原因码、隔离年龄、重放率和重复计数,按严重度设置门禁。这样既不会让少量坏记录阻塞全批,也不会把异常隐藏成数据正常。
常见错误
- 打开
ignoreCorruptFiles就算完成 → 解析能继续但记录可能永久消失 → 把可恢复记录写入带原因码的隔离区。 - 失败记录直接放一张可编辑表 → 原始证据被篡改 → 原始层不可变,修复只产生新版本。
- 重放直接再插入事实表 → 重复计入指标 → 用
event_id唯一键和提交日志做幂等。 - 所有错误都阻断整批 → 少量坏数据拖垮新鲜度 → 按阶段、严重度和隔离能力选择 fail 或 quarantine。
- 只统计失败数量 → 无法定位来源和规则回归 → 按原因码、schema 版本、来源和时间分布监控。
- 修复后绕过验证 → 新载荷可能引入第二个错误 → 重放必须重新执行完整验证链。
追问及应对
隔离比例突然从 0.2% 升到 8%,要不要继续发布?
先按原因码和来源拆分。如果是单一供应商的可恢复字段缺失,暂停该来源并继续处理其他来源;如果是 schema 解析或签名错误,冻结下游发布并回滚规则版本。阈值应绑定业务损失和历史基线,不能只看绝对百分比。
如何保证重放不会改变已经结算的订单?
把原始事件和更正事件分开,事实表保存版本或有效时间。结算快照锁定输入代;修复事件进入补偿流程,由财务规则决定是否生成调整单,而不是静默覆盖历史金额。
隔离区也包含 PII,支持人员如何排查?
默认只展示脱敏字段和原因码,原始载荷用短期授权访问并记录审计日志。删除请求通过 event_id 关联原始对象、隔离记录和派生索引,删除后保留不可逆的审计摘要。
worker 在目标写入后崩溃,状态仍是待重放怎么办?
重试前按幂等键检查目标表和提交日志;已存在则补写状态,不再重复业务副作用。若两者都没有,重新提交。状态迁移要使用可重试的条件更新,避免把未知结果标记为失败或成功。
什么时候应该丢弃而不是长期隔离?
只有记录不可恢复、无合规保留要求且业务方明确接受损失时才丢弃。即使丢弃,也要保存数量、原因、批次和策略版本的审计摘要,并让监控和数据质量报告可见。