LabHub
学习 学习路径 课程

数据流水线

schema — 现在不校验,以后加倍还

在 LabHub 中继续学习

一句话总结

假设外部传入的数据都是字符串,只有在明确验证和转换后才能放入有类型的表格中。

概念图: 原始(staging) · 净化(clean) · 拒绝(reject) · 其理由

为什么需要这个?

实际操作的原始数据无一例外都很乱。同一天的数据以四个格式输入,金额附有货币符号和千分位点,状态值大小写和前后空格各不相同,必填项为空。

如果试图将这个直接放入有类型的表格中,就会失败。然后通常会出现两种错误的应对方式。要么跳过失败的行(数据不知不覺消失),要么把所有列都变成文本(将问题推迟到下游)。

怎么行动

标准的结构有三层。

  1. 原始(staging) — 所有列都是字符串。保持原样保存。绝对不修改。
  2. 净化(clean) — 只带合格的行类型进入。
  3. 拒绝(reject) — 包含未能通过的行为和其理由

这种结构的核心是保存定律。精制件数和拒绝件数的和必须与原始件数完全相同。如果有任何没有的行,那就是悄无声息消失的数据,是管道中最危险的事故。

留下拒绝理由也是不能妥协的。没有理由就丢弃的行为,以后没有人能恢复。像“amount无效”、“email无效”这样短也不错,一定要一起写。

转换规则也必须是明确的。如果日期格式有多个,则用正则表达式识别每个格式,并应用不同的解析规则。交给自动推理的话03/04/2025根据3月4日还是4月3日,会安静地出错。

将金额中的空字符串改为0也是常见的错误。**没有值和0是不同的。**如果结算金额为空,那不是0韩元的结算,而是信息遗漏,一旦填充为0,这一事实就会被抹去。

在现场相遇的样子

将模式约定留作代码是有帮助的。将精简表格的列和类型列表记录在单独的表格中,以后有人更改列类型时,下游管道可以立即检测到。模式更改是原先悄无声息地发生,几天后以奇怪的数字出现的类别的错误。

选择格式也值得注意。CSV可以在任何地方读取,但没有类型信息,分隔符逃避很脆弱。JSON表示嵌套,但体积很大。Parquet等列导向格式同时包含类型和统计信息,压缩率好,对分析工作负载有利。原始保存以CSV或JSON为准,分析用数据以列格式保存的分组配置很常见。

安全更改结构表的规则

管道图的模式是生产者和消费者之间的合同。只要改变一方就 因为对面会碎掉,所以先确定在哪个方向兼容。

兼容方向 意思 允许的变更
子兼容(backward) 新消费者读取旧数据 删除字段,添加有默认值的字段
上位兼容(forward) 旧消费者读取新数据 添加字段,删除选择字段
完全兼容(full) 双方都 仅添加·删除有基本值的选择字段

在流媒体中,以下游兼容性为基础。首先提升消费者,生产者 因为以后上传就可以了。如果倒序的话,新数据就会遇到以前的消费者。 会碎掉。

**添加必填字段总是打破的变更。**给出默认值并选择 放入后,所有生产者填完后,那时必须上传。分为两个步骤。 是正规。

文件格式决定性能

形式 结构 适合的地方 注意
CSV 人看得见的小量 无类型。编码·分隔符地狱
JSON Lines 结构化数据流式加载 大而慢
Avro 流媒体、活动 与模式注册表一起
Parquet 分析查询 写入量大。不适合小文件

热导向(Parquet)在分析中快速的原因是只读所需的列。列 50个中只用3个的提问很常见,行为倾向是读完50个。这里还有十个单位 压缩效果好,体积也变小了。

但是Parquet **如果有很多小文件反而会变慢。**每个文件都包含元数据 需要读取,而且可能比实际数据大。以128MB~1GB为目标进行打包。

分区遵循查询模式

s3://lake/events/dt=2026-09-06/hour=14/part-0001.parquet
                 └─ 날짜로 자르는 질의가 대부분이면 이렇게

如果选择错误的分区键,**所有问题都会覆盖整个内容。**相反,如果分区太细的话 出现小文件问题。计算一天有多少个文件出来后再确定。

日期dt=2026-09-06像这样用一个字符串放在一起是year=/month=/day=罗 通常比分割更容易。范围查询易于使用,目录深度也浅。

下次实习要做的事情

仅由字符串组成的staging.orders_raw对进行分析,找出缺陷,将日期、金额和状态规范化后,将数据分为精炼表和拒绝表,并记录下来,确认是否成立保存法则。