全文共 3,730 字 预计阅读 11 分钟
bg

bg10.5.小复习

我们又回来辣

前面我们已经把整个大数据的基本的这些组件过了一遍,现在来进行一个整理吧。本篇涵盖的内容包括:hdfs、对象存储、parquet、mr、spark、kafka、flink、hive、iceberg、yarn、airflow、olap 查询引擎

1、hdfs

hdfs 是 分布式文件系统,它把一个大文件切成固定大小的 block(128MB),每块存多副本,分散在多台机器上。其中,namenode 是常驻在 内存的 jvm 堆,存储元数据,datanode 也是 jvm 进程,负责读、写、存储对应的 block,并定期向 namenode 发心跳和块汇报。一般来说 hdfs 是 3 副本,一个放入 写入客户端的 datanode,另两个放在另一机架的两个不同节点上。

hdfs 的特性就是一次写入,追加为主,这样就可以顺序读写。要随机读写,一般还是要 hbase(nosql)或者 oltp。

hdfs 主要是为了解决 单台机器的 磁盘容量上限不足、容易丢数据、io 吞吐有限 问题

hdfs 最大的问题就是 namenode 内存容易耗尽

2、对象存储

对象存储的本质就是 能用 http 接口访问的 分布式 kv 系统 ,k 就是对象的路径,v 就是对象存储的内容。相比 hdfs,对象存储没有 namenode 集中存储元数据、用纠偏码替代 3 副本,来解决 hdfs 强依 datanode、运维重,小文件压力大、机架副本机制耦合集群,公有云难以托管的痛点。实现了:无限水平扩容、非 / 半结构化数据低成本持久归档、脱离 hadoop 只要联网就可以读写。

对象存储的问题就在于 list 操作(给定前缀,查询这个前缀下所有对象的 key,并没有真目录,明显慢于 hdfs 的目录遍历)比较慢、rename 操作非原子(把旧 key 的数据复制到新 key -> 删除原始旧对象)。

3、parquet

parquet 就是列存,是为了解决行存:IO 浪费、压缩效果差(不同类数据一起压缩)、无法谓词下推 的三大痛点。列存通过:

  • 只读需要的列
  • 同列同类型,高压缩比
  • 编码(重复值建立字典,存储值变编码而非值 + 连续相同的值存值 * 重复次数 ) + 压缩(对编码后的字节加通用压缩算法 snappy gzip 这种)
  • 把文件切成 row group + footer(记录 row group 的元数据),记载每个 group 的 min / max,来实现行组裁剪(谓词下推 就是不把不要的数据读到内存之后过滤,不全量从磁盘读数据到内存)。row group 对应 spark 的一个 读 task,column chunk 是一个 row group 内一块连续存放的数据,page 是 column chunk 再切成页,是 解压缩 的最小单元。 parquet 的问题在于:
  • 每个 parquet 都有 footer,就会有小文件爆炸
  • row group 设定的 太大或太小 会有并行度低 / 压缩差、元数据爆炸
  • schema 演进受限,不能改 列名和数据类型,因为 footer 只读,最终只能改 iceberg 的 catalog

4、mapreduce

mapreduce 解决的痛点就是 有了 分布式文件存储,那就需要有个引擎去把 数据的切分、任务分发到哪个机器、机器挂了怎么办、中间结果怎么在机器间传输 这些任务给抽象调。

mr 就是通过 map 来读 kv,shuffle 来聚合 kv、reduce 来做计算,实现 程序员只要关注 map 和 reduce 两个函数,分发、容错、shuffle 由框架兜底。

mr 的问题在于:

  • 所有作业结果存 hdfs,io 负担,迭代型作业 io 直接爆炸
  • 什么都得写成 map 和 reduce 两个函数,复杂逻辑代码冗长
  • 容错靠 任务级重跑

5、spark

spark 就是为了解决这个痛点而提出的,它的解决方案是:

  • 用 dag 对 pipeline 规划,必要处才 shuffle
  • 中间结果放 executor 内存,读内存很快
  • 提供多种算子的链式组合,编程方便
  • 记录 rdd 血缘,丢了只重算丢失的分区

rdd 特性:

  • 本质是各 executor 的一批分区,一个分区就是一个 task 处理单位
  • rdd 不可变,任何操作生成新的 rdd 代价: 窄依赖重跑便宜,宽依赖重跑整个 shuffle,触发大面积重算,解决是直接 checkpoint 落盘 hdfs,让血缘恢复回到 checkpoint 恢复。

其中,窄依赖就是 rdd 的父亲分区只有一个,宽依赖就是一个 rdd 有多个父分区。

spark 是通过 Driver 进程(dag、切 stage、生成 task、派给 executor,汇总结果),executor(jvm 进程,用线程池并发跑多个 task)

stage 是任务调度的基本单元,按照 shuffle(宽依赖) 来切分的。 shuffle 的减少是 spark 调优的第一要义:shuffle 会有需要 落盘开销、网络 IO 开销、排序序列化开销。

catalyst 负责生成 spark 的执行计划。aqe 可以在运行时根据每个 shuffle 真是产出的数据大小动态调整计划,避免某个 key 分区过慢的木桶效应。

spark 的问题就是:shuffle 分区数过少/多 会有单分区多大,oom / 小文件、不适合低延迟 / 流处理

6、kafka

kafka 主要是为了解决以下痛点而诞生的:1、海量日志,传统 MQ 吞吐不足 2、多下游消费需要频繁复制 3、消息丢失、故障恢复、日志回溯 4、难以集群化扩容

kafka 的本质就是 分布式追加日志,通过磁盘顺序写来保证高吞吐;通过消息不删,保障可回溯;通过消费者自己记录 offset 来保证多方独立消费互不干扰。当然,这种架构的代价,就是有序性被限制在了分区,不适用低延迟点对点的场景,精确一次的语义比较复杂。

kafka 的抽象概念比如 producer、consumer、consumer group、topic 这些,物理上,一个节点就是一个 broker;topic 被分为多个 partition。

一个 partition 可以有多个 replica,分布在不同 broker 上,分为 leader 和 follower。其中,isr 就是同步副本集,follower 超过落后 leader 时间的阈值就会被踢出。high watermark 就是 isr 中所有副本都已复制到的最小 offset。

kafka 的快来源于 pagecache(热数据直接读内存) + 批量 fsync、零拷贝(把日志从磁盘 -> 消费者,用 sendfile 让数据直接从 page cache -> 网卡,跳过了 内核 -> 用户态 -> 内核 的拷贝和反序列化「因为一般消息进入用户态就需要反序列化」)、批量发送 + 压缩。

最多一次:先提交 offset 再处理,这样会丢数据,但不重复。

至少一次:先处理再提交 offset,这样会重复,但是不丢失数据。

精确一次:依赖 幂等生产者和事务,eos 只在 kafka 生态内部闭环实现,比如 flink 的 kafkasink。

kraft 和 zk 主要负责的是 分布式调度,比如 broker 的管理、controller 的选举(负责 partition 的 replica 分配, isr...)、元数据的存储。kraft 主要是消除了 zk 的外部依赖。

rebalance 会在消费者加入 / 离开;订阅的分区数变化时,对消费者组内 分区 -> 消费者 的分配关系进行重排,但当消费者一批消息处理太久,超过阈值,coordinator 就会把它踢出消费者组,出发 rebalance,但等它下次 poll 又会要加回来,最终导致 rebalance 风暴。

7、flink

flink 的出现主要是为了解决 流式计算中 低延迟 + 强一致 + 时间语义 storm 和 spark streaming 不能三大的痛点。storm 做到了逐条处理和低延迟,但是容错语义只能到最少一次,同一条数据会被处理多次,下游会有重复结果。spark streaming 的问题就在于它天然按照 数据到达的时间切分,难以处理乱序。

flink 集群有三部份组成:Client(负责提交作业)、JobManager(Dispatcher 托管 web ui,为每个作业启动 JobMaster、ResourceManager 管理 task slot、JobMaster 管理每个 JobGraph 的执行、 Checkpoint Coordinator 通知 source 插入 barrier)、TaskManager(负责执行 task)

task slot 是 flink 中资源调度的最小单元,隔离内存,但共享 cpu。

flinksql / datastream -> streamgraph(客户端侧的算子拓扑) -> jobgraph(算子合并,窄依赖让同一个线程,避免 落盘开销)-> executiongraph(并行版本)->物理执行

时间语义:event time、processing time、ingestion time(进入 flink source)。watermark 就是插在数据流中的一个标记,含义是 event time <=t 的 数据基本已经到达了,后续不应该再出现 event time <=t 的 数据了。

watermark = (到目前为止观察到的最大事件时间) - (允许的最大乱序时长)

在一个时间窗口 [start,end) 中,当 watermark 推进到 >= end,就判定 end 为止的数据都到齐了,于是触发窗口计算,输出结果,清理窗口状态。此后迟到的数据都默认丢弃。

state backend 决定状态(也就是中间计算结果)存在 内存 hashmap / 磁盘(序列化后存在 内嵌 rocksdb)

checkpoint 会周期性地对所有算子(source、transform、sink)的状态进行全局一致的快照,持久化到 oss / hdfs 上面。作业失败时,把所有算子的状态回滚到最新一次成功的快照,让 source 重放该快照点之后的数据。savepoint 手动触发,一般用于运维。

如何在不停数据流的时候,拍下一个全局一致的快照:通过 jobmanager 中的 checkpoint coordinator 周期性通知 source,在数据流里插入 编号为 n 的 barrier,每个算子一收到 barrier n,就对自己当前的状态做快照,把 barrier n 转发给下游,当全部流到 sink 的时候, 编号为 n 的 checkpoint 就完成了。然而,多输入算子就需要等待所有的 barrier 到达通道,这就叫对齐,它会引入等待延迟。

8、hive

hive 的出现是为了解决 mr 和 hdfs 刚出现时需要 java 编程重复造轮子的痛点,它通过统一封装通用计算逻辑避免重复写 mr、将 hdfs 抽象为结构化表来通过 sql 操纵海量数据。

9、iceberg

iceberg 主要是为了解决 hive 的几个痛点:没有 acid、在对象存储上 rename 不原子、分区靠目录,list 很慢、没有 schema / 分区 演进。

iceberg 本质就是通过 元数据文件 把一批 parquet 文件组织成带事务语义的 表。

具体来说,iceberg 是一个树:catalog(本质就是一个指针,表名到 metadata 的指针) -> metadata(一个 json,记录表的 schema、分区规格、所有历史快照的列表、快照 id) -> manifest list(每个快照对应一个,列出本快照由哪些 manifest 组成,为每个 manifest 记录分区取值范围,便于跳过) -> manifest(列出 数据 / 删除文件,为每个文件记录:属于哪个分区、有多少行、各自的 min/max、null 数,这里取代了 LIST) -> datafile

snapshot 表示某一时刻,这张表由哪些数据文件构成。为了解决 rename 的不原子问题,iceberg 选择将书记处文件、manifest、manifest list、metadata 都提前写好,最后只需要通过 cas 将把 catalog 中表名 -> metadata 的指针原子的替换掉即可,这样就可以让文件原子的出现。

iceberg 通过把 “向对象存储发成千上万的 LIST http” 替换为 读几个 manifest 文件,文件清单从原本运行时列目录,变成提交时写在元数据中 来解决 LIST 慢的问题。又通过文件级裁剪来解决小文件定位的问题。

同时,hive 要专门写一个 dt 字段,iceberg 直接按照 day(ts) 来进行分区,查询时对 ts 进行自动的 days(ts) 的变化来推导出要扫的分区。分区演进也不会动历史数据,老数据依旧按照之前的 day 分区,新写入的按照 hour 分区,查询规划时会对不同 spec 的文件分别裁剪。

同时,iceberg 为每列分配一个一个 id,读写按照 id 对应,改名就是 改 id -> 名字的映射,数据文件不动。加列就是老文件没有这个 id 的数据,那就返回 null。这样也就支持了 schema 演进。

当然,iceberg 也会有问题:

  • 依旧有小文件问题,需要 compaction
  • 元数据从会膨胀
  • 多引擎/高并发写 因为 cas 会出现大量反复的重试

10、yarn

yarn 的出现是为了解决 旧版 hadoop 中 jobtracker 既需要管理全集群资源,又需要调度和监控每一个 mr 任务。yarn 就是通过把资源管理和作业的生命周期管理拆开,分为 resourcemanager(分为 scheduler 来按策略地分配资源, applicationsmanager 受理作业提交,为每个作业启动container 来跑 am)、nodemanager(每个机器一个,定期向 rm 心跳,按命令启动监控 container,管理 cpu/内存)、applicationmaster(每个作业一个,向 rm 申请资源拿到 container 后联系 nm 拉起任务进程,作业结束就退出)、container(资源的抽象单位,实际计算跑在 container)

11、airflow

airflow 的出现是为了解决 linux 中 crontab 定时 + shell 脚本看不清楚 作业与作业之间的依赖关系、重跑复杂、失败无兜底的痛点。

airflow 就是通过 把任务组织成 dag,将任务抽象为节点,依赖抽象为节点间的有向边。它的核心组件包括:scheduler(调度器,不断解析 dag,判断哪些任务此刻能跑)、executor(把该跑的任务真正发出去执行)、metadata db(存所有任务的状态)

12、olap 查询引擎

olap 查询引擎面向分析,其查询有这样几种特点:触碰几乎所有行、但少数几列,并发低,事务弱。标准:几十亿行,秒级别返回。olap 主要是通过 mpp + 向量化来实现的。

mpp 就是把数据切成片,分给很多机器,让他们同时算,最终汇总。向量化执行就是查询算子每次处理一批而不是一行。

对于一个查询,首先有 列裁剪、分区裁剪、与编码压缩,把要读的数据砍到最小,然后有 MPP 把查询拆到多节点并行,节点间 exchange,再有向量化,按列将数据向量化,降低传统火山模型(每行都要走一遍算子链的函数调用、算子间传递)函数调用的开销,并将数据按顺序进 L1/L2 cache,吃满 cpu 缓存,最后还可以用上 simd,一条 cpu 指令同时处理一个向量寄存器里的多个值。这三条促成了 olap 查的快。

同时,通过跳数索引,每个数据块存 min/max,与预聚合,来实现 跳过算、或者说提前算。

存算分离的架构:Server 层负责解析、优化、生成分布式计划;然后下发 fragment 交给计算组进行 MPP + 向量化;存储层对象存储通过 parquet + partition 方便计算组读取。

Back to Blog