84 KiB
title, book_kind, book_number, book_part, weight, breadcrumbs
| title | book_kind | book_number | book_part | weight | breadcrumbs |
|---|---|---|---|---|---|
| 批处理 | chapter | 11 | III | 311 | false |
带有太强个人色彩的系统无法成功。当最初的设计完成并且相对稳健时,不同的人开始以自己的方式进行试验,真正的考验才开始。
高德纳
到目前为止,本书的大部分内容都在讨论 请求(request)和 查询(query),以及相应的 响应(response)或 结果(result)。许多现代数据系统都默认采用这种处理方式:你请求某样东西,或者发出一条指令,系统便尽量快速地给出答案。
浏览器请求网页、服务调用远程 API,以及数据库、缓存、搜索索引等许多系统,都是这样工作的。我们称它们为 在线系统(online system)。这类系统通常以响应时间作为主要性能指标,而且往往需要具备容错能力,才能保证高可用性。
然而,有些计算规模太大,或者需要处理的数据太多,无法放在一次交互式请求中完成。例如,你可能需要训练 AI 模型,把大量数据从一种形式转换成另一种形式,或者在非常大的数据集上进行分析。我们把这类任务称为 批处理(batch processing)作业,相应的系统有时也称为 离线系统(offline system)。
批处理作业读取一组输入数据(只读),并产生一组输出数据(每次运行都从头生成)。它通常不会像读写事务那样修改现有数据。因此,输出是由输入 衍生(derived)而来的(参见“权威记录系统与衍生数据”):如果对输出不满意,只需把它删除,调整作业逻辑,再运行一次。把输入视为不可变数据,并避免产生副作用(例如写入外部数据库),不仅能让批处理作业获得良好的性能,还会带来其他好处:
-
如果代码中引入了错误,导致输出有误或遭到破坏,只要回滚到先前版本的代码并重新运行作业,输出便能恢复正确。更简单的办法是把旧输出保存在另一个目录中,需要时直接切换回去。大多数对象存储和开放表格式(参见“云数据仓库”)都支持这种称为 时间旅行(time travel)的功能。大多数支持读写事务的数据库却不具备这种性质:如果错误代码把坏数据写进数据库,回滚代码并不能修复已经写入的数据。这种从错误代码中恢复的能力称为 容忍人为失误(human fault tolerance)1 。
-
由于回滚很容易,功能开发可以比“犯错就会造成不可逆损害”的环境推进得更快。这种 尽量减少不可逆操作 的原则有利于敏捷软件开发 2 。
-
同一组文件可以作为许多不同作业的输入,其中也包括监控作业:它们计算指标,并检查某项作业的输出是否具备预期特征,例如将其与上一次运行的输出比较,衡量两者之间的差异。
-
批处理框架能够高效利用计算资源。虽然 OLTP 数据库和应用服务器等在线数据系统也能成批处理数据,但完成同样工作所需的资源可能昂贵得多。
批处理也会带来一些挑战。在大多数框架中,只有整个作业运行完毕,其他作业才能继续处理它的输出。批处理还可能效率不高:输入数据发生任何变化——哪怕只有一个字节——都意味着批处理作业必须重新处理整个输入数据集。尽管存在这些局限,批处理仍在许多场景中证明了自己的价值,我们将在“批处理用例”中再次讨论这些场景。
一项批处理作业可能要运行很长时间:几分钟、几小时,甚至几天。作业也可能按固定周期调度运行,例如每天一次。其主要性能指标通常是吞吐量,即单位时间内能够处理多少数据。有些批处理系统遇到故障时只会中止并重新启动整个作业;另一些则具备容错能力,即使某些节点崩溃,作业也能顺利完成。
Note
批处理之外还有一种选择,即 流处理。流处理作业不会在处理完当前输入后结束,而是继续监视输入,并在输入发生变化后不久加以处理。我们将在第 12 章讨论流处理。
在线系统与批处理系统之间的界限并不总是分明:一条长时间运行的数据库查询看起来就很像批处理过程。不过,批处理有一些独特的性质,使其成为构建可靠、可伸缩且可维护应用的重要构件。例如,它经常用于 数据集成,也就是把多个数据系统组合起来,完成单个系统无法独自完成的工作。“数据仓库”中讨论的 ETL 就是一例。
现代批处理深受 MapReduce 影响。Google 于 2004 年发表了这种批处理算法 3 ,随后 Hadoop、CouchDB 和 MongoDB 等多个开源数据系统都实现了它。MapReduce 是一种相当底层的编程模型,没有数据仓库中的并行查询执行引擎那么精巧 4 5 。刚问世时,MapReduce 让普通商用硬件能够达到的处理规模迈上了一个新台阶;如今它已经基本过时,Google 也不再使用 6 7 。
如今,批处理更多由 Spark、Flink 之类的框架或数据仓库查询引擎完成。与 MapReduce 一样,这些系统高度依赖分片(参见第 7 章)和并行执行,但它们的缓存和执行策略精巧得多。随着系统逐渐成熟,运维方面的问题基本得到解决,关注点也转向了易用性。数据流 API、查询语言和数据框 API 如今都得到了广泛支持,作业与工作流编排同样日趋成熟。Oozie、Azkaban 等以 Hadoop 为中心的工作流调度器,已经被 Airflow、Dagster 和 Prefect 等更通用的方案所取代;后者支持各种批处理框架和云数据仓库。
云计算已经十分普及,批处理的存储层也正从 HDFS、GlusterFS 和 CephFS 等分布式文件系统(DFS)转向 S3 之类的对象存储。BigQuery、Snowflake 等可伸缩的云数据仓库,则进一步模糊了数据仓库与批处理系统之间的界限。
为了直观地理解批处理,本章先从单台机器上的标准 Unix 工具入手,再研究如何把数据处理扩展到分布式系统中的多台机器。我们将看到,分布式批处理框架与操作系统很相似,同样拥有调度器和文件系统。随后,我们会考察编写批处理作业时使用的几种处理模型,最后讨论常见的批处理用例。
使用 Unix 工具的批处理
假设有一台 Web 服务器,每处理一个请求,它都会在日志文件末尾追加一行。以 nginx 默认的访问日志格式为例,其中一行可能如下所示:
216.58.210.78 - - [27/Jun/2025:17:55:11 +0000] "GET /css/typography.css HTTP/1.1"
200 3377 "https://martin.kleppmann.com/" "Mozilla/5.0 (Macintosh; Intel Mac OS X
10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/137.0.0.0 Safari/537.36"
(这其实是一行,只是为了便于阅读才在这里折成了多行。)这一行包含许多信息。要解释它,需要先查看日志格式的定义:
$remote_addr - $remote_user [$time_local] "$request"
$status $body_bytes_sent "$http_referer" "$http_user_agent"
这行日志表示:在 UTC 时间 2025 年 6 月 27 日 17:55:11,服务器收到了来自客户端 IP 地址 216.58.210.78、请求文件 /css/typography.css 的请求。用户没有经过身份认证,因此 $remote_user 被设为连字符(-)。响应状态码为 200(即请求成功),响应大小为 3,377 字节。浏览器是 Chrome 137;它之所以加载这个文件,是因为网址 * https://martin.kleppmann.com/* 的页面引用了该文件。
解析日志听起来或许像是一个刻意编造的例子,实际上却是许多现代科技公司的关键工作,从广告数据管道到支付处理无所不包。事实上,日志处理正是 MapReduce 得以迅速普及、并推动“大数据”浪潮的重要原因之一。
简单日志分析
许多工具都能读取这些日志文件,生成漂亮的网站流量报告。不过为了练习,我们来用基本的 Unix 工具自行构建一个。假设你想找出网站上最受欢迎的五个页面,可以在 Unix shell 中执行下面的命令:
cat /var/log/nginx/access.log | #1
awk '{print $7}' | #2
sort | #3
uniq -c | #4
sort -r -n | #5
head -n 5 #6
-
读取日志文件。(严格来说,这里的
cat并非必需,因为可以把输入文件直接作为参数传给awk。不过这样写能让线性管道显得更清楚。) -
按空白字符把每一行拆成字段,只输出第七个字段,而它恰好就是请求的 URL。在前面的示例中,这个 URL 是 /css/typography.css。
-
按字母顺序
sort请求 URL 列表。如果某个 URL 被请求了 n 次,那么排序后的文件中就会连续出现 n 个相同的 URL。 -
uniq命令通过检查相邻两行是否相同,滤掉输入中重复的行。-c选项让它同时输出计数器:对于每个不同的 URL,它会报告该 URL 在输入中出现了多少次。 -
第二个
sort按每行开头的数字(-n)排序,也就是按 URL 的请求次数排序;然后以逆序(-r)返回结果,让最大的数字排在最前面。 -
最后,
head只输出输入的前五行(-n 5),并丢弃其余内容。
这一系列命令的输出大致如下:
4189 /favicon.ico
3631 /2016/02/08/how-to-do-distributed-locking.html
2124 /2020/11/18/distributed-systems-and-elliptic-curves.html
1369 /
915 /css/typography.css
如果你不熟悉 Unix 工具,前面的命令行可能显得有些晦涩,但它的能力非常强大。它能在几秒钟内处理数 GB 的日志,而且可以轻松修改分析方式来满足需要。例如,如果不希望报告中包含 CSS 文件,只需把 awk 参数改为 '$7 !~ /\.css$/ {print $7}';如果想统计最常出现的客户端 IP 地址,而不是最常访问的页面,则把参数改为 '{print $1}',以此类推。
本书没有足够篇幅详细介绍 Unix 工具,但它们非常值得学习。令人惊讶的是,仅用 awk、sed、grep、sort、uniq 和 xargs 的某种组合,几分钟内就能完成许多数据分析,而且性能也出奇地好 8 。
命令链与自定义程序
除了使用 Unix 命令链,你也可以写一个简单的程序来完成同样的工作。例如,Python 程序可能如下所示:
from collections import defaultdict
counts = defaultdict(int) #1
with open('/var/log/nginx/access.log', 'r') as file:
for line in file:
url = line.split()[6] #2
counts[url] += 1 #3
top5 = sorted(((count, url) for url, count in counts.items()), reverse=True)[:5] #4
for count, url in top5: #5
print(f"{count} {url}")
-
counts是一个散列表,保存每个 URL 出现次数的计数器;计数器的默认值为 0。 -
从每一行日志中取出按空白字符分隔的第七个字段,作为 URL(因为 Python 数组从 0 开始计数,所以数组下标是 6)。
-
把当前日志行中 URL 对应的计数器加一。
-
按计数器的值对散列表内容降序排列,并取出前五项。
-
输出这五项。
这个程序没有 Unix 管道那么简洁,但也相当容易读懂;喜欢哪一种,部分取决于个人偏好。不过,除了表面的语法差异,两者的执行流程也大不相同。如果在大文件上运行这项分析,区别就会显现出来。
排序与内存聚合
Python 脚本在内存中维护一个 URL 散列表,把每个 URL 映射到它出现的次数。Unix 管道没有这样的散列表,而是依靠排序 URL 列表;同一个 URL 出现多次时,它只是在列表中重复多次。
哪种方法更好?这取决于不同 URL 的数量。对于大多数中小型网站,你大概可以把所有不同的 URL 及其计数器放进 1 GB 左右的内存。这个作业的 工作集(working set;即作业需要随机访问的内存量)只取决于不同 URL 的数量:即使日志中同一个 URL 出现了一百万次,散列表所需的空间仍然只是一个 URL 加一个计数器。只要工作集足够小,内存散列表就能很好地工作——即便在笔记本电脑上也是如此。
另一方面,如果作业的工作集大于可用内存,排序方法就有一个优势:它可以高效利用磁盘。这与“日志结构存储”中讨论的原理相同:先在内存中对数据块排序,并把它们作为段文件写入磁盘,再将多个有序段合并成一个更大的有序文件。归并排序采用顺序访问模式,在磁盘上表现很好(参见“顺序与随机写入”)。
GNU Coreutils(Linux)中的 sort 工具会把无法装入内存的数据自动溢写到磁盘,还会自动利用多个 CPU 核并行排序 9 。这意味着前面的简单 Unix 命令链可以轻松扩展到大型数据集,而不会耗尽内存。此时的瓶颈很可能是从磁盘读取输入文件的速度。
Unix 工具的局限在于,它们只能在单台机器上运行。如果数据集大到无法装入单机内存或本地磁盘,就会遇到问题——这正是分布式批处理框架的用武之地。
分布式系统中的批处理
运行前面 Unix 工具示例的机器,需要由几个组件协同处理日志数据:
-
通过操作系统的文件系统接口访问的存储设备。
-
决定进程何时运行、如何为其分配 CPU 资源的调度器。
-
一系列 Unix 程序,其
stdin和stdout通过管道连接在一起。
分布式数据处理框架中也存在同样的组件。事实上,你可以把 分布式处理框架(distributed processing framework)看作一种分布式操作系统:它们拥有文件系统和作业调度器,程序则通过文件系统或其他通信通道相互发送数据。
分布式文件系统
操作系统提供的文件系统由若干层组成:
-
最底层的块设备驱动程序直接与磁盘通信,让上层能够读写原始数据块。
-
块设备层之上是页缓存,它把最近访问的数据块保存在内存中,以加快再次访问的速度。
-
文件系统层封装了块 API,把大文件拆分成块,并维护 inode、目录和文件等元数据。例如,ext4 和 XFS 是 Linux 上常用的两种实现。
-
最后,操作系统通过名为虚拟文件系统(VFS)的统一 API,把不同文件系统暴露给应用。无论底层采用哪一种文件系统,应用都可以用同样的方式读写数据。
分布式文件系统(distributed filesystem)的工作方式与此非常相似。文件同样被拆分成块,只不过这些块分布在许多机器上。分布式文件系统的块通常比本地文件系统大得多:HDFS(Hadoop 分布式文件系统)的默认块大小为 128 MB,JuiceFS 和许多对象存储使用 4 MB 的块,而 ext4 的块只有 4,096 字节。块越大,需要跟踪的元数据就越少;对于 PB 级数据集,这种差异十分显著。相对于读取数据块所需的时间,大块也能降低寻道开销所占的比例。
大多数物理存储设备无法写入不完整的数据块,因此即使数据没有填满一个块,操作系统也必须为写入使用整个块。分布式文件系统的块更大,而且通常构建在操作系统文件系统之上,所以没有这种要求。例如,一个 900 MB 的文件采用 128 MB 的分块时,会由 7 个占用 128 MB 的块和 1 个占用 4 MB 的块组成。
读取分布式文件系统中的块,需要向集群中存放该块的机器发出网络请求。每台机器都运行一个守护进程,并公开一套 API,让远程进程能够读写以文件形式存储在其本地文件系统中的数据块。HDFS 把这些守护进程称为 DataNode,GlusterFS 则称为 glusterfsd。本书统一把它们称为 数据节点(data node)。
分布式文件系统还实现了分布式版本的页缓存。由于数据块以文件形式存放在数据节点上,读写操作会经过每个数据节点的操作系统,其中就包含内存页缓存。因此,经常读取的数据块会被缓存在数据节点的内存中。有些分布式文件系统还实现了更多缓存层,例如 JuiceFS 提供的客户端缓存和本地磁盘缓存。
ext4 和 XFS 等文件系统会跟踪空闲空间、文件块位置、目录结构和权限设置等存储元数据。分布式文件系统同样需要记录文件分散在哪些机器上、具有什么权限等信息。Hadoop 通过 NameNode 服务维护集群元数据;DeepSeek 的 3FS 则使用元数据服务,并把数据持久化到 FoundationDB 之类的键值存储中。
文件系统层之上是 VFS;在批处理中,与之最接近的是分布式文件系统的协议。分布式文件系统必须公开某种协议或接口,让批处理系统能够读写文件。它充当一套可插拔接口:只要实现该协议,任何分布式文件系统都可以接入。例如,MinIO、Cloudflare R2、Tigris 和 Backblaze B2 等许多存储系统都采用了 Amazon S3 API。支持 S3 的批处理系统,也就能够使用其中任何一种存储。
有些分布式文件系统提供与 POSIX 兼容的接口,在操作系统的 VFS 看来,它们与其他文件系统并无不同。这类系统通常通过用户空间文件系统(FUSE)或网络文件系统(NFS)协议接入 VFS。NFS 或许是最著名的分布式文件系统协议。它最初的用途,是让多个客户端读写单台服务器上的数据。近来,AWS Elastic File System(EFS)和 Archil 等文件系统提供了可伸缩得多、但仍与 NFS 兼容的分布式实现。NFS 客户端依旧只连接一个端点,不过这些系统会在内部与分布式元数据服务和数据节点通信,完成数据读写。
[!TIP] 分布式文件系统与网络存储 分布式文件系统以 无共享 原则为基础(参见“共享内存、共享磁盘与无共享架构”),不同于网络附加存储(NAS)和存储区域网络(SAN)架构采用的 共享磁盘 方法。共享磁盘存储由集中式存储设备实现,往往需要定制硬件和光纤通道等专用网络基础设施。相比之下,无共享方法不需要特殊硬件,只需用普通的数据中心网络连接计算机即可。
许多分布式文件系统构建在普通商用硬件上。这种硬件比较便宜,但故障率也高于企业级硬件。为了容忍机器和磁盘故障,文件块会复制到多台机器上。这样也便于调度器更均匀地分配工作负载,因为任务可以在任何存有其输入数据副本的节点上运行。这里的复制可以像第 6 章所述,在多台机器上保存若干份完整副本;也可以采用 Reed–Solomon 码等 纠删码(erasure coding),以低于完整复制的存储开销恢复丢失的数据 10 11 12 。这些技术与 RAID 很相似,后者在连接到同一台机器的多个磁盘之间提供冗余;区别在于,分布式文件系统通过普通的数据中心网络访问和复制文件,不需要特殊硬件。
对象存储
Amazon S3、Google Cloud Storage、Azure Blob Storage 和 OpenStack Swift 等 对象存储(object storage)服务,已经成为批处理作业中分布式文件系统的常用替代方案。事实上,两者的界限有些模糊。正如上一节和“以对象存储为后端的数据库”中所述,用户空间文件系统(FUSE)驱动程序可以让用户把 S3 之类的对象存储当作文件系统使用。JuiceFS 和 Ceph 等分布式文件系统实现也同时提供对象存储与文件系统 API。不过,两类系统的 API、性能和一致性保证可能大相径庭。即使某个系统看似实现了所需的 API,采用之前也必须仔细确认它的实际行为符合预期。
对象存储中的每个对象都有一个 URL,例如 s3://my-photo-bucket/2025/04/01/birthday.png。URL 的主机部分(my-photo-bucket)表示存放对象的存储桶(bucket),后面的部分则是对象的 键(本例中为 /2025/04/01/birthday.png)。存储桶的名称在全局范围内唯一,而每个对象的键在所属存储桶内必须唯一。
对象通过 get 调用读取,通过 put 调用写入。与文件系统中的文件不同,对象写入后是不可变的;若要更新,只能像键值存储一样,通过 put 完整重写整个对象。Azure Blob Storage 和 S3 Express One Zone 支持追加写入,但大多数对象存储都不支持。对象存储也没有 fopen、fseek 之类的文件句柄 API。
对象看起来似乎按目录组织,但这有些误导,因为对象存储根本没有目录的概念。路径结构只是一种约定,斜杠本身也是对象键的一部分。这种约定允许你按特定前缀请求对象列表,实现类似目录列表的效果。不过,按前缀列出对象与文件系统列出目录有两点区别:
-
前缀
list操作类似于 Unix 系统上的递归ls -R:它会返回所有以该前缀开头的对象,其中也包括子路径下的对象。 -
对象存储中不可能存在空目录。假如删除
s3://my-photo-bucket/2025/04/01下的所有对象,那么对s3://my-photo-bucket/2025/04调用list时,01就不会再出现。常见的做法是用一个零字节对象来表示空目录,例如创建空对象s3://my-photo-bucket/2025/04/01,这样即使所有子对象都已删除,它仍会保留下来。
分布式文件系统通常支持硬链接、符号链接、文件锁和原子重命名等常见文件操作,对象存储则没有这些功能。它们一般不支持链接和锁,重命名也不是原子的,而是先把对象复制到新键,再删除旧对象。如果想重命名一个“目录”,就必须逐一重命名其中的每个对象,因为目录名本身是对象键的一部分。
第 4 章讨论的键值存储针对小值(通常只有几 KB)以及频繁、低延迟的读写进行了优化。相比之下,分布式文件系统和对象存储通常针对大对象(数 MB 至数 GB)和频率较低、规模较大的读取进行了优化。不过最近,对象存储也开始支持更加频繁、规模更小的读写。例如,S3 Express One Zone 如今可以达到个位数毫秒延迟,其定价模型也更接近键值存储。
分布式文件系统与对象存储还有一个区别:HDFS 等分布式文件系统可以把计算任务安排在存有特定文件副本的机器上运行。任务就能直接读取本地文件,无需通过网络传输;如果任务的可执行代码远小于它要读取的文件,这样可以节省大量带宽。对象存储通常把存储与计算分开。这种做法可能消耗更多带宽,但现代数据中心网络的速度很快,因而往往可以接受。存储与计算解耦后,CPU、内存等计算资源还可以独立于存储容量进行伸缩。
分布式作业编排
操作系统的类比同样适用于作业编排。执行 Unix 批处理作业时,总要有某个组件真正运行 awk、sort、uniq 和 head 这些进程。它必须把一个进程的输出传到另一个进程的输入,为每个进程分配内存,在 CPU 上公平地调度并执行各进程的指令,实施内存和 I/O 边界,等等。在单台机器上,这些工作由操作系统内核负责;在分布式环境中,它们则由作业编排器承担。
批处理框架会向编排器的调度器发送运行作业的请求。启动作业的请求包含如下元数据:
-
要执行的任务数量;
-
每项任务所需的内存、CPU 和磁盘资源;
-
作业标识符;
-
访问凭据;
-
输入、输出数据等作业参数;
-
GPU 或磁盘类型等必要的硬件信息;
-
作业可执行代码所在的位置。
Kubernetes 和 Hadoop YARN(Yet Another Resource Negotiator)13 等编排器会结合这些信息与集群元数据,通过下列组件执行作业:
- 任务执行器(Task Executor)
-
集群中的每个节点都运行着执行器守护进程,例如 YARN 的 NodeManager 或 Kubernetes 的 kubelet。执行器负责运行作业任务,发送心跳来表明自己仍然存活,并跟踪节点上的任务状态与资源分配。执行器收到启动任务的请求后,会取得作业的可执行代码,并运行命令来启动任务。然后,它会监视进程,直到进程结束或失败,再相应更新任务状态元数据。
许多执行器还会与操作系统配合,同时提供安全隔离与性能隔离。例如,YARN 和 Kubernetes 都使用 Linux cgroups。这样既能防止任务访问无权访问的数据,也能防止任务过度使用资源、影响同一节点上其他任务的性能。
- 资源管理器(Resource Manager)
-
编排器的资源管理器保存每个节点的元数据,包括可用硬件(CPU、GPU、内存和磁盘等)、任务状态、网络位置、节点状态以及其他相关信息。因此,资源管理器能够提供集群当前状态的全局视图。资源管理器的中心化特性可能同时形成可伸缩性和可用性瓶颈。YARN 使用 ZooKeeper、Kubernetes 使用 etcd 来存储集群状态(参见“协调服务”)。
- 调度器(Scheduler)
-
编排器通常有一个中心化的调度子系统,负责接收启动、停止作业或查询作业状态的请求。例如,调度器可能收到一项请求:使用特定的 Docker 镜像,在配有某种 GPU 的节点上启动一项包含 10 个任务的作业。调度器根据请求中的信息和资源管理器保存的状态,决定把哪些任务放在哪些节点上运行。随后,它会把分配结果通知任务执行器,由执行器开始运行任务。
虽然各个编排器使用的术语不尽相同,但几乎所有编排系统中都能找到这些组件。
Note
有些调度决策需要由应用专用的调度器来完成,以便考虑特定需求,例如在查询量达到某个阈值时自动扩展只读副本。中心调度器与应用专用调度器共同决定任务的最佳执行方式。YARN 把这种子调度器称为 ApplicationMaster,Kubernetes 则称为 operator。
资源分配
调度器在作业编排中扮演着格外棘手的角色:面对需求相互竞争的作业,它必须找出分配集群有限资源的最佳方式。从根本上说,调度决策需要在公平与效率之间取得平衡。
设想一个由五个节点组成的小型集群,总共有 160 个 CPU 核可用。集群调度器收到了两项作业请求,每项都希望使用 100 个核来完成工作。怎样调度才最好?
-
调度器可以同时为每项作业运行 80 个任务,等先前的任务完成后,再启动两项作业各自剩余的 20 个任务。
-
调度器也可以先运行一项作业的所有任务,等到有 100 个核可用时,再开始运行第二项作业。这种策略称为 成组调度(gang scheduling)。
-
一项作业请求会比另一项先到。调度器必须决定是把 100 个核全部分配给先到的作业,还是留下一些资源,为尚未到来的作业做准备。
这个例子非常简单,却已经暴露出许多艰难的权衡。以成组调度为例:如果调度器不断预留 CPU 核,直到 100 个核能够同时使用,那么一些节点就会闲置,集群的资源利用率也会下降;如果其他作业也试图预留 CPU 核,甚至还可能发生死锁。
另一方面,如果调度器只是等待 100 个核空闲,其他作业又可能在此期间抢走这些核。集群或许会在很长时间内都凑不出 100 个可用核,从而导致 饥饿(starvation)。调度器还可以 抢占(preempt)第一项作业的部分任务,终止它们来为第二项作业腾出空间。不过,抢占任务同样会降低集群效率,因为被终止的任务稍后需要重新启动、重新执行。
现在再设想一下,调度器必须为数百乃至数百万项这样的作业请求作出分配决策。要找到最优解似乎根本不可行。事实上,这个问题是 NP-hard 的;也就是说,除了规模最小的例子之外,计算最优解所需的时间长得令人无法接受 14 15 。
因此,实际的调度器会采用启发式方法,作出虽非最优、但还算合理的决策。常见算法包括先进先出(FIFO)、主导资源公平(DRF)、优先级队列、基于容量或配额的调度,以及各种装箱算法。这些算法的细节超出了本书范围,但调度确实是一个十分有趣的研究领域。
工作流调度
本章开头的 Unix 工具示例由若干命令串联而成。分布式批处理也经常采用相同的模式:一项作业的输出需要成为另一项或多项作业的输入,而一项作业又可能有多个输入,分别由其他作业产生。这种作业结构称为 工作流(workflow),也称为作业的 有向无环图(directed acyclic graph,DAG)。
Note
在“持久化执行与工作流”中,我们见过能够持久执行一系列步骤的工作流引擎,这些步骤通常会发出 RPC。在批处理的语境中,“工作流”有不同含义:它是一系列批处理过程,每个过程都接收输入数据、产生输出数据,通常不会向外部服务发出 RPC。持久化执行引擎通常会比批处理系统在每个请求中处理更少的数据,不过两者之间的界限并不十分清晰。
采用多项作业组成的工作流可能有几个原因:
-
如果一项作业的输出需要成为多项其他作业的输入,而且这些下游作业由不同团队维护,最好先让第一项作业把输出写到一个所有下游作业都能读取的位置。每当数据更新时,消费它的作业就可以安排运行,也可以按照其他时间表运行。
-
你可能需要把数据从一种处理工具传给另一种。例如,一项 Spark 作业把数据输出到 HDFS,随后由 Python 脚本触发 Trino SQL 查询(参见“云数据仓库”),继续处理 HDFS 文件,并把结果输出到 S3。
-
有些数据管道本身就需要多个处理阶段。例如,某个阶段需要按一个键分片数据,而下一阶段需要按另一个键分片,那么第一个阶段就可以按照第二阶段需要的方式对输出数据分片。
在 Unix 工具示例中,连接一项命令输出与另一项命令输入的管道只使用一个很小的内存缓冲区,并不会把数据写入文件。如果缓冲区已满,生产数据的进程就必须等待,直到消费进程从缓冲区中读走一部分数据,才能继续输出——这是一种 背压。Spark、Flink 等批处理执行引擎支持类似的模型,可以把一项任务的输出直接传给另一项任务;如果两项任务运行在不同机器上,数据就通过网络传输。
不过在工作流中,更常见的做法是让一项作业把输出写入分布式文件系统或对象存储,再由下一项作业从那里读取。这样可以使作业彼此解耦,在不同时间运行。如果一项作业有多个输入,工作流调度器通常要等到产生这些输入的所有作业都成功完成,才会运行消费这些输入的作业。
YARN ResourceManager 等编排框架中的调度器,以及 Spark 的内置调度器,都不会管理完整的工作流;它们只按单项作业进行调度。为了处理多次作业执行之间的依赖,人们开发了 Airflow、Dagster 和 Prefect 等工作流调度器。维护大量批处理作业时,这些调度器提供了非常有用的管理功能。许多数据管道的工作流通常包含 50 至 100 项作业;在大型组织中,还可能有许多团队运行不同的作业或工作流,跨越多个系统读取彼此的输出。管理这样复杂的数据流,离不开相应的工具支持。
故障处理
批处理作业往往会运行很长时间。一项包含许多并行任务、长时间运行的作业,很可能会在途中遇到至少一次任务失败。正如“硬件与软件故障”和“不可靠的网络”中所讨论的,这可能由许多原因造成,其中包括硬件故障(普通商用硬件上尤其常见)和网络中断。
任务无法完成的另一个原因,是调度器可能有意抢占(终止)它。当系统设置多个优先级时,抢占尤其有用:低优先级任务运行成本较低,高优先级任务则要付出更高成本。只要还有空闲计算容量,就可以运行低优先级任务;但如果一项高优先级任务到来,低优先级任务随时可能遭到抢占。这类价格较低的低优先级虚拟机,在 Amazon EC2 中称为 竞价实例(spot instance),在 Azure 中称为 竞价虚拟机(spot virtual machine),在 Google Cloud 中则称为 可抢占实例(preemptible instance)16 。
批处理通常用于时效性要求不高的作业,因此很适合使用低优先级任务和竞价实例,以降低运行成本。实质上,这些作业利用了原本会被闲置的计算资源,从而提高集群利用率。不过,这也意味着调度器更可能终止这些任务:抢占发生的频率要高于硬件故障 17 。
由于批处理作业每次运行都从头生成输出,任务失败要比在线系统中的故障容易处理:系统可以删除失败执行留下的不完整输出,再把任务安排到另一台机器上重新运行。不过,只因一项任务失败就重跑整个作业会非常浪费。因此,MapReduce 及其后继系统让并行任务彼此独立,从而可以按单项任务的粒度重试工作 3 。
如果一项任务的输出要在工作流中成为另一项任务的输入,容错就会更加棘手。MapReduce 的办法是始终把这类中间数据写回分布式文件系统,并等待写入任务顺利完成,之后才允许其他任务读取数据。即使在抢占频繁的环境中,这种做法也能正常工作,但它要向分布式文件系统写入大量数据,效率可能很低。
Spark 把中间数据保存在内存中,必要时“溢写”到本地磁盘,只把最终结果写入分布式文件系统。它还会跟踪中间数据的计算过程,以便在数据丢失时重新计算 18 。Flink 采用另一种方法,定期为任务状态的快照建立检查点 19 。我们将在“数据流引擎”中再次讨论这个话题。
批处理模型
我们已经了解了分布式环境如何调度批处理作业。现在把注意力转向批处理框架实际处理数据的方式。最常见的两种模型是 MapReduce 和数据流引擎。虽然在实践中,数据流引擎已经基本取代 MapReduce,但理解 MapReduce 的工作原理仍然很有用,因为许多现代批处理框架都深受它的影响。
MapReduce 和数据流引擎已经逐渐支持多种编程模型,其中包括底层编程 API、关系查询语言和数据框 API。丰富的选择使应用工程师、分析工程师、业务分析师,甚至不具备技术背景的员工,都能够出于各种用途处理企业数据。我们将在“批处理用例”中讨论这些用途。
MapReduce
MapReduce 的数据处理模式与“简单日志分析”中的 Web 服务器日志示例非常相似:
-
读取一组输入文件,并把它们拆分成 记录。在 Web 服务器日志示例中,每条记录就是日志的一行(也就是说,
\n是记录分隔符)。在 Hadoop MapReduce 中,输入文件存储在 HDFS 之类的分布式文件系统,或 S3 之类的对象存储中。系统可以使用多种文件格式,例如 Apache Parquet(列式格式,参见“列式存储”)或 Apache Avro(行式格式,参见“Avro”)。 -
对每条输入记录调用 mapper 函数,从中提取键和值。在 Unix 工具示例中,mapper 函数就是
awk '{print $7}':它提取 URL($7)作为键,并把值留空。 -
按键对所有键值对排序。在日志示例中,这一步由第一个
sort命令完成。 -
调用 reducer 函数,遍历排序后的键值对。如果同一个键出现多次,排序会让这些记录在列表中彼此相邻,因此不必在内存中保存大量状态,就能轻松合并这些值。在 Unix 工具示例中,reducer 由
uniq -c命令实现,它负责统计具有相同键的相邻记录数。
这四个步骤可以由一项 MapReduce 作业完成。第 2 步(map)与第 4 步(reduce)需要编写定制的数据处理代码。第 1 步(把文件拆成记录)由输入格式解析器负责。第 3 步的 sort 在 MapReduce 中是隐式的——mapper 的输出总会在传给 reducer 之前排序,所以不需要自行实现。排序是批处理中的基础算法,我们将在“混洗数据”中再次讨论。
要创建一项 MapReduce 作业,需要实现 mapper 和 reducer 两个回调函数,其行为如下:
- Mapper
-
每条输入记录都会调用一次 mapper,其任务是从输入记录中提取键和值。对于每条输入,它可以生成任意数量的键值对(也可以一个都不生成)。mapper 不会把某条输入记录的状态保留到下一条,因此每条记录都能独立处理。
- Reducer
-
MapReduce 框架接收 mapper 产生的键值对,收集属于同一个键的所有值,再用一个遍历该值集合的迭代器调用 reducer。reducer 可以生成输出记录,例如同一 URL 的出现次数。
在 Web 服务器日志示例中,第 5 步还有第二个 sort 命令,用于按请求次数排列 URL。在 MapReduce 中,如果需要第二个排序阶段,可以再编写一项 MapReduce 作业,把第一项作业的输出作为第二项的输入。从这个角度看,mapper 的作用是准备数据,把它转换成适合排序的形式;reducer 的作用则是处理已经排好序的数据。
[!TIP] MapReduce 与函数式编程 MapReduce 虽然用于批处理,其编程模型却来自函数式编程。Lisp 最早引入了 map 和 reduce(也称 fold),把它们作为列表上的高阶函数;后来,Python、Rust 和 Java 等主流语言也采用了这些函数。包括 SQL 提供的操作在内,许多常见的数据处理操作都可以建立在 MapReduce 之上。两个函数乃至函数式编程整体所具备的一些重要性质,恰好能为 MapReduce 所用。map 与 reduce 可以相互组合,这非常适合数据处理(正如 Unix 示例所示)。map 还天然易于并行,因为每项输入都独立处理;reduce 则可以并行处理不同的键。
实际上,使用原始 MapReduce API 实现复杂的处理作业非常困难而且费力——例如,作业使用的任何连接算法都必须从头实现 20 。与较新的批处理器相比,MapReduce 的速度也很慢。其中一个原因是,基于文件的 I/O 无法形成作业流水线:上游作业完成之前,下游作业不能开始处理它的输出。
数据流引擎
为了解决 MapReduce 的一些问题,人们开发了几种新的分布式批处理执行引擎,其中最著名的是 Spark 18 21 和 Flink 19 。它们的设计方式各有不同,却有一个共同点:把整个工作流作为一项作业处理,而不是拆成彼此独立的子作业。
这些系统显式建模了数据流经多个处理阶段的过程,因此称为 数据流引擎(dataflow engine)。与 MapReduce 一样,它们提供底层 API,通过反复调用用户定义的函数,每次处理一条记录;同时也提供 连接(join)和 分组(grouping)等高层算子。它们通过对输入分片来并行执行工作,并把一项任务的输出复制到另一项任务,作为后者的输入;如果两项任务运行在不同机器上,复制就通过网络完成。与 MapReduce 不同,算子不必严格交替扮演 map 和 reduce 的角色,而可以以更加灵活的方式组合。
这些数据流 API 通常使用关系模型风格的构件来表达计算:按照某个字段的值连接数据集,按键对元组分组,根据某项条件过滤数据,以及通过计数、求和等函数聚合元组。这些操作在内部通过下一节讨论的混洗算法实现。
这类处理引擎建立在 Dryad 22 和 Nephele 23 等研究系统的基础之上。与 MapReduce 模型相比,它们有几个优点:
-
排序之类代价高昂的工作只需在确有必要之处执行,不必默认出现在每个 map 阶段与 reduce 阶段之间。
-
如果一连串算子都不改变数据集的分片方式(例如 map 或 filter),它们就可以合并成一项任务,从而减少复制数据的开销。
-
工作流中的所有连接和数据依赖都经过显式声明,因此调度器可以总览全局,了解各处需要哪些数据,并据此优化数据局部性。例如,它可以尝试把消费某项数据的任务放在产生该数据的任务所在机器上,这样便能通过共享内存缓冲区交换数据,无需通过网络复制。
-
算子之间的中间状态通常只需保存在内存中,或写入本地磁盘;这样所需的 I/O 少于写入分布式文件系统或对象存储,因为后者还要把数据复制到多台机器上,并在每个副本所在机器上写入磁盘。MapReduce 已经对 mapper 的输出采用了这种优化,数据流引擎则把这个思想推广到了所有中间状态。
-
算子可以在输入准备就绪后立即开始执行,无需等前一阶段全部结束,下一阶段才能开始。
-
现有进程可以重复用于运行新的算子,从而减少启动开销;相比之下,MapReduce 会为每项任务启动一个新的 JVM。
数据流引擎能够实现与 MapReduce 工作流相同的计算,而且得益于上述优化,执行速度通常快得多。
混洗数据
我们已经看到,本章开头的 Unix 工具示例和 MapReduce 都以排序为基础。批处理器必须能够对 PB 级的数据集排序,而这些数据不可能装入一台机器。因此,它们需要一种输入和输出都经过分片的分布式排序算法,这种算法称为 混洗(shuffle)。
[!NOTE] 混洗不是随机 “shuffle” 这个词容易引起误解:把一副扑克牌洗牌,会得到随机顺序;这里所说的混洗却会产生排好序的结果,其中不含任何随机性。
混洗是批处理器的一项基础算法,连接与聚合都要使用它。MapReduce、Spark、Flink、Daft、Dataflow 和 BigQuery 24 都实现了可伸缩的高性能混洗算法,以处理大型数据集。下面以 Hadoop MapReduce 中的混洗为例 25 ,不过本节的概念也适用于其他系统。
{{< xref fig="11-1" page="/ch11" anchor="fig_batch_mapreduce" >}}图 11-1{{< /xref >}} 展示了一项 MapReduce 作业中的数据流。假设作业的输入已经分片,三个分片分别标为 m 1、m 2 和 m 3。例如,每个分片可以是 HDFS 中的单独文件,也可以是对象存储中的单独对象。同一数据集的全部分片,可以集中放在同一个 HDFS 目录中;在对象存储的存储桶里,它们也可以使用相同的键前缀。
{{< fig num="11-1" id="fig_batch_mapreduce" src="/fig/ddia_1101.png" caption="一个包含三个 mapper 和三个 reducer 的 MapReduce 作业。" class="ddia-figure ddia-figure--standard" width="2880" height="1920" />}}
框架会为每个输入分片启动一项单独的 map 任务。任务读取分配给它的文件,每次把一条记录传给 mapper 回调函数。计算的 reduce 端同样经过分片。map 任务的数量由输入分片数决定,而 reduce 任务的数量则由作业作者配置,两者可以不同。
mapper 的输出由键值对组成。框架必须保证:如果两个不同的 mapper 输出了相同的键,这些键值对最终会由同一个 reducer 任务处理。为此,每个 mapper 都会在本地磁盘上为每个 reducer 分别创建一个输出文件。例如,{{< xref fig="11-1" page="/ch11" anchor="fig_batch_mapreduce" >}}图 11-1{{< /xref >}}中的文件 m 1, r 2 由 mapper 1 创建,包含发往 reducer 2 的数据。mapper 输出一个键值对时,通常会对键进行哈希,据此决定把它写入哪个 reducer 文件,这与“按键的哈希分片”相似。
mapper 在写入这些文件的同时,还会在每个文件中按键排列键值对。这里可以采用“日志结构存储”中见过的技术:先在内存的有序数据结构中收集一批键值对,再把它们写成有序段文件,然后逐步将较小的段合并成较大的段。
每个 mapper 完成后,reducer 会连接到它,并把属于自己的有序键值对文件复制到本地磁盘。reduce 任务取得所有 mapper 输出中属于自己的那一份后,会像归并排序一样合并这些文件,同时保持排序顺序。这样一来,即使具有相同键的键值对来自不同的 mapper,最终也会彼此相邻。随后,每个键都会调用一次 reducer 函数,并向它传入一个迭代器,用于返回该键对应的所有值。
reducer 函数产生的记录会顺序写入文件,每项 reduce 任务对应一个文件。{{< xref fig="11-1" page="/ch11" anchor="fig_batch_mapreduce" >}}图 11-1{{< /xref >}}中的 r 1、r 2 和 r 3 就构成了作业输出数据集的三个分片,它们会被写回分布式文件系统或对象存储。
MapReduce 在 map 与 reduce 阶段之间执行混洗,而现代数据流引擎和云数据仓库则要精巧得多。BigQuery 等系统优化了混洗算法,尽量把数据保存在内存中,或把数据写入外部排序服务 24 。这类服务既能加快混洗,又能通过复制混洗数据来提高容错能力。
JOIN 与 GROUP BY
下面来看有序数据如何简化分布式连接与聚合。为便于说明,我们仍以 MapReduce 为例,不过这些概念适用于大多数批处理系统。
{{< xref fig="11-2" page="/ch11" anchor="fig_batch_join_example" >}}图 11-2{{< /xref >}} 展示了批处理作业中一个典型的连接示例。左侧是一份事件日志,记录已登录用户在网站上的操作,这些记录称为 活动事件(activity event),也叫 点击流数据(clickstream data);右侧则是用户数据库。这个例子可以看作星型模式的一部分(参见“星型与雪花型:分析模式”):事件日志是事实表,用户数据库则是其中一张维度表。
{{< fig num="11-2" id="fig_batch_join_example" src="/fig/ddia_1102.png" caption="用户活动日志与用户画像数据库的连接。" class="ddia-figure ddia-figure--wide" width="2880" height="1238" />}}
如果想结合用户数据库中的信息来分析活动事件——例如利用用户资料中的出生日期,了解某些页面更受年轻用户还是年长用户欢迎——就需要连接这两张表。假设两张表都大到必须分片,该如何计算这种连接?
可以利用 MapReduce 的一个性质:无论键值对最初位于哪个分片,混洗都会把具有相同键的键值对集中到同一个 reducer。这里可以把用户 ID 作为键。因此,我们可以编写一个 mapper 遍历用户活动事件,以用户 ID 为键,输出页面浏览的 URL,如{{< xref fig="11-3" page="/ch11" anchor="fig_batch_join_reduce" >}}图 11-3{{< /xref >}}所示。另一个 mapper 逐行遍历用户数据库,提取用户 ID 作为键、用户出生日期作为值。
{{< fig num="11-3" id="fig_batch_join_reduce" src="/fig/ddia_1103.png" caption="基于用户 ID 的排序合并连接。若输入数据集由多个文件分片组成,可并行启动多个 mapper 处理。" class="ddia-figure ddia-figure--wide" width="2880" height="1308" />}}
混洗随后确保 reducer 函数能够同时访问某位用户的出生日期,以及该用户的所有页面浏览事件。MapReduce 甚至可以安排记录的顺序,让 reducer 总是先看到用户数据库中的记录,接着再按时间戳顺序看到活动事件;这种技术称为 二次排序(secondary sort)25 。
这样,reducer 就能轻松完成实际的连接逻辑。第一个值应当是出生日期,reducer 先把它保存在局部变量中,再遍历具有相同用户 ID 的活动事件,输出每个浏览过的 URL 以及浏览者的出生日期。reducer 会一次处理某个用户 ID 的全部记录,所以任何时刻只需在内存中保存一条用户记录,也完全不必发出网络请求。这种算法称为 排序合并连接(sort-merge join),因为 mapper 的输出按键排序,而 reducer 随后会合并连接两侧的有序记录列表。
工作流中的下一项 MapReduce 作业可以继续计算每个 URL 的浏览者年龄分布。它首先用 URL 作为键混洗数据;排序完成后,reducer 会遍历同一个 URL 的所有页面浏览记录(其中包含浏览者的出生日期),为各个年龄段维护浏览次数计数器,并为每条页面浏览记录递增相应的计数器。这样就实现了 分组(group by)操作和聚合。
查询语言
多年来,分布式批处理的执行引擎日趋成熟。如今,基础设施已经足够稳健,可以在超过 10,000 台机器的集群上存储和处理数 PB 数据。既然如此规模的批处理系统如何实际运行已经大体得到解决,人们就把注意力转向了编程模型的改进。
MapReduce、数据流引擎和云数据仓库都采用 SQL 作为批处理的通用语言。这是很自然的选择:传统数据仓库本来就使用 SQL,数据分析和 ETL 工具也已经支持 SQL,而且开发者和分析师无不熟悉 SQL。
与手写 MapReduce 作业相比,查询语言接口除了所需代码更少,还有一个明显的优点:它们可以交互使用。你可以编写分析查询,再从终端或图形界面运行。这种交互查询方式很高效也很自然,业务分析师、产品经理、销售和财务团队等人员都可以借此在批处理环境中探索数据。尽管它不是经典形式的批处理,SQL 支持仍然使探索性查询成为分布式批处理系统的一项适用场景。
高级查询语言不仅提高了使用系统的人所能达到的生产率,也能从机器层面提高作业执行效率。正如“云数据仓库”所述,查询引擎负责把 SQL 查询转换成在集群中执行的批处理作业。从查询转换到语法树、再转换成物理算子的过程,让引擎有机会优化查询。Hive、Trino、Spark 和 Flink 等查询引擎都提供基于代价的查询优化器,可以分析连接输入的特征,自动决定哪一种算法最适合眼前的任务。优化器甚至可以改变连接的顺序,以尽量减少中间状态 19 26 27 28 。
SQL 是最流行的通用批处理查询语言,不过其他语言仍用于一些专门场景。Apache Pig 是一种以关系算子为基础的语言,允许用户逐步描述数据管道,而不是把所有逻辑写成一个庞大的 SQL 查询。数据框(见下一节)具有类似特征,Morel 则是受 Pig 影响的一种较新语言。还有一些用户采用 jq、JMESPath 或 JsonPath 等 JSON 查询语言。
在“图数据模型”中,我们讨论了如何用图来建模数据,以及如何用图查询语言遍历图中的边和顶点。许多图处理框架也支持通过查询语言进行批计算,例如 Apache TinkerPop 的 Gremlin。我们将在“批处理用例”中进一步讨论图处理的用例。
[!TIP] 批处理与云数据仓库正在收敛 过去,数据仓库运行在专用硬件设备上,为关系数据提供 SQL 分析查询。相比之下,MapReduce 等批处理框架的目标是提供更强的可伸缩性与灵活性:它们支持使用通用编程语言编写处理逻辑,并允许读写任意数据格式。
随着时间推移,两者变得越来越相似。现代批处理框架如今支持用 SQL 编写批处理作业,并通过 Parquet 等列式存储格式和经过优化的查询执行引擎,在关系查询上取得了良好性能(参见“查询执行:编译与向量化”)。与此同时,数据仓库迁移到云端后,变得更加可伸缩(参见“云数据仓库”),并实现了许多与分布式批处理框架相同的调度、容错和混洗技术。其中许多系统也使用分布式文件系统。
正如批处理系统采用 SQL 作为处理模型一样,云数据仓库也采用了数据框等其他处理模型(见下一节)。例如,Google Cloud BigQuery 提供 BigQuery DataFrames 库,Snowflake 的 Snowpark 则与 pandas 集成。Airflow、Prefect 和 Dagster 等批处理工作流编排器同样能够集成云数据仓库。
当然,并非所有批处理作业都容易用 SQL 表达。PageRank 等迭代图算法、复杂的机器学习以及许多其他任务,都很难用 SQL 编写。AI 数据处理包含图像、视频和音频等非关系型多模态数据,同样很难用 SQL 完成。
此外,云数据仓库不擅长某些工作负载。使用列式存储格式时,逐行计算的效率较低;在这种情况下,最好使用数据仓库的其他 API,或者改用批处理系统。云数据仓库也往往比其他批处理系统昂贵。大型作业改在 Spark 或 Flink 等批处理系统上运行,可能更加经济。
归根结底,应该用批处理系统还是数据仓库来处理数据,取决于成本、便利性、实现难度、可用性等因素。大多数大型企业都拥有多套数据处理系统,因而可以灵活作出选择;小型公司则往往只靠一套系统就已足够。
数据框
随着数据科学家和统计学家开始使用分布式批处理框架进行机器学习,他们发现现有的处理模型用起来很麻烦,因为他们习惯的是 R 和 pandas 中的数据框模型(参见“数据框、矩阵与数组”)。数据框与关系数据库中的表相似:它由许多行组成,同一列中的所有值都具有相同类型。用户不必编写一条庞大的 SQL 查询,而是调用与关系算子对应的函数,执行过滤、连接、排序、分组等操作。
最初,数据框操作一般在本机内存中执行,因而只能处理单台机器能够容纳的数据集。数据科学家希望继续使用熟悉的数据框 API,与批处理环境中的大型数据集交互。Spark、Flink 和 Daft 等分布式数据处理框架为满足这种需求,也采用了数据框 API。不过,本地数据框通常带有索引并且有序,分布式数据框一般却不具备这些性质 29 。因此,把程序迁移到批处理框架后,其性能可能出人意料。
数据框 API 看起来与数据流 API 相似,但具体实现各不相同。pandas 会在数据框方法被调用时立即执行操作;Apache Spark 则会先把所有数据框 API 调用转换成查询计划,经过查询优化,再在分布式数据流引擎上执行工作流。这样便能改善性能。
Daft 等框架甚至同时支持客户端与服务端计算:规模较小的内存操作在客户端执行,规模较大的数据集和计算则在服务端执行。Apache Arrow 等列式存储格式提供了统一的数据模型,可由客户端与服务端的执行引擎共享。
批处理用例
了解批处理如何工作之后,我们来看看它在各种应用中怎样发挥作用。批处理作业非常适合成批处理大型数据集,却不适合低延迟场景。因此,只要数据量很大、数据新鲜度又不重要,通常就能看到批处理作业。听起来这似乎是一个很大的局限,但事实证明,相当多的数据处理都符合这一模型:
-
会计与库存核对通常成批进行,企业借此检查交易是否与银行账户和库存相符 30 。
-
制造业的需求预测由周期运行的批处理作业计算 31 。
-
许多金融系统同样以批处理为基础。例如,美国的银行网络几乎完全依靠批处理作业运行 34 。
下面几节将讨论几乎每个行业都能见到的一些批处理用例。
提取—转换—加载(ETL)
“数据仓库”介绍了 ETL 和 ELT:数据处理管道从生产数据库抽取数据,对其进行转换,再把结果加载到下游系统。本节用“ETL”同时指代 ETL 与 ELT 工作负载。这类工作负载经常由批处理作业完成,尤其是在下游系统为数据仓库时。
批处理作业天然具有并行性,因而非常适合数据转换。许多数据转换工作负载都易于并行:过滤数据、投影字段,以及其他许多常见的数据仓库转换,都可以并行完成。
批处理环境还配有稳健的工作流调度器,能够轻松安排、编排和调试 ETL 数据管道作业。发生故障时,调度器通常会重试作业,以消除可能出现的暂时性问题。作业若反复失败,就会被标记为失败,开发者可以很容易地看到数据管道中的哪项作业停止了工作。Airflow 等调度器甚至内置了 MySQL、PostgreSQL、Snowflake、Spark 和 Flink 等数十种流行系统的数据源、数据接收端及查询算子。调度器与数据处理系统的紧密集成简化了数据集成。
我们还看到,批处理作业出了问题之后很容易排查和修复;这一点在调试数据管道时尤其有用。可以直接检查有问题的文件,找出错误所在;修复 ETL 批处理作业后,再重新运行即可。例如,输入文件可能不再包含转换作业所要使用的某个字段。数据工程师发现字段缺失后,可以更新转换逻辑,或者修改产生该输入的作业。
过去,数据管道通常由一支数据工程团队统一管理,因为要求开发产品功能的其他团队编写和管理复杂的批处理数据管道并不公平。近来,批处理模型和元数据管理的改进,使组织中的工程师更容易参与和管理自己的数据管道。数据网格(data mesh)35 36 、数据契约(data contract)37 和 数据编织(data fabric)38 等实践提供了标准与工具,帮助团队安全地发布数据,供组织中的任何人使用。
如今,数据管道与分析查询不仅开始共享处理模型,也开始共享执行引擎。许多批处理 ETL 作业与读取其输出的分析查询,会运行在同一套系统上。数据管道转换和分析查询都以 SparkSQL、Trino 或 DuckDB 查询来执行,已经十分常见。这样的架构进一步模糊了应用工程、数据工程、分析工程和业务分析之间的界限。
分析
在“分析型与事务型系统”中,我们看到分析查询(OLAP)经常扫描大量记录,并执行分组与聚合。这样的工作负载可以与其他批处理工作负载一起,在批处理系统中运行。分析师编写 SQL 查询,由查询引擎执行,并读写分布式文件系统或对象存储。表与文件之间的映射、名称和类型等表元数据,则通过 Apache Iceberg 等表格式和 Unity 等目录服务管理(参见“云数据仓库”)。这种架构称为 数据湖仓(data lakehouse)39 。
与 ETL 一样,SQL 查询接口的改进使许多组织如今也使用 Spark 等批处理框架进行分析。这类查询模式分为两种:
-
预聚合查询(pre-aggregated query):把数据汇总成 OLAP 多维数据集或数据集市,以加快查询(参见“物化视图与多维数据集”)。预聚合数据可以在数据仓库中查询,也可以推送到 Apache Druid 或 Apache Pinot 等专用的实时 OLAP 系统。预聚合通常按固定周期进行,这类工作负载由“工作流调度”中讨论的工作流调度器管理。
-
即席查询(ad hoc query):用户运行查询来回答特定业务问题、调查用户行为、调试运行故障,以及完成其他许多工作。在这种场景中,响应时间很重要。分析师会反复运行查询,在得到响应、进一步了解正在研究的数据后,继续调整查询。能够快速执行查询的批处理框架,可以减少分析师的等待时间。
SQL 支持还让批处理框架能够与电子表格和 Tableau、Power BI、Looker、Apache Superset 等数据可视化工具集成。例如,Tableau 提供 SparkSQL 和 Presto 连接器;Apache Superset 则支持 Trino、Hive、Spark SQL、Presto 等许多最终会执行批处理作业来查询数据的系统。
机器学习
机器学习(ML)经常用到批处理。数据科学家、机器学习工程师和 AI 工程师使用批处理框架来探索数据规律、转换数据并训练机器学习模型。常见用途包括:
- 特征工程:对原始数据进行过滤和转换,使其能够用于训练模型。预测模型通常要求输入数值数据,因此工程师必须把文本或离散值等其他形式的数据转换成所需格式。
- 模型训练:训练数据是批处理过程的输入,训练所得的模型权重则是输出。
- 批量推理:训练好的模型可以对大量数据成批进行预测,适用于数据集很大而又不要求实时返回结果的情况,其中也包括在测试数据集上评估模型的预测效果。
批处理框架为这些场景提供了专门的工具。例如,Apache Spark 的 MLlib 和 Apache Flink 的 FlinkML 都内置了丰富的特征工程工具、统计函数和分类器。
推荐引擎、排序系统等机器学习应用也大量使用图处理(参见“图数据模型”)。许多图算法可以表述为:每次沿一条边遍历,将一个顶点与相邻顶点连接起来以传播某些信息,并不断重复,直到满足某项条件——例如没有更多的边可以继续遍历,或者某项指标已经收敛。
批量同步并行(bulk synchronous parallel,BSP)计算模型 40 已经成为批量处理图数据时的常用模型。Apache Giraph 20 、Spark 的 GraphX API 和 Flink 的 Gelly API 41 等系统都实现了这一模型。它也称为 Pregel 模型,因为 Google 的 Pregel 论文推广了这种图处理方法 42 。
批处理也是 大语言模型(LLM)数据准备与训练的重要组成部分。网站等原始文本输入通常存放在分布式文件系统或对象存储中,必须先经过预处理才能用于训练。适合交给批处理框架完成的预处理步骤包括:
- 从 HTML 中提取纯文本,并修复格式损坏的文本;
- 检测并删除质量低劣、内容无关或重复的文档;
- 对文本进行分词(拆分成单词),再将其转换成嵌入向量,也就是每个单词的数值表示。
Kubeflow、Flyte 和 Ray 等批处理框架就是为这类工作负载设计的。例如,OpenAI 在 ChatGPT 的训练过程中就使用了 Ray 43 。这些框架内置了与 PyTorch、TensorFlow、XGBoost 等大语言模型和 AI 库的集成,并直接支持特征工程、模型训练、批量推理和微调(针对特定用例调整基础模型)等操作。
最后,数据科学家还经常在 Jupyter 或 Hex 等交互式笔记本中对数据进行实验。笔记本由若干 单元格(cell)组成,每个单元格都是一小段 Markdown、Python 或 SQL。按顺序执行这些单元格,可以生成电子表格、图表或数据。许多笔记本通过数据框 API 使用批处理系统,或者用 SQL 查询这类系统。
对外提供衍生数据
批处理作业经常用于构建预计算或衍生数据集,例如商品推荐、面向用户的报表,以及机器学习模型所需的特征。这些数据集通常由生产数据库、键值存储或搜索引擎对外提供。不论使用哪种系统,预计算数据最终都必须从批处理系统的分布式文件系统或对象存储,回到承载线上流量的数据库中。
最直观的选择或许是直接在批处理作业中调用数据库客户端库,逐条写入数据库服务器。这种方法确实能用——前提是防火墙规则允许批处理环境直接访问生产数据库——但基于以下几个原因,它并不是个好主意:
- 每条记录都发起一次网络请求,要比批处理任务的正常吞吐量慢几个数量级。即使客户端库支持批量写入,性能也很可能不理想。
- 批处理框架通常会并行运行许多任务。如果所有任务都以批处理应有的速率同时写入同一个输出数据库,数据库很容易不堪重负,查询性能也可能随之下降,继而给系统其他部分带来运行故障 44 。
- 批处理作业通常会对作业输出提供干净利落的“全有或全无”保证:如果作业成功,那么即使途中有些任务失败并经过重试,结果也等同于每个任务都恰好一次执行所产生的输出;如果整个作业失败,则不会产生任何输出。然而,从作业内部写入外部系统会产生无法用这种方式隐藏的外部可见副作用。这样一来,就不得不考虑尚未完整的作业结果被其他系统看到的问题;任务失败并重新启动时,也可能重复写入失败执行已经产生的输出。
更好的方案是让批处理作业把预计算数据集推送到 Kafka 主题之类的流中,我们将在第 12 章中进一步讨论这种做法。Elasticsearch 等搜索引擎、Apache Pinot 和 Apache Druid 等实时 OLAP 系统、Venice 等衍生数据存储 45 ,以及 ClickHouse 等云数据仓库,都内置了从 Kafka 摄取数据的能力。让数据经过流系统,可以缓解上述问题中的一部分:
- 流系统针对顺序写入进行了优化,因此更适合批处理作业的大批量写入负载;
- 流系统还可以在批处理作业与生产数据库之间充当缓冲区。下游系统可以限制自己的读取速率,确保仍有足够能力承载线上流量;
- 同一个批处理作业的输出可以由多个下游系统消费;
- 流系统可以充当批处理环境与生产网络之间的安全边界:它可以部署在所谓的 DMZ(隔离区)网络中,位于批处理网络与生产网络之间。
不过,让数据经过流并不会自动解决前面提到的“全有或全无”保证问题。为此,批处理作业必须在完成时通知下游系统:作业已经完成,数据现在可以对外提供。流消费者则必须能够在收到完成通知之前,让接收到的数据对查询保持不可见,就像采用 读已提交(read committed)隔离级别时未提交的事务一样(参见“读已提交”)。
另一种模式在数据库初始化时更为常见:直接在批处理作业 内部 构建一个全新的数据库,再把分布式文件系统、对象存储或本地文件系统中的文件批量加载到该数据库中。许多数据系统都提供了相应的批量导入工具,例如 TiDB 的 Lightning 工具,以及 Apache Pinot 和 Apache Druid 的 Hadoop 导入作业。RocksDB 也提供了从批处理作业批量导入 SST 的 API。
通过批处理构建数据库并批量导入数据,速度非常快,也更便于系统在不同数据集版本之间进行原子切换。另一方面,由批处理作业构建全新数据库时,要增量更新数据集会比较困难。如果同时需要初始化和增量加载,通常会采用混合方案。例如,Venice 支持混合存储,既可以进行基于行的批量更新,也可以切换整个数据集。
本章小结
本章探讨了批处理系统的设计与实现。我们先从经典的 Unix 工具链(awk、sort、uniq 等)入手,以此说明排序、计数等基本的批处理原语。
接着,我们把规模扩大到分布式批处理系统。我们看到,批处理式 I/O 以不可变、有界的输入数据集为处理对象并生成输出数据,因而可以在不产生副作用的情况下重跑和调试。为了处理文件,批处理框架主要由三部分组成:决定作业在何时何地运行的编排层,持久化数据的存储层,以及实际处理数据的计算层。
我们了解了分布式文件系统与对象存储如何通过按块复制、缓存和元数据服务来管理大文件,以及现代批处理框架如何通过可插拔 API 与这些系统交互。我们还讨论了编排器如何在大型集群中调度任务、分配资源和处理故障,并比较了调度单个作业的作业编排器,与管理一组依赖图作业整个生命周期的工作流编排器。
我们考察了几种批处理模型,首先是 MapReduce 及其经典的 map 和 reduce 函数,随后转向 Spark、Flink 等数据流引擎;它们的数据流 API 更易使用,性能也更好。为了理解批处理作业如何扩展,我们还介绍了混洗(shuffle)算法——这项基础操作使分组、连接和聚合成为可能。
随着批处理系统逐渐成熟,关注点也转向了易用性。SQL 等高级查询语言和数据框 API 降低了批处理作业的使用门槛,也让它们更容易优化。查询优化器会把声明式查询转换成高效的执行计划。
最后我们回顾了批处理常见用例:
- ETL 数据管道:通过定时工作流在不同系统之间提取、转换和加载数据;
- 分析:批处理作业既支持预聚合的仪表板,也支持即席查询;
- 机器学习:批处理作业负责准备和处理大规模训练数据集;
- 用批处理输出填充面向生产流量的系统:通常经由流或批量加载工具,把衍生数据提供给用户。
下一章我们将转向流处理,其中的输入是 无界的(unbounded):作业依然存在,但它的输入是永无止境的数据流。由于任何时刻都可能有更多工作到来,作业永远不会完成。我们会看到,流处理与批处理在某些方面相似,但“流是无界的”这一假设也会在很大程度上改变系统的构建方式。
脚注
参考文献
-
Nathan Marz. How to Beat the CAP Theorem. nathanmarz.com, October 2011. Archived at perma.cc/4BS9-R9A4 ↩︎
-
Molly Bartlett Dishman and Martin Fowler. Agile Architecture. At O'Reilly Software Architecture Conference, March 2015. ↩︎
-
Jeffrey Dean and Sanjay Ghemawat. MapReduce: Simplified Data Processing on Large Clusters. At 6th USENIX Symposium on Operating System Design and Implementation (OSDI), December 2004. ↩︎
-
Shivnath Babu and Herodotos Herodotou. Massively Parallel Databases and MapReduce Systems. Foundations and Trends in Databases, volume 5, issue 1, pages 1--104, November 2013. doi:10.1561/1900000036 ↩︎
-
David J. DeWitt and Michael Stonebraker. MapReduce: A Major Step Backwards. Originally published at databasecolumn.vertica.com, January 2008. Archived at perma.cc/U8PA-K48V ↩︎
-
Henry Robinson. The Elephant Was a Trojan Horse: On the Death of Map-Reduce at Google. the-paper-trail.org, June 2014. Archived at perma.cc/9FEM-X787 ↩︎
-
Urs Hölzle. R.I.P. MapReduce. After having served us well since 2003, today we removed the remaining internal codebase for good. twitter.com, September 2019. Archived at perma.cc/B34T-LLY7 ↩︎
-
Adam Drake. Command-Line Tools Can Be 235x Faster than Your Hadoop Cluster. aadrake.com, January 2014. Archived at perma.cc/87SP-ZMCY ↩︎
-
sort: Sort text files. GNU Coreutils 9.7 Documentation, Free Software Foundation, Inc., 2025. ↩︎ -
Michael Ovsiannikov, Silvius Rus, Damian Reeves, Paul Sutter, Sriram Rao, and Jim Kelly. The Quantcast File System. Proceedings of the VLDB Endowment, volume 6, issue 11, pages 1092--1101, August 2013. doi:10.14778/2536222.2536234 ↩︎
-
Andrew Wang, Zhe Zhang, Kai Zheng, Uma Maheswara G., and Vinayakumar B. Introduction to HDFS Erasure Coding in Apache Hadoop. blog.cloudera.com, September 2015. Archived at archive.org ↩︎
-
Andy Warfield. Building and operating a pretty big storage system called S3. allthingsdistributed.com, July 2023. Archived at perma.cc/7LPK-TP7V ↩︎
-
Vinod Kumar Vavilapalli, Arun C. Murthy, Chris Douglas, Sharad Agarwal, Mahadev Konar, Robert Evans, Thomas Graves, Jason Lowe, Hitesh Shah, Siddharth Seth, Bikas Saha, Carlo Curino, Owen O'Malley, Sanjay Radia, Benjamin Reed, and Eric Baldeschwieler. Apache Hadoop YARN: Yet Another Resource Negotiator. At 4th Annual Symposium on Cloud Computing (SoCC), October 2013. doi:10.1145/2523616.2523633 ↩︎
-
Richard M. Karp. Reducibility Among Combinatorial Problems. Complexity of Computer Computations. The IBM Research Symposia Series. Springer, 1972. doi:10.1007/978-1-4684-2001-2_9 ↩︎
-
J. D. Ullman. NP-Complete Scheduling Problems. Journal of Computer and System Sciences, volume 10, issue 3, June 1975. doi:10.1016/S0022-0000(75)80008-0 ↩︎
-
Gilad David Maayan. The complete guide to spot instances on AWS, Azure and GCP. datacenterdynamics.com, March 2021. Archived at archive.org ↩︎
-
Abhishek Verma, Luis Pedrosa, Madhukar Korupolu, David Oppenheimer, Eric Tune, and John Wilkes. Large-Scale Cluster Management at Google with Borg. At 10th European Conference on Computer Systems (EuroSys), April 2015. doi:10.1145/2741948.2741964 ↩︎
-
Matei Zaharia, Mosharaf Chowdhury, Tathagata Das, Ankur Dave, Justin Ma, Murphy McCauley, Michael J. Franklin, Scott Shenker, and Ion Stoica. Resilient Distributed Datasets: A Fault-Tolerant Abstraction for In-Memory Cluster Computing. At 9th USENIX Symposium on Networked Systems Design and Implementation (NSDI), April 2012. ↩︎
-
Paris Carbone, Stephan Ewen, Seif Haridi, Asterios Katsifodimos, Volker Markl, and Kostas Tzoumas. Apache Flink™: Stream and Batch Processing in a Single Engine. Bulletin of the IEEE Computer Society Technical Committee on Data Engineering, volume 38, issue 4, December 2015. Archived at perma.cc/G3N3-BKX5 ↩︎
-
Mark Grover, Ted Malaska, Jonathan Seidman, and Gwen Shapira. Hadoop Application Architectures. O'Reilly Media, 2015. ISBN: 978-1-491-90004-8 ↩︎
-
Jules S. Damji, Brooke Wenig, Tathagata Das, and Denny Lee. Learning Spark, 2nd Edition. O'Reilly Media, 2020. ISBN: 978-1492050049 ↩︎
-
Michael Isard, Mihai Budiu, Yuan Yu, Andrew Birrell, and Dennis Fetterly. Dryad: Distributed Data-Parallel Programs from Sequential Building Blocks. At 2nd European Conference on Computer Systems (EuroSys), March 2007. doi:10.1145/1272996.1273005 ↩︎
-
Daniel Warneke and Odej Kao. Nephele: Efficient Parallel Data Processing in the Cloud. At 2nd Workshop on Many-Task Computing on Grids and Supercomputers (MTAGS), November 2009. doi:10.1145/1646468.1646476 ↩︎
-
Hossein Ahmadi. In-memory query execution in Google BigQuery. cloud.google.com, August 2016. Archived at perma.cc/DGG2-FL9W ↩︎
-
Tom White. Hadoop: The Definitive Guide, 4th edition. O'Reilly Media, 2015. ISBN: 978-1-491-90163-2 ↩︎
-
Fabian Hüske. Peeking into Apache Flink's Engine Room. flink.apache.org, March 2015. Archived at perma.cc/44BW-ALJX ↩︎
-
Mostafa Mokhtar. Hive 0.14 Cost Based Optimizer (CBO) Technical Overview. hortonworks.com, March 2015. Archived on archive.org ↩︎
-
Michael Armbrust, Reynold S. Xin, Cheng Lian, Yin Huai, Davies Liu, Joseph K. Bradley, Xiangrui Meng, Tomer Kaftan, Michael J. Franklin, Ali Ghodsi, and Matei Zaharia. Spark SQL: Relational Data Processing in Spark. At ACM International Conference on Management of Data (SIGMOD), June 2015. doi:10.1145/2723372.2742797 ↩︎
-
Kaya Kupferschmidt. Spark vs Pandas, part 2 -- Spark. towardsdatascience.com, October 2020. Archived at perma.cc/5BRK-G4N5 ↩︎
-
Ammar Chalifah. Tracking payments at scale. bolt.eu.com, June 2025. Archived at perma.cc/Q4KX-8K3J ↩︎
-
Nafi Ahmet Turgut, Hamza Akyıldız, Hasan Burak Yel, Mehmet İkbal Özmen, Mutlu Polatcan, Pinar Baki, and Esra Kayabali. Demand forecasting at Getir built with Amazon Forecast. aws.amazon.com.com, May 2023. Archived at perma.cc/H3H6-GNL7 ↩︎
-
Jason (Siyu) Zhu. Enhancing homepage feed relevance by harnessing the power of large corpus sparse ID embeddings. linkedin.com, August 2023. Archived at archive.org ↩︎
-
Avery Ching, Sital Kedia, and Shuojie Wang. Apache Spark @Scale: A 60 TB+ production use case. engineering.fb.com, August 2016. Archived at perma.cc/F7R5-YFAV ↩︎
-
Edward Kim. How ACH works: A developer perspective --- Part 1. engineering.gusto.com, April 2014. Archived at perma.cc/F67P-VBLK ↩︎
-
Zhamak Dehghani. How to Move Beyond a Monolithic Data Lake to a Distributed Data Mesh. martinfowler.com, May 2019. Archived at perma.cc/LN2L-L4VC ↩︎
-
Chris Riccomini. What the Heck is a Data Mesh?! cnr.sh, June 2021. Archived at perma.cc/NEJ2-BAX3 ↩︎
-
Chad Sanderson, Mark Freeman, B. E. Schmidt. Data Contracts. O'Reilly Media, 2025. ISBN: 9781098157623 ↩︎
-
Daniel Abadi. Data Fabric vs. Data Mesh: What's the Difference? starburst.io, November 2021. Archived at perma.cc/RSK3-HXDK ↩︎
-
Michael Armbrust, Ali Ghodsi, Reynold Xin, and Matei Zaharia. Lakehouse: A New Generation of Open Platforms that Unify Data Warehousing and Advanced Analytics. At 11th Annual Conference on Innovative Data Systems Research (CIDR), January 2021. ↩︎
-
Leslie G. Valiant. A Bridging Model for Parallel Computation. Communications of the ACM, volume 33, issue 8, pages 103--111, August 1990. doi:10.1145/79173.79181 ↩︎
-
Stephan Ewen, Kostas Tzoumas, Moritz Kaufmann, and Volker Markl. Spinning Fast Iterative Data Flows. Proceedings of the VLDB Endowment, volume 5, issue 11, pages 1268-1279, July 2012. doi:10.14778/2350229.2350245 ↩︎
-
Grzegorz Malewicz, Matthew H. Austern, Aart J. C. Bik, James C. Dehnert, Ilan Horn, Naty Leiser, and Grzegorz Czajkowski. Pregel: A System for Large-Scale Graph Processing. At ACM International Conference on Management of Data (SIGMOD), June 2010. doi:10.1145/1807167.1807184 ↩︎
-
Richard MacManus. OpenAI Chats about Scaling LLMs at Anyscale's Ray Summit. thenewstack.io, September 2023. Archived at perma.cc/YJD6-KUXU ↩︎
-
Jay Kreps. Why Local State is a Fundamental Primitive in Stream Processing. oreilly.com, July 2014. Archived at perma.cc/P8HU-R5LA ↩︎
-
Félix GV. Open Sourcing Venice -- LinkedIn's Derived Data Platform. linkedin.com, September 2022. Archived at archive.org ↩︎
