领域逐渐熟悉,谈不上精通,但至少不会像前几章看的有点懵逼了,hh
引言
批量处理,顾名思义就是一次性处理一批数据的技术,对应的系统称之为离线系统(offline systems),目前我所在的搜索建库离线架构组,负责的一部分工作就是维护数据批量处理的通路。批量处理的执行流程,其实在日常中经常使用到,流程为 输入数据 -> 逻辑1 -> 逻辑2 .. -> 产出数据,其中产出数据不会覆盖输入数据,且各个逻辑一般均为幂等的。比如说需要分析一下一批日志中的请求 ip 的信息,列出请求数量最高的 ip ,一般就会用 Unix 管道将一批 CLI 命令串起来,样例为 cat *.log.* | awk '...' | sort | uniq -c | sort -n | head -10 。但是如果输入数据的量级较大,本地执行的耗时就长了,并且所需的计算资源单机也不能满足,此时就需要上分布式数据批量处理系统了,比如用 Hadoop 跑 MapReduce 任务。
分布式数据存储
之前的章节提到过,对于分布式系统来说,本质上就是协调分布在不同机器上的实例共同来共同执行任务。在本地通过 Unix 管道来处理数据时,上游进程写到管道中的数据是被缓存在系统内存中的,下游进程直接从内存中消费上游产出的数据。那么在切换至分布式系统时,上游任务产出的数据存哪儿呢?一种方式是存储在分布式文件系统中(Distributed Filesystem,DFS),例如使用 Hadoop 跑 MR 任务时使用的是 HDFS 。一些 DFS 实现了操作系统层面的 virtual filesystem(VFS),这样用户就可以通过 ls / cat 等命令,或者底层的 fopen / fseek ,来读取其中的文件。对于存储在磁盘上的文件,是按照 block 进行拆分存储在实际的存储设备上的,对于存储在 DFS 上的文件,也是被拆分为多个块,存储在不同的机器上,并且是多副本的,每一个块对应所在机器的信息,有统一的 metadata 系统来存储和管理,例如在 Hadoop 中叫 NameNode 。在我日常使用 HDFS 时,比较常用的还是通过 hadoop fs 命令,例如 -ls / -get / -mkdir , 等来管理 HDFS 上的文件。对于一些常用的目录,也会将其 mount 到物理机上,和普通的目录和文件一样通过 Linux 命令来查看和管理。
另一种实现多机器共享文件数据的存储为对象存储(Object Stores),这种模式有点类型于 kv 存储,只不过支持存储 value 的量级较大,比如可以存视频 / 二进制 bin 等。对象存储数据的 key 一般可以通过添加 / 来实现类似于目录层级的风格。虽然在底层数据存储的方式上和 DFS 类似,都是将数据切切块进行分布式存储,但是在使用时,对象存储的数据一般是通过覆盖全部的值来进行修改的,而不能像 DFS 文件一样只改一部分,另外就是没法完全模拟目录,比如 ls 时是类似于普通的递归 ls ,以及没法表示单独的目录。
就个人使用体感来说,在数据访问的场景,如果用 DFS ,需要用它指定的 client 来访问,比如 HDFS 需要使用 hadoop ,比较麻烦,但是如果是对象存储的话,数据可以直接通过 http 来访问,就是说直接通过 wget / curl 命令来下载即可。
任务执行与调度
除了产出数据的存储外,还需要一个系统来对计算进行调度,这个系统一般有这样几个组件:
Task Executor ,部署在任务执行的机器上,用于管理实际任务的启停以及收集当前机器的资源信息,例如 k8s 的 kubelet
Resource Manager ,用于存储和管理机器的资源信息、任务执行状态、节点状态等信息,例如 k8s 是使用 etcd 来存储的
Scheduler ,用于基于当前的系统资源信息来调度任务的执行,并且执行为 app 定制调度的逻辑,例如 k8s 的 operator
对于任务的调度,一般是 NP-hard 的,所以常常会使用一些启发式(heuristic)的策略。“启发式” 的含义为,虽然不是数据最优解,但是可以在尽量短的时间内给出一个足够好的解。对于任务调度场景来说,调度时用于分析的参数,例如资源信息、任务状态,都是实时变化的,如果调度算法执行的慢,那么就会有类似于 “刻舟求剑” 的问题,尤其是如果需要调度的任务和执行任务的机器的量级非常大的情况下。一种比较有意思的解法是 Bin-packing(装箱算法),目标是将任务经量紧地装到机器中,其中 Best-fit 为将任务调度到剩余资源最接近的机器上,这样可以减少资源碎片;First-fit 为将任务调度到第一台满足其执行资源需求的机器上,这样调度速度快。另外还可以基于 Priority Queue 来做优先级调度,打分的方式有多种,例如等待时间越长优先级越高。
在任务执行时,一般都是按照一个特定的 workflow 来执行,workflow 中各个节点产出的数据需要被下游节点接收,对于 MapReduce 来说,会将中间节点的处理结果写 DFS ,而对于 Spark 来说,会将其存在内存/文件中。我虽然对 Spark 没啥了解,但是看这里的描述,感觉像是假如一个 workflow 中的上下游任务是在一台机器上的话,那么就通过内存/本地文件来传递处理的数据,假如是在不同机器上的话,还是需要使用 DFS 的。
目前搜索离线建库这边,在跑批量建库数据产出任务时,会使用一个叫做 compass 的系统,来执行由多个 MR 任务及脚本任务组成的 DAG 。这样做是因为数据处理的逻辑比较复杂,分了很多阶段,单独一个 MR 任务搞不定。不过正是由于拆了很多个阶段,就导致这些阶段产出的中间数据都需要存 HDFS ,需要使用大量的存储资源。25 年开始已经陆续将一部分库种的倒排索引构建部分的任务迁到了一个新的系统中,对于每一个分片的数据,数据处理 DAG 里的所有任务都会在一台机器上来执行,这样 DAG 中各个阶段产出的中间数据,如果没有特殊需求,都可以直接存在本地,这样可以极大节省 DFS 的资源。由于我没有具体负责这块系统,只是在对接建库业务时使用过,所以不太清楚具体的实现细节,只了解一个大概。
MapReduce
MapReduce 是一个比较经典的批量数据处理模式,一次 MapReduce 主要的流程如下:
首先读取输入数据,然后特定的分割规则(一般是 \n)将数据进行拆分,形成一批分片
然后调用 mapper 函数,对各个分片进行处理,生成一批 key - value 键值对
接着根据特定的策略对 key 进行排序,并再次进行拆分,形成一批分片,一般相同的 key 在同一分片上(shuffle)
最后调用 reducer 函数,遍历分片上的 key-value 进行处理,产出最终的结果
一个比较有意思的使用场景为使用 MapReduce 比较两份量级较大的 url list 的 diff 。首先我们假定 url.list.A 在 HDFS 上的路径为 /usr/build/iteration-A/* ,url.list.B 在 HDFS 上的路径为 /usr/build/iteration-B/* ,两者都是一行一个 url 。
接下来在 mapper 阶段,产出 key -value 的策略为,如果输入分片来自于路径 /usr/build/iteration-A/* ,则产出的数据为 url \t A ,否则为 url \t B 。这个路径的信息可以从环境变量 map_input_file 中获取。在对产出的数据进行排序时,我们要求排序的规则为,首先根据第一列进行排序,如果第一列的值相同,则根据第二列排:
-D stream.num.map.output.key.fields=2指定产出数据的前几列作为 key ,这里是前两列-D num.key.fields.for.partition=1指定基于前几列 key 来进行分片,这里为了让同一个 url 落在一个分片上,所以指定仅基于第一列 url 。-D mapred.text.key.comparator.options="-k1,2"指定基于哪几列 key 来排序,这里由于需要确保 url 相同时,A 在 B 的前面,所以需要将这两列都作为排序的参考。设置使用的排序算法为
-D mapred.output.key.comparator.class=org.apache.hadoop.mapred.lib.KeyFieldBasedComparator,指定分片策略为-partitioner org.apache.hadoop.mapred.lib.KeyFieldBasedPartitioner
此时假如原始 url 的样例数据为:
1 | url.list.A |
则预期 reducer 输入的数据样例为:
1 | url1 A |
此时 reducer 的逻辑实现的就比较简单了,按行处理数据,如果当前的 url 和前一个 url 不一样(注意处理边界情况),则根据前一个 url 是否的第二列是否出现过 A 或 B 来判断是只在 A 或 B 中存在,还是都存在,为不同的情况打上不同的标签,比如说都存在用标签 C ,否则为 A 或 B,例如:
1 | url1 #C |
最后为了便于查看,可以设置 -outputformat org.apache.hadoop.mapred.lib.SuffixMultipleTextOutputFormat ,这样会根据 # 后的值,将第一列放到不同的文件中,这样就可以按需下载了,文件名称样例为:
1 | part-xxx-A |
其实在日常使用中,用的特别多的场景主要是用 MR 来跑并行计算任务,在这种场景下,一般只有 Mapper ,没有 Reducer 。MR 任务支持将 shell 脚本,而非 jar 包,作为 mapper 来执行,通过 hadoop streaming 提交任务即可。其中一个常见的业务场景为,将建库任务产出到 HDFS 上的 url.list 作为 MR 的输入,在 mapper 中执行 rpc client 请求流式建库服务来重刷 url 数据。使用 MR 跑并发任务的一个好处是,不用自己去管理任务的执行状态,并且可以方便地通过 -set-map-capacity 参数来调整任务执行的并行度。我 22 年入职后第一次接触 MR 的任务,就是这种只有 mapper 的任务,用来做批量数据迁移。另外假如并行执行的任务需要处理的数据量较大且为 CPU 密集型的话,本地单机跑的速度肯定没有分布式跑的快。
其它
书中还提到了 Dataflow Engines ,相较于 MapReduce ,业务在实现节点的业务逻辑时,不需要严格实现 mapper 和 reducer ,更灵活。并且和 MR 任务不同的是,数据处理的中间结果需要落 DFS ,而是一般可以通过内存/本地文件进行传递。对于本地数据的批量科学计算,可以使用 DataFrame 结构,它是一系列 row 的集合,对于这个集合的每一列,值的类型是相同的。
数据批量处理的应用场景有很多,比如数据仓库中的 ETL(Extract-Transform-Load)、数据分析和机器学习。之前看《Database Internals》 的时候,了解到其实新索引的构建也算是 Batch Processing ,因为它是基于一个固定的全量数据来生成新的数据(索引),由于一开始就知道全集,所以在构建 B+树的时候,先对数据进行排序(一般是外部排序),然后就可以直接通过 Bottom-Up ,自底向上构建树结构,避免频繁地页分裂。