数据落盘之后,分析问题才刚开始
上一章从 Kafka 导出了播放事件。文件完整保存之后,一个自然问题是:能否把所有 JSON 文件放进目录,再让 SQL 引擎直接查询?对于临时分析,这当然可行。但如果每天持续写入、字段逐渐变化、多个任务并发更新,目录里哪些文件属于当前结果,便需要更明确的规则。
例如,日报重算产生了一批新文件,旧文件还在原目录里。查询程序如果把两批都读进去,会重复计算;先删除旧文件再写新文件,中途失败又可能只剩半份结果。字段改名之后,旧文件中的列如何解释?查询昨天的版本时,到底应该读哪些文件?这些都不是“文件已经保存成功”能够回答的问题。
本篇沿着三层组织展开:文件格式如何表达数据,数据模型如何表达业务,表元数据如何组织一组文件。我们使用真实 Spark 和 Iceberg 运行实验,先验证正确性,再讨论空间和查询计划,不把小样例的文件大小直接当作生产成本承诺。
文件格式、表格式与计算引擎
CSV、JSON 和 Parquet 主要回答一份文件内部如何编码记录。CSV 直观、通用,但类型与转义需要额外约定;JSON 能表达嵌套结构,同时携带较多字段名称;Parquet 面向列式数据布局,提供类型、编码、压缩和读取所需的元信息。Parquet 官方概览介绍了它的定位。
Iceberg 则管理一个分析表由哪些数据文件及元数据组成,并通过快照描述表版本。它不替代负责执行查询的引擎,也不等同于保存文件的对象存储。Spark 可以读取 Iceberg 表,也可以直接读取 Parquet;这两种读取方式的输入组织与一致性边界不同。
存储介质又是另一层。本机磁盘、分布式文件系统和对象存储各有接口与故障模型。本实验使用独立本地目录作为 warehouse,目的是观察表和文件的关系;没有部署对象存储服务,也没有测试跨机器文件可见性或网络分区。
把三层分开以后,许多比较问题会变得清楚:Parquet 与 Iceberg 不是同一层的替代关系,Spark 与 Iceberg 也不是两个互斥的数据库。评估方案时,应分别确认文件编码、表管理、执行引擎和存储介质,而不是只选择一个听起来最完整的产品名称。
列式布局为什么适合一类分析负载
日报经常只需要视频 ID、事件时间和观看时长,而事件可能还包含设备、网络、页面、地域等字段。按行组织时,相邻字节通常对应一条记录的多个字段;按列组织时,同一列的值可以被放在相邻区域。对于只读取少量列、扫描大量记录的查询,后者有机会减少读取与解码的工作。
同一列中的值通常具有相近类型和分布,也有利于选择编码与压缩方式。例如重复的视频分类、递增的时间值与随机字符串,其压缩表现可能不同。最终空间占用仍取决于数据分布、编码、压缩算法、行组大小以及文件数量,不能只根据“列式”两个字推导固定压缩比。
列式文件并不是所有访问方式的最优选择。如果每次根据一个主键读取完整单行,或频繁执行非常小的更新,成本结构与大范围聚合不同。更新分析表时还可能需要重写文件或处理删除信息。应该用实际读写模式评估,而不是认为分析文件格式可以无条件替代业务数据库。
本实验让 CSV 与 Parquet 保存同一组六个字段,并将输出分区收敛为一个,减少文件数量不同对比较的干扰。CSV 使用 Spark 默认输出选项,Parquet 使用默认压缩;这比较的是两套明确配置下的落盘结果,没有将所有差异归因于单一因素。
一条播放记录在数仓里是什么
播放明细是事实数据,表达一次已经发生的业务记录;视频标题、作者、分类等描述信息属于维度。建模前首先要确定事实粒度:一行代表一次播放、一次心跳,还是一个用户对一个视频一天的汇总?粒度一旦混杂,同一条聚合 SQL 就可能把不同层次的数量相加。
本系列一行合法明细对应一个唯一播放事件,观看毫秒数可以在这个粒度上累加。视频维度一行对应一个视频 ID。示例维度是静态合成数据,因此可以用简单关联;真实维度如果变化,还需要决定用事件发生时的分类还是查询时的最新分类。
这个选择会直接影响历史指标。视频昨天属于“技术”,今天改为“生活”,昨天日报应该随之改变吗?如果答案是否,就需要保留维度历史及有效时间,或在事实中保存当时的必要属性。如果只存最新分类,再要求历史报表稳定,数据模型本身就缺少信息。
关联时也要检查维度键是否唯一。一张所谓维度表出现两个相同视频 ID,明细与它关联后可能被放大两倍。作业可以顺利结束,聚合值却全部偏大。因此,维度唯一性和关联前后行数检查应成为计算契约,不能等到用户发现指标翻倍才追查。
分层是为了管理规则的复用
数仓中常见原始层、明细层、汇总层和应用层的划分。命名可以不同,重点是每层承担什么职责。原始层尽量保留输入事实及来源信息;明细层统一字段、清洗和去重;汇总层按稳定维度减少数据粒度;应用层面向具体查询或报表。
如果每份日报都直接从原始事件重新实现去重和坏数据过滤,规则很容易漂移。同一事件在某张报表计一次,在另一张计两次,最后只能在汇报时人工解释。把共享规则集中到明细层,可以让多个指标在相同输入语义上计算。
分层也会增加复制、存储和调度成本。一个简单项目不必机械地创建几十张名称规范但没有实际职责的表。可以从原始文件、去重明细和日报汇总三层开始,只有出现明确复用或服务需求时再增加一层。本实验采用这个最小结构,让每一层都能找到对应输入与验收。
原始层并不意味着永远保存所有信息。敏感字段、访问权限和保留期限应在接入时就有规则;本实验使用合成标识,不涉及真实个人数据。真实平台应根据用途保留必要字段,同时保证排障和重算所需的数据不会在未知情况下消失。
分区裁剪解决什么,又不能解决什么
分区把数据按某个规则组织成相对独立的集合。查询昨天的数据时,如果引擎能够依据条件排除其他日期,便可以减少候选文件。这个收益来自查询条件与数据布局之间的匹配,不是给表加上任意分区字段就自动产生。
按用户 ID 这类高基数字段逐值分区,可能制造大量很小的目录和文件;只按日期分区,在某一天数据极大时又可能仍需扫描很大范围。设计时应结合常用过滤条件、每日数据规模与写入批次,避免只追求逻辑上看起来整齐。
还要区分目录级或表级文件裁剪、文件内部统计跳过和列投影。它们发生在不同层次,共同影响最终读取量。查询执行计划可以提供部分证据,但仅看到一个过滤条件出现在计划中,不等于已经证明扫描字节显著下降。应结合实际文件数、输入字节和运行指标验证。
默认一万条样例只覆盖很短时间,本篇用它验证分区表达和表元数据,不声称测得跨日期裁剪收益。若要评估收益,应生成多个日期、固定查询范围并观察实际扫描。把没有覆盖的数据分布写进结论,是分析实验中常见的证据越界。
小文件为什么会拖慢分析
每个文件都需要被发现、规划和打开,还会携带自己的元数据。一个文件只有几条记录时,管理它的成本可能比真正读取数据还明显。大量小批次写入、按高基数字段分区或过多写入并行度,都可能形成小文件。
但解决方法也不是把所有数据强制写成一个巨大文件。那会压低读取并行度,增加单个写入任务的压力,也使局部更新和重写更昂贵。合适的目标文件大小与压缩后体积、执行资源和查询方式有关,应观察分布而非只看平均值。
整理文件需要消耗计算资源,还可能与正常写入并发。表格式提供管理这些变化的机制,但维护任务仍需调度、限流与验收。删除旧快照、清理孤立文件和合并小文件是不同操作,不能因为目录看起来很乱就手工删除未知文件。
本实验没有持续数天制造小文件风暴,也没有测量对象存储请求费用。文章中的小文件讨论是机制与设计分析;实际测量部分只覆盖给定样例的文件体积、数据正确性和表演进。这样的边界使读者能够复现结论,而不是复制一组无法比较的性能数字。
快照把“当前表”变成可定位版本
将目录当表时,查询可能面对正在变化的文件集合。表元数据则可以记录一个明确版本由哪些文件构成,读者据此读取该版本。Iceberg 的快照、清单和数据文件分工,可参考可靠性说明与表规范。
快照为排查提供了具体锚点:某次日报基于哪个表版本生成,字段演进前有哪些数据,重算前后的输入是否相同。只保存“昨天跑过一次”这样的描述,很难在数据持续追加后重新解释结果;保存快照身份可以把问题缩小到一个明确输入集合。
历史查询仍受保留策略约束。旧快照及其需要的数据文件被正常清理之后,不能继续假设它们永远可读。设置清理策略时,需要将最长查询、补数窗口、审计和恢复要求考虑进去。快照提供版本机制,组织仍需决定要保留哪些版本以及保留多久。
同样,恢复表指针与撤销所有外部影响并不是一个操作。如果某个错误汇总已经被同步到另一个系统,回到旧快照不会自动修改那里的副本。完整回滚应列出消费者和下游结果的修正动作,第六篇会继续讨论结果发布与补数。
字段演进与分区演进
业务增加设备类型后,新事件可能带有 device,旧记录没有这个值。一个合理的查询应能同时读取新旧数据,并明确旧值如何表现。不能仅因为新文件有更多列,就依靠位置猜测所有历史列的含义。
Iceberg 用字段身份管理 Schema 演进,并允许表的分区规格演进。改变分区规格时,旧数据文件不会因此立即全部重写,新旧布局可以通过元数据一起被规划。Iceberg 演进文档解释了这些机制。
机制支持并不代表任意业务变化都安全。把“观看毫秒”改成“观看秒”而保持同一字段名称,属于语义变化;引擎即使能顺利读取,也无法替业务发现单位错了。应采用新字段、明确转换或版本化契约,并在消费者升级时核对指标。
字段新增实验也应检查已有数据,而不仅是查看表定义。本次增加可空字符串字段后,旧的一万行仍应存在,新增字段都为空;同时用增加字段前记录的快照 ID 查询,确认历史版本仍可读取。这两类检查分别验证当前兼容性与历史访问。
最小实验与实际结果
按照实验说明先完成 Kafka 接入,再执行:
python3 run.py spark --work /absolute/path/to/lab-work
脚本从真实 Kafka 导出的原始文件读取数据,过滤非法记录,按事件 ID 去重,写入 Iceberg 播放明细和日报表。日报、五分钟窗口结果都与独立生成器答案逐组比较。随后写出相同字段的 CSV 与 Parquet,并执行字段新增、历史快照和分区规格演进。
| 验证项 | 本次观察 |
|---|---|
| 去重后的合法明细 | 10,000 行 |
| 日报与窗口参考答案 | 均一致 |
| CSV 数据文件体积 | 550,999 字节 |
| Parquet 数据文件体积 | 114,085 字节 |
新增 device 后旧行空值数量 |
10,000 |
| 演进前快照读取行数 | 10,000 |
体积统计仅累加实际数据文件,不包括元数据、校验文件和运行日志;两份输出字段相同,但格式和默认压缩不同。这个结果说明本样例在当前配置下 Parquet 占用更少空间,不能推出所有业务都会得到同样比例,也不能据此推导查询延迟提升比例。
完整 SQL、依赖版本、快照与结果摘要均随实验包提供。快照 ID 是一次运行产生的身份,读者应使用自己运行生成的值,不能把文章中的旧 ID 复制到新建表上。
为什么实验保存查询与文件两类证据
SQL 返回一万行,证明查询路径读到了预期数量,但不能单独说明磁盘上没有多余文件;文件目录中存在产物,也不能说明查询已使用它们。实验因此同时保存输出结果、表元数据和文件体积,并分别核对。遇到版本切换或维护任务时,这种分层证据可以帮助区分“数据尚未提交”“当前快照没有引用”和“查询条件没有命中”,避免把不同问题都解释成文件丢失。
从能查询到能够维护
一个可维护的分析表应让使用者知道粒度、主键或去重身份、时间与单位、维度语义、版本和保留范围。执行层可以改变,存储介质可以迁移,但这些业务契约必须持续成立。否则每一次性能优化都可能悄悄改变报表含义。
接下来,我们保持输入和指标不变,转向 Spark 如何执行 SQL。只有在结果含义已经固定之后,比较不同 Join 和 Shuffle 计划才有意义;否则所谓更快的方案,可能只是少算了一部分数据。