11. 批处理系统
派生数据
1. 组合的数据系统
前面我们讨论了分布式系统所有主要注意事项,但是这些讨论都只包含一个数据库。事实上数据系统是复杂的,通常需要以多种方式访问和处理数据,并且一个数据库往往无法同时满足所有不同的需求。因此应用程序需要使用若干不同的数据存储区、索引、高速缓存、分析系统等的组合,并实现数据从一个存储系统移动到另一个存储系统。
接下来我们将讨论如何将**不同数据系统(具有不同数据模型,并针对不同访问模式进行过优化)**整合至一致的应用程序体系结构中。
整合不同系统是大型应用中最为关键的任务之一。
1.1 记录系统与派生数据系统
存储和处理数据的系统可以分成两类:
- 记录系统: 真实数据系统,拥有数据的权威版本
- 派生数据系统:
- 从其他系统中获取已有数据并转换或处理而来
- 非规范化数值、索引、物化视图都属于派生数据类别
- 派生数据是冗余的,但是其对于获取良好的读取查询性能至关重要
作这样的区分是为了明确系统中的数据流: 明确系统中的某部分有哪些输入和输出,以及它们之间的依赖关系。通过弄清楚数据的来龙去脉,来帮助厘清复杂的系统架构。
2. 批处理系统
按照交互的不同方式,我们可以区分出三种不同类型的系统:
- 在线服务(在线系统)
- 服务等待客户请求,服务尽快处理并发回响应
- 响应时间通常是服务性能的衡量指标,而可用性同样重要
- 批处理系统(离线系统)
- 接收大量的输入数据,运行一个作业来处理输出,并产出数据
- 作业执行需要一定时间,用户通常不对等待作业完成
- 吞吐量是这类系统的衡量标准
- 流处理系统(近实时系统): 介于在线与离线系统之间
接下来我们将介绍 MapReduce 和其他一些批处理算法和框架。
3. MapReduce 和分布式文件系统
MapReduce 有点像分布在数千台机器上的 UNIX 工具,与 UNIX 命令行工具类似,这是一个相当直接、蛮力、却有效的神奇组件:
- 需要一个或多个输入,并产生一个或多个输出
- 不会修改输入,处理生成输出外没有任何副作用
- 输出文件以序列方式一次性写入
- UNIX 使用 stdin 和 stdout 作为输出和输出,MapReduce 作业在分布式文件系统上读写文件
3.1 分布式文件系统
分布式文件系统有很多,譬如:
- Hadoop 的 HDFS
- GlusterFS 和 Quantcast File System(QFS)
- 诸如 Amazon S3,Azure Blob 存储和 OpenStack Swift 对象存储服务也有很多相似之处
与网络连接存储 NAS 和 存储区域网络 SAN 架构的共享磁盘方法相比,HDFS 基于无共享原则。共享磁盘存储由集中式存储设备实现,通常使用定制硬件和特殊网络基础设施。而无共享方法则只需要通过传统数据中心网络连接的计算机。
HDFS 有以下几个部分组成:
- HDFS 的每台机器上运行着一个是守护进程,并会开放一个网络服务以允许其他节点访问存储在该机器上的文件
- 名为 NameNode 的中央服务器会跟踪哪个文件块存储在哪台机器上
- 考虑机器和磁盘的容错,文件块被复制在多台机器上
- 计算时计算任务会就近安排在存储所需文件的机器上,避免大量的数据传输成本(被称为: 将计算靠近数据)
3.2 MapReduce 的分布式执行
MapReduce 分为 map 和 reduce 两个过程,map 的输出始终会在排序之后在传递给 reducer:
- Mapper
- 每个输入记录都会调用一次 mapper 程序
- 任务是从输入记录中提取关键字和值
- 每个记录可以生成任意个键值对
- Reducer:
- MpaReduce 框架使用由 mapper 生成的键值对,收集属于同一个关键字的所有值,并使用迭代器调用 reducer 以使用该值集合。
MapReduce 的并行化是基于分区的,作业的输入通常是 HDFS 的一个目录。输入目录中的每个文件或文件块都被视为一个分区。每个分区由一个单独的 map 任务处理。MapReduce 调度器会尝试在输入文件副本的某台机器上运行 mapper 任务(计算靠近数据),避免数据的复制,提高访问的局部性。

如上图所示,分布式的计算过程分成了如下几个步骤:
- 代码分发: 将 map 任务所运行的应用程序代码分配到运行任务的节点上
- map 任务执行: 启动 map 任务,每次读取一条记录传递给回调函数 mapper,输出键值对
- 键值排序:
- 为了确保具有相同关键字的所有键值对都在相同的 reducer 任务中处理,框架使用关键字的哈希值来确定哪个 reducer 任务接收特定的键值对
- 键值对必须排序,排序是分阶段进行的
- 首先每个 map 任务都基于关键字哈希值,按照reducer对输出进行分块
- 每个分块都被写入 mapper 程序所在本地磁盘上的已排序文件,使用的计数类似于 SSTable 和 LVS-Trees
- shuffle:
- mapper 读取完输入并写入经过排序的输出文件后,reducer 与 mapper 相连接,并按照其分块从 mapper 中下载排序后的键值对文件
- reduce 从 mapper 获取文件并将它们合并在一起,同时保持数据的排序
- reducer 执行: reducer 通过关键字和迭代器进行调用,迭代器逐步扫描所有具有相同关键字的记录,并将输出结果写入 HDFS 中
3.3 MapReduce 工作流
单个 MapReduce 作业可以解决的问题有限,将 MapReduce 作业链接到工作流非常普遍。一个作业的输出将成为下一个作业的输入。Hadoop MapReduce 对工作流没有任何特殊的支持,作业的链接是通过目录名隐式完成的。每个MapReduce 作业的输出被写入到临时文件,下一个 MapReduce 命令从临时文件中读取(这被称为中间状态实体化)。
工作流中的一个作业只有在先前的作业成功完成时才能开始,为了处理这些作业执行之间的依赖关系,已经开发了各种 Hadoop 的工作流调度器,包括: Oozie,Azkaban, Luigi, Airflow 和 Pinball。Hadoop 的各种高级工具(如 Pig,Hive,Cascading,Crunch 和 FlumeJava) 则支持设置多个 MapReduce 阶段的工作流,这些不同的阶段会被恰当的自动链接起来。
3.4 Reduce 端的 Join 与分组
数据集通常存在关联,在批处理的背景下讨论 join,我们主要是解决数据集内在关联的所有事件。
join 的最简单实现是逐个遍历活动事件,并在远程服务器上的用户数据库中查询每个遇到的用户 ID。但是这一个方案性能非常差,原油是:
- 吞吐量受到数据库服务器的往返时间限制,本地缓存的有效性将很大程度上取决于数据的分布
- 同时运行的大量并行查询很容易是数据库不堪重负
- 远程数据库中的数据可能会发生变化,查询远程数据库意味着增加批处理作业的不确定性
因此更好的方式是获取用户数据库的副本(ETL),并将其放入与用户活动事件日志相同的分布式文件系统,并使用 MapReduce 进行数据合并。
排序-合并join
以用户 ID 为键执行的 reduce 端排序-合并join 如下图所示:

经过 mapper 对键进行排序分区的结果是所有活动事件和用户ID相同的用户记录在reducer 的输入中彼此相邻。MapReduce 作业甚至可以对记录进行排序,一边 reducer 会首先看到用户数据库中的记录,然后按时间戳排序查看活动事件。这种技术称为次级排序。经过排序和分区后,reducer 可以很容易执行真正的 join 逻辑。
上面这个算法被称为排序-合并-join,因为 mapper 的输出是按关键字排序的,然后 reducer 将来自join两侧的已排序记录列表合并在一起。mapper 和排序过程确保了所有需要的数据已经预先排列好了,所以 reducer 可以相当简单。
使用 MapReduce 编程模型将计算中的物理网络通信部分(从正确的机器获取数据) 从 应用逻辑(处理数据) 中分离出来。这种分离与数据库的典型使用形成鲜明对比: 从数据库中获取数据的请求经常发生在应用程序代码的深处。
由于 MapReduce 能够处理所有的网络通信,因此它也避免了在应用程序中处理局部故障,例如某个节点的崩溃: MapReduce 会在不影响应用程序逻辑的情况下透明的重试失败任务。
处理数据倾斜
3.5 map 段join
之前我们描述了 reduce 端 join 即 mapper 负责准备输入数据: 从每个输入记录中提取关键字和值,将键值对分配给 reducer 分区,并按关键字排序。
Reduce 端 join 方法的优点是不需要对输入数据做任何假设。然而不利的一面是,所有这些排序,复制到 reducer 以及合并 reducer 输入可能会是非常昂贵的操作,这取决于可用缓冲区。
另一方面如果可以对输入数据进行某些假设,则可以通过使用所谓的 map 端join 来加快速度。这种方法使用了缩减版的 MapReduce ,其中没有 reducer,也没有排序。相反每个 mapper 只需从分布式文件系统中读取输入文件块,然后处理输出即可。
map 端 join 实现有如下几种方法:
- 广播哈希 join:
- 适用: 大数据集和小数据集join,尤其是小数据集可以全部加载到每个 mapper 内存中
- 应用: pig(replicated join),hive(MapJoin)
- 分区哈希 join
- 适用: 两个 join的输入具有相同数量的分区
- 算法: 根据相同的关键字和相同的哈希函数将记录分配至分区中
- 应用: hive 中的 bucketed map join
3.6 对比
map 端或 reduce 端join 的不同选择会影响到输出结构:
- reduce 端 join 的输出按 join 关键字进行分区和排序
- map 端 join 的输出按照与大数据集相同的方式进行分区和排序
正如讨论的, map 端 join 也存在对输入数据集的大小、排序和分区方面的假设。在优化 join 策略时,了解分布式文件系统中的物理布局非常重要: 仅仅知道编码格式和数据存储目录的名称是不够的,还必须知道数据分区数量,以及分区和排序的关键字。
4. 批处理系统的优缺点
4.1 批处理系统与分布式数据库
4.2 中间状态实体化
4.3 批处理系统与流处理系统
4.4 图与迭代处理
5. 总结
分布式批处理框架需要解决两个主要问题:
- 分区:
- MapReduce 中 mappper 根据输入文件进行分区,mapper 的输出被重新分区、排序、合并成一个可配置数量的 reducer 分区
- 目的是把所有相关数据放在一起
- 除非必要,MapReduce 的数据流引擎都尽量避免排序,但它们采用了类似的分区方法
- 容错:
- MapReduce 需要频繁写入磁盘,这使得可以从单个失败任务中轻松恢复,而无需重启整个作业,但在无故障下会减慢执行速度
- 数据流引擎执行较少的中间状态实体化并保留更多内存,这意味着如果节点出现故障,他们需要重新计算更多的数据。确定性运算符减少了需要重新计算的数据量
分布式批处理引擎有一个有意限制的编程模型: 回调函数(mapper 和 reducer) 被设定为无状态,并且除了指定输出之外没有外部可见的任何副作用。这个限制使得框架隐藏了抽象背后的一些困难的分布式问题,从而在面对崩溃和网络问题时,可以安全的重试任务,并丢弃任何失败任务的输出。
得益于这样的框架,批处理作业中的代码无须考虑容错机制: 框架可以保证作业的最终输出与没有发生错误的情况相同。而在线服务在处理用户请求时,将写入数据库作为请求的副作用,与之相比批处理的可靠性语义要强大的多。
批处理系统的显著特点是读取一些输入数据,至关重要的是输入数据是有界的: 数据大小固定已知。因为有界所以总是可以知道作业何时结束。在下一章中我们 将转向流处理,其输入是无界的。即作业的输入是永无止境的数据流。我们将看到流处理与批处理有很多相似之处,但是流数据无界的假设也会深刻改变我们设计系统的方法。