WRITING / 2026.10.02

大数据专题(六):让指标可信,批流对账、历史补数与故障恢复

把 Kafka、Spark、Flink 和 Iceberg 实验串成可验证的数据链路,明确实时与离线口径,完成缺失检测、历史补数、幂等重跑和结果验收。

所有任务都是绿色,指标仍可能是错的

接入程序成功退出,批任务没有异常,实时作业也在运行,为什么业务仍会看到少算或多算的播放量?因为运行状态描述的是程序执行过程,指标正确性描述的是输入、规则和结果之间的关系,两者并不等价。

一份日报可以成功处理一个空目录;一个关联可以成功放大重复维度;一个实时窗口可以按预定策略排除迟到事件;一次重跑也可以成功把同一天的数据追加两遍。程序都完成了自己收到的指令,错误却发生在指令对应的业务契约上。

本篇把前五章的真实实验连起来。我们先区分实时与离线的比较范围,再删除教学表中的一天汇总,执行两次相同补数并核对结果。删除仅发生在专用实验 warehouse 内,不涉及博客写作库或线上业务数据。

从指标差异发现、输入定位到补数和结果验收

“一致”必须先有共同参照

两个数字相等,只能说明这两个数字相等。它们可能同时漏掉同一类事件,也可能一个多算、另一个少算而偶然抵消。真正的对账需要说明比较对象:哪些时间范围、哪些业务键、哪些合法记录、哪个去重规则、哪个维度版本和哪个计算口径。

本系列的独立参考计算器以事件 ID 去重,将合法观看毫秒数按日期和窗口累加。Spark 与 Flink 的输出分别与它比较,而不是只让两个引擎互相比。这样既能发现执行实现差异,也能避免“两个系统引用同一错误结果”的简单循环证明。

参照也有成本和适用范围。完整集合比较适合一万条确定性样例;更大数据通常需要按日期、分区、业务键分层核对,用数量、总量、摘要和抽样逐步缩小差异。摘要相同也不是无条件的数学证明,但可以作为工程排查的一部分。

最终仍应保留能够解释具体差异的输入身份。只有“昨天少了百分之一”而无法定位少在哪些视频、哪个小时、哪个接入分区,修复工作就容易变成盲目重跑。对账的设计目标不仅是报警,更是让问题可以被调查。

实时暂定值与离线校正值

第五篇的小样例清楚展示了两个合法答案:实时窗口接纳 a、b 和允许迟到的 c,结果为三;完整离线输入还包含过迟的 d,结果为四。重复的 a 在两边都只计一次。差异有具体事件作为解释,不是一个模糊的“最终一致”。

对实时实现进行正确性验收时,应将它与相同接纳范围的参考结果比较,确认三确实由预期事件组成。对最终业务指标进行验收时,则需要将完整输入纳入离线校正,明确四什么时候成为可见结果。两项检查承担不同职责。

如果业务要求实时结果也最终修订到四,可以设计延长保留、专门修正流或离线覆盖等方式,但每种方式都增加状态、输出协议和使用端复杂性。不能只在对账脚本里忽略差异,同时对外宣称实时指标已经包含全部迟到数据。

本文选择保留实时主结果及迟到诊断,完整离线结果作为最终核对依据。实验输出明确区分两种状态;真实产品应在接口或看板中表达更新时间与结果阶段,避免用户拿尚未完成的数据和昨日最终值直接比较。

把指标写成一个可检查的契约

一个指标至少应定义事件粒度、唯一身份、时间范围、时区、过滤条件、聚合方式和空值处理。观看时长的单位也要固定,毫秒和秒混用会产生数量级错误,而这类错误常常不会触发计算引擎异常。

维度关联应说明缺失值和重复键如何处理。使用最新维度与使用事件发生时维度,可能分别满足不同报表需求;不能让两个作业各自选择,然后要求结果自然一致。对于更新后的历史分类,还需要明确是否允许重述过去的指标。

结果接口应说明数值是增量还是绝对值。实时窗口输出二后再输出三,表示最新值变为三;若服务层将两条都当作增量,就会得到五。补数同理:重算日报后应替换受影响键的结果,而不是把整份新日报再次累加到旧日报。

契约还应附带版本。过滤有效播放的规则从一秒改成三秒,是业务口径变化,不应悄悄覆盖旧定义。可以生成新版本的指标,或按明确范围重算历史并通知使用者。技术上的表结构兼容不能替代业务上的指标可比性。

四类质量信号如何组合

完整性关注应到的数据是否到达。可以从源端批次或结束位点、接入清单、合法事件数和目标分组数建立检查。没有明确“应到范围”的情况下,单看目标有数据就难以证明完整。

唯一性关注重复业务身份是否被多次应用。除了统计重复数量,还应区分完全相同的重传和同 ID 不同内容的冲突。前者通常可以按规则去重,后者需要隔离或版本政策。把两者都计入一个“重复率”可能掩盖更严重的上游问题。

及时性关注数据落后业务发生多久,以及已发布结果多久没有更新。消费积压、事件时间落后和结果更新时间都是有用信号,但它们不表达同一件事。一个任务可能没有积压,却一直没有获得新数据;另一个任务可能仍有积压,但已满足某个历史补数范围。

有效性关注字段和业务约束。本实验将负观看时长作为非法记录,并在诊断中保留它。生产规则还可能包括时间范围、枚举值、关联键和单位检查。规则应可追踪、可解释,不能让静默过滤成为“数据质量已经改善”的唯一证据。

这些信号需要一起看。输入突然下降而错误率也下降,可能不是质量变好,而是上报整体失效;重复率上升但最终指标稳定,可能说明去重在发挥作用,也可能意味着上游重试风暴正在消耗资源。质量监控应连接业务结果与链路过程。

发现差异后,先定位再重跑

第一步冻结比较范围。记录时间、分区、offset 或快照,以及当前口径版本,避免排查期间输入继续变化,让差异不断移动。对于实时系统,可以抓取一个已提交结果快照,再用对应边界做比较。

第二步从粗到细拆分。先按日期和业务域,随后按小时、视频或接入分区,找出差异集中在哪里。计数一致但时长不一致,应优先检查数值字段、单位和聚合规则;某些视频全部缺失,则可能与维度关联或过滤条件有关。

第三步检查差异类型。缺失、重复、错误归组、错误值和发布延迟需要不同修复方式。直接全量重跑可能暂时让数字变化,却无法说明根因已经解决;如果上游规则仍然错误,下一次任务还会重复同样问题。

第四步确认修复输入。若原始事件仍在日志或原始层中,可以重新计算;若源头从未产生数据,重新执行下游无法补出事实。恢复计划应写清楚可以恢复的范围,以及必须由业务或上游重新提供的信息。

补数任务需要稳定的输入和输出身份

一次补数应有独立运行 ID,绑定输入范围、代码或规则版本和目标范围。相同业务日期的正常任务与补数任务可能同时运行,若二者都无条件覆盖结果,就可能产生旧任务覆盖新任务的问题。

本机实验按日期与视频组成稳定汇总键,重算后通过 Iceberg MERGE INTO 写入绝对值。匹配到现有键时更新,没有该键时插入。再次使用相同输入执行,相同键得到相同值,从而具备本实验所需的幂等语义。具体写入能力可参考 Iceberg Spark 写入文档。

这项语义有一个重要边界:仅更新或插入不会自动删除“旧结果里存在、新结果里已经没有”的键。本实验恢复的是被删除的同一份日数据,没有模拟口径收缩导致分组消失。真实重算需要按目标分区整体替换,或明确删除已不属于新结果的旧键。

并发发布还需要更强的协调。可以为结果附加版本、采用条件更新或发布指针,让较旧补数不能覆盖较新口径结果。本文没有实现多写者并发控制测试,因此不把单任务重复运行的幂等性扩大为所有并发写入的正确性。

计算成功与结果发布应有不同阶段

如果边计算边覆盖用户正在读取的汇总,用户可能看到一半新值、一半旧值。将计算结果先写到隔离位置,完成数量、摘要和业务规则检查后,再通过目标系统支持的机制发布,可以让失败更容易定位和回滚。

在表格式内部,提交和快照为结果版本提供支持;跨表或跨系统时,仍需定义发布的可见性边界。日报表、榜单缓存和导出文件不一定能够自动原子切换。应明确哪个对象代表一次发布完成,以及消费者如何识别对应版本。

回滚也应针对明确对象。切回旧表版本不等于撤销已发送通知,恢复旧代码不等于修改已生成指标,重新消费消息也不等于删除旧错误结果。把代码、数据和对外发布分别记录,才能在事故中选择正确恢复动作。

教学实验通过专用目录和表保留输入、阶段结果与最终结果,便于逐步比较。它没有向真实用户发布分析看板,也没有操作业务数据库;这种隔离使我们可以安全删除一份汇总,观察修复流程,而不会把教学故障施加给线上系统。

最小实验:删除一天汇总并连续补数两次

前面的存储步骤已经生成一万条合法明细和一百个视频日报分组。进阶脚本先在实验表中删除样例日期的汇总,将删除后的结果导出;随后基于原始明细重新计算,执行第一次合并,再用相同输入执行第二次合并。

python3 advanced.py /absolute/path/to/lab-work

脚本读取两次补数结果,逐个日期和视频比较播放次数、观看毫秒数,并与独立生成器的参考答案比较。这里检查的是最终业务值,而不是只判断两条 SQL 都返回成功。

阶段 本次实际结果
删除实验日期的日报 汇总结果为 0 行
第一次补数 恢复 100 个视频分组
第二次重复补数 仍为相同 100 个分组
与独立参考答案比较 两次均完全一致
与删除前合法明细对应 共 10,000 条唯一播放事件

这验证了确定输入、稳定键、绝对值更新下的重复补数行为。没有测试并发写者、源数据在补数中改变或业务口径变更,因此这些情形不在本次结论范围内。实验代码和命令位于源码包,实际记录位于结果摘要。

将实时差异纳入验收

针对小规模迟到样例,参考计算器分别计算被实时接纳的事件集合和完整合法集合。前者为三次播放、三百毫秒,后者为四次播放、四百毫秒;两者的差异恰好是一条进入迟到侧输出的事件。

这个过程给出了一种有解释的差异报告:主结果少了一条,但能够指出具体事件、原因和修正路径。它比设置一个百分比容忍区间更有信息。如果业务允许一定时延,差异可能在预期范围内;如果要求最终完整,则应等待离线校正并验证修正后的结果。

对于默认一万条基础数据,实时窗口与完整离线窗口一致,因为样例的乱序范围符合配置,并通过控制记录推进时间。这个成功不意味着所有未来输入也都一致;一旦上游晚到分布超过约定,侧输出与差异监控就应让变化变得可见。

对账报告也应包括输入规模和比较阶段。只写“批流一致率百分之百”,却不说明是否过滤了迟到、是否只比较非空键、是否忽略缺失窗口,很容易制造过强结论。实验报告保留明确分组结果,避免把一个汇总百分比当成完整证据。

日常运行时怎样观察这条链路

接入侧关注生产失败、分区积压、保留窗口和最老事件年龄;计算侧关注输入输出速率、任务失败、状态大小、Checkpoint 时间与恢复;存储侧关注文件数量、数据体积、快照与维护任务;结果侧关注更新时间、缺失分组和对账差异。

指标之间应有可跟踪关系。看到日报没有更新,可以沿着本次运行 ID 找到输入范围、Spark 执行记录和输出版本;看到实时结果停滞,可以区分没有新事件、Watermark 停滞和事务尚未提交。孤立地堆积仪表盘,很难缩短定位路径。

告警阈值需要结合业务节奏。夜间低流量不应自动被当成上报故障,促销期间的突发积压也不一定意味着无法恢复。可以先建立历史分布与明确的服务目标,再逐步收紧告警,而不是给每个指标统一设置一个静态数字。

本实验记录了进程、作业、消费范围、快照和结果文件,但没有部署长期监控平台,也没有运行数周统计分位数。它验证的是信息能否被采集、关联和用于一次真实故障恢复;生产告警质量需要在持续运行中另外评估。

六篇文章共同建立的工程边界

Kafka 提供可保留与重放的事件日志,Parquet 改善一类分析数据的组织,Iceberg 管理表版本与演进,Spark 执行有界计算,Flink 维护持续流的时间与状态。任何一个组件都不能独立回答业务指标是否完整、唯一、及时和可解释。

这些问题最终需要通过共同契约连接:事件有稳定身份,输入有明确范围,计算保留证据,结果具备版本,修复能够重跑且可核对。这样,当某个数字异常时,我们才能从结果追溯到事实,并通过一次范围明确的修复恢复可信状态。

阅读与复现实验可以从专题入口开始,依次查看环境和运行说明、各章源码与结果摘要。先在确定性小数据上验证机制,再逐步增加规模、输入变化与故障范围,比直接复制一个庞大平台配置更容易获得可以解释的认识。

大数据工程:从业务事件到可信指标

  1. 大数据专题(一):从业务数据库到数据平台,一条分析链路如何形成
  2. 大数据专题(二):Kafka 数据接入,分区、消费进度与可靠重放
  3. 大数据专题(三):从文件到分析表,Parquet、数仓分层与 Iceberg
  4. 大数据专题(四):Spark 批处理,一条 SQL 如何变成分布式计算
  5. 大数据专题(五):Flink 实时计算,乱序事件如何得到正确结果
  6. 大数据专题(六):让指标可信,批流对账、历史补数与故障恢复 · 当前文章

按时间浏览