从两个看似简单的需求开始
视频平台的运营提出两个需求:每天早上查看前一天各视频的播放次数与观看时长;活动期间查看每五分钟更新的热门视频榜单。业务数据库里已经存着视频信息,客户端也能上报播放事件。第一版似乎只需要建一张表,再写两条聚合 SQL。
当数据不多、查询不频繁时,这样做完全合理。真正的分界线并不是“超过多少万条就必须用大数据”,而是这两类工作是否开始争夺同一组资源。播放事件持续写入,报表需要扫描历史,运营不断增加筛选条件,临时查询又可能关联用户与视频信息。一条统计查询的扫描范围与在线接口的延迟预算,开始发生冲突。
本系列用合成视频播放事件搭建一个教学实验,讨论这条链路如何形成。所有用户和视频标识均由程序生成,不代表真实业务、生产规模或个人项目成绩。我们会真正启动组件、观察执行与恢复,但单机结果不外推到大型集群。
先定义“数什么”,再决定“用什么”
“播放量”不是一个天然唯一的数字。点击播放算一次、开始解码算一次、播放超过三秒算一次、一次会话结束算一次,都可能是合理业务定义,却不能混在一起比较。没有定义事件边界,后面的精确一次、批流一致和对账都找不到共同对象。
本实验将一条合法的 play 事件视为一次待统计的播放记录,用 event_id 标识事件身份,watch_ms 表示该记录对应的观看毫秒数。同一事件被重复传输时沿用同一个 ID;用户再次播放则产生新的 ID。这个定义用于讲解计算机制,不替代真实播放器的会话、心跳和有效播放规则。
还要明确时间。event_time 表示事件发生时间,ingest_time 表示样例中的到达时间,两者都保存 UTC 毫秒时间戳。日报按上海时区的自然日分组;实时统计按不重叠的五分钟事件时间窗口分组。这里的五分钟榜单是固定窗口榜单,与“任意当前时刻向前滚动五分钟”不同,后者需要滑动窗口或其他增量维护方式。
视频分类放在独立维度数据里,观看时长采用整数毫秒累加,避免浮点求和引入无关差异。榜单先按播放量降序,再按视频 ID 升序,给并列结果一个确定顺序。口径中的这些细节,比先画出十几个技术组件更重要:它们决定什么样的测试结果才算正确。
在线事务和分析查询为什么需要分工
在线接口通常知道要操作的少量记录。例如,按视频 ID 查询标题,更新某位用户的收藏状态,提交一次支付状态转换。它们关心操作的事务边界、并发冲突、索引访问以及稳定的响应时间。业务峰值到来时,单次请求的工作量最好保持可预测。
分析查询则经常从一大片记录中提取少量答案。计算一天的播放量,输入可能是当天全部事件,输出却只有几百个聚合分组。查询增加一个维度,可能改变关联方式和中间数据量。分析系统通常需要提高扫描、压缩、并行计算与重复查询的效率,而不是把每次查询都限制为几行数据。
这不意味着关系数据库不能做分析,也不意味着分析系统完全没有事务。两者的功能存在重叠,区别在于负载与优化重点。对于规模较小的站点,一个只读副本、合理索引和定时汇总表就可能足够。应先测量业务库的扫描行数、锁等待、资源利用和请求延迟,再判断是否需要拆分。
拆分之后也不是免费获得性能。你需要维护数据传输、重复处理、数据延迟、权限和运维监控。每增加一个系统,都增加了一个可能独立失败、升级或积压的边界。因此本系列的架构不是要求所有项目从第一天就照搬,而是用一个完整案例说明这些边界为何出现。
大数据首先是一组约束
数据量是约束之一,但不是唯一约束。同样一亿条事件,只做每晚一次扫描,与要求每条新事件几秒内反映到榜单,设计难点不同。字段类型稳定与每天变化、可以重算与原始数据无法保留、单一使用者与多团队共享,也会影响系统形态。
可以先写出四项预算:允许的数据延迟、需要保存的历史范围、查询的扫描范围与并发量、可以承担的计算和运维成本。然后再问,现有系统在哪一项预算上无法满足要求。这样得出的组件选择能够解释原因,也更容易在需求变化后重新评估。
例如,日报允许数小时延迟,夜间批任务能够充分利用顺序读取与批量聚合;活动榜单要求更快反馈,维护持续更新的状态更合适。两者可以共享原始事件和指标定义,而采用不同执行方式。共享数据并不要求共享同一段代码,更不自动保证结果相同。
将系统拆成五个职责
接入层负责接受事件、检查基础格式并将事件可靠交给下游。这里首先要回答确认成功意味着什么,以及下游短暂停止时数据放在哪里。Kafka 在实验中提供可保留、可独立读取的分区日志;它的职责是承接事件流,不负责替我们定义播放量。
存储层负责保存原始数据与可查询的数据表。原始记录为重放、纠错和审计提供依据;清洗后的明细为多个指标共享规则;汇总结果减少服务查询时的扫描。文件、表元数据和指标之间不能直接画等号,第三篇会分别解释 Parquet 与 Iceberg 的位置。
计算层把输入转为业务结果。Spark 批处理读取一个有边界的数据集合,生成日报与离线窗口结果;Flink 持续处理事件流,维护去重和窗口状态。它们的名字不是语义承诺:任务成功只能证明执行过程结束,仍需要检查输入是否完整、逻辑是否正确、输出是否可见。
服务层将结果交给使用者。教学实验通过 SQL 与 JSON 快照读取结果,不再加入一个独立的在线分析数据库。真实产品可以根据查询延迟、并发和维度组合评估 ClickHouse、Doris 或其他服务方式,但不能因为计算结果已经写成文件,就假设用户查询会自动变快。
治理与观测贯穿上述各层。你需要知道一个指标来自哪批输入、采用哪个口径版本、何时更新、是否缺数、能否重算。治理不只是末尾加一张“数据血缘图”,而是每个边界都保留足够身份与证据,让故障能够被解释和修复。
批处理与流处理如何共享一个问题
批处理面对一个有限集合。开始任务前明确输入范围,处理结束后能够计算总数和摘要,再决定是否发布结果。一个批次可以是昨天的事件,也可以是某段 offset 范围,不必绑定自然日。关键是输入边界可描述、可复现。
流处理面对持续到来的记录。它无法等到“所有事件结束”再计算,必须定义何时输出阶段性结果,以及之后到达的数据如何影响已有结果。状态、事件时间进度、迟到策略和故障恢复因此成为核心问题。实时并不等于立即得知完整事实,而是在不完整输入下持续给出有明确语义的答案。
同一条事件晚到十分钟,批任务在数据到齐后重算可能会包含它,实时窗口则可能已经关闭。两个结果不一致未必说明某个引擎算错,也可能说明输入边界不同。第六篇会先对齐“哪些事件允许进入本次比较”,再比较数字,避免用错误的对账方式制造告警。
本系列保留两条计算路径,目的是让这些差异可见。原始事件进入 Kafka 后,一路导出为有界文件,交给 Spark 和 Iceberg;另一路由 Flink 生成窗口更新,再以事务方式写回 Kafka。核对脚本用相同键提取最终窗口值,与独立参考结果比较。
如何理解熟悉的生态名称
Hadoop 常用来指一组围绕分布式存储、资源调度和计算的技术,也经常被宽泛地当作大数据生态的代称。学习时应拆开具体职责:文件系统解决存储组织,资源管理器安排计算资源,计算引擎执行任务。知道产品名不足以判断一条数据到底经过了哪些边界。
Hive 使数据仓库能够以表和 SQL 的方式组织分析;Spark 提供分布式计算能力;Flink 强调有状态的数据流处理;Kafka 保存与传递事件流;Iceberg 管理分析表的元数据与演进。它们有交叉能力,不能简单用“谁替代谁”的一句话概括。Spark 官方概览和 Iceberg 入门文档可作为理解计算与表格式分工的入口。
本机实验选择 Spark 本地执行、Iceberg 本地文件目录、单节点 Kafka 和本地 Flink 集群,主动缩小安装与网络问题的范围。这里出现了真实执行计划、Checkpoint 和事务输出,但没有测试跨机器容灾、磁盘损坏或多租户资源隔离。把本机多任务运行称为生产分布式容错,会扩大证据能够支持的范围。
最小实验:生成数据和独立答案
先下载完整源码包,阅读实验环境说明。所有章节共享一个独立工作目录,原始数据和运行产物放在该目录,博客服务器只分发代码与图解,不执行这些计算任务。
python3 generate.py ./lab-data
默认样例包含一万条唯一合法事件,额外复制其中一百条,再加入一条观看时长为负数的非法记录,因此输入文件共有 10,101 行。生成器采用固定种子,输出视频维度以及按同一口径计算的日报和五分钟窗口参考结果。后续实验重新读取真实 Kafka 导出的记录,分别经过 Spark 与 Flink,最终与这些答案比较。
参考计算器必须尽量独立。它使用 Python 字典按事件 ID 去重、按日期和窗口键累加,不直接调用 Spark SQL 或复用 Flink 的窗口实现。如果实现和测试共用同一个有缺陷的聚合函数,两个结果相同也不能说明计算正确。独立实现不能消除所有共同误解,但能降低同一代码错误同时污染两边的概率。
默认数据集中在短时间范围内,是为了快速运行正确性实验。它不能代表一天的峰谷分布,也不适合直接评估业务库和大数据平台的成本差异。跨自然日、边界时间、坏字段和大规模倾斜另由定向用例验证。读者应将“默认演示数据”和“覆盖边界的测试集合”分开理解。
怎样定义一次有效验收
第一层检查输入。应能说明写入多少条、合法多少条、重复多少条、错误多少条,事件 ID 是否存在内容冲突。只检查总行数不够:丢掉一条并重复另一条,总数仍然相同。对于确定性小样例,可以比较完整事件集合;对于更大的数据,再按分区和业务键分组检查数量与摘要。
第二层检查过程。Kafka 消费位置是否到达记录的结束位置,Spark 是否发生了预期的关联和 Shuffle,Flink 是否完成了 Checkpoint,故障后是否确实恢复了状态。日志中没有异常不能替代这些检查,因为任务可能正常处理了一个错误的输入范围。
第三层检查结果。日报与窗口聚合应与参考答案一致,榜单的并列顺序应稳定,重复执行补数不能再次累加。对于实时结果,需要区分首次输出、允许迟到后的更新和离线校正值;不同阶段的数字各自有意义,不能只留下最后一个截图。
最后检查资源与边界。本次启动的进程应能停止,数据目录应可定位,依赖版本应可复现,机器环境应被记录。如果只有作者电脑上某个隐藏服务存在时才能成功,实验就还没有形成可以交付的工程材料。
给结果附上一份最小契约
一个可以交付的指标不应只有数值。至少还要附上输入范围、统计时区、去重规则、口径版本、生成时间和结果状态。调用方据此判断数据是否足够新、能否用于比较,以及重算之后应该替换哪一份旧结果。否则,同一个接口今天返回实时暂定值、明天返回离线修订值,使用者可能把正常校正误认为业务波动。
这种契约可以先保存在一个简单的结果清单中,不必从第一天就搭建庞大的治理平台。重要的是每次运行都产生它,失败时也能解释停在哪个阶段。等到团队和链路规模扩大,再把这些已存在的信息接入统一目录与监控,通常比事后猜测数据来源更可靠。
从需求倒推演进顺序
如果只需要每天几个固定指标,可以先从业务库只读查询、定时任务和汇总表开始。出现历史扫描影响在线负载时,再考虑数据复制与独立分析存储;多个下游需要独立消费、回放和追赶时,事件日志的价值才更清楚;需要持续维护低延迟状态时,再进入有状态流处理。
每一步都应说明增加了什么能力,也增加了什么责任。引入事件日志后需要管理保留范围与积压;引入分析表后需要处理字段和分区演进;引入流计算后需要面对状态、迟到与输出一致性。架构图上多一个方框,只代表多了一个工具位置,不能证明对应责任已经得到履行。
阅读下一篇时,可以带着一个具体问题:如果消费程序已经把事件写到文件,却还没提交消费位置就崩溃,恢复后发生什么?这个小故障足以把接入、重放、去重和结果发布连接起来,也是整套实验从“能跑”走向“能解释”的第一步。