第一章分了几类,介绍了一下数据系统架构的 Trade-Offs 。
概览
首先给出了 data-intensive 系统的定义,也就是对于数据的管理是这个系统首要的任务,包括数据的存储和读写,更具体一点说,包括确保数据的持久性以及读写的高可用。与之对比的是 compute-intensive ,也即专注对数据进行高性能地并行计算。后面各个章节的核心估计也是围绕着存储和读写的工程实现细节展开的,可以说是非常期待了。虽然 23 年还是 24 年的时候看了一遍,但是除了一些核心概念,其它的内容基本忘了 QAQ。
接下来列出了几个数据系统的核心概念:
database :
数据库,用于存储和管理数据,以供其它应用访问。也就是存数据的地方。
典型的像是 MySQL 、MongoDB 就不用提了。对于 MongoDB ,我目前应该算是有着相对丰富的运维经验。当然还有专门存时序数据的 Prometheus ,可以执行复杂查询的 Elasticsearch 。甚至是 Kafka 这种消息队列可能也算是一种数据库?给一个 partition 和 offset ,就可以把对应的数据 seek 出来,也可以根据时间戳 seek 后面的数据,因为它本地也会基于 offset 和时间戳维护索引。目前建库的相关倒排索引服务使用的 Kafka Topic ,被配置为 compact 类型,也即会存全量的数据,每次重启后,会读 topic 加载全量的数据。当前这种方式会严重影响 Kafka 的 IO 性能,因为从 earliest 开始读会污染 page cache ,同时 Kafka 也存在一些丢数据的问题,在某些异常情况下,止损时设置了 unclean leader 导致丢数据。后面针对这些问题,服务维护的团队也做了一些优化,例如会 dump 一份重启前的数据,这样重启后只需要从 dump 时刻最新的 offset 开始消费就可以了。最新版的服务使用的 Topic 甚至不是 compact 的,因为流式写 Kafka 的数据会额外再别的地方存一份,由其它的服务做 compact ,这样倒排索引服务加载数据时,对于 Kafka 部分只需要加载 compact 之后最新的数据即可,这个是根据数据的时间戳信息来 seek 的。这种设计方案其实跟经典数据库存储中的写 WAL 然后和存量数据 merge 是类似的。
cache :
用于保存执行成本较高的操作的执行结果,用于加速下一次的读。
根据之前的经验,除了加速读外,某种程度上还可以用于做降级,不过应该不算是严格意义上的 cache ,在读的目标不可用的时候,直接用本地的 cache 来兜底。22 年我入职没多久做的一件事情,就是给建库各个模块访问的远程词典加本地的 rocksdb cache 。建库的各个业务都有几个远程 kv 词典用于存储一些配置,其中涉及到一些比较重要的路由信息。由于远程 kv 字典不稳定,如果读失败就会导致路由信息缺失,模块不知道下游是啥,进而导致相关业务断流。加了 rocksdb cache 后,会将读到的 kv 字典内容写本地 rocksdb ,当远程字典读失败时,降低读本地的 rocksdb 。同时给 rocksdb 中的数据配置了数据过期时间,如果一个业务长期没有数据推送的话,那么 rocksdb 中业务相关的 cache 就会自然过期。上线后,后续在远程字典故障的情况下,做了非常稳定的兜底。
不过在使用 cache 时,需要注意的是如何保证 cache 中的数据和原数据是一致的。目前建库这边的配置中心服务,会在本地缓存 MySQL 中全量的配置信息,5 分钟同步一次。在配置上线时,总是会 sleep 5 分钟后才触发相关服务拉取最新的配置。在进行配置回滚止损的时候,也需要人工触发配置中心服务的 cache reload ,否则读到的配置信息是老的。
23 年接收北京大搜那儿做的纳海离线计算系统业务的时候,发现这个系统也做了算子级别的 cache ,也即将算子的输入字段和输出字段存到一个存储中,其中输入数据会算一个 hash 。每次计算的时候,在执行计算拓扑前,会根据建库 url 加载所有配置了需要使用的 cache 算子的输入和输出字段对应的值,如果新的计算中使用的输入字段的 hash 没变,那么直接会跳过这个算子,复用之前的输出字段。不过当算子中请求的远程服务有变更,需要重刷数据的话,我有点不太记得这种情况是咋处理的了,只记得在请求的时候有一个专门的开关,支持禁用 cache 。如果让我来设计的话,我应该会给算子再加一个输入字段,这个字段由计算系统来固定生成。在远程服务变更上线前改一下这个字段的值,例如 model_v1 –> model_v2 ,强制让 cache 失效。找大模型咨询了一下,一种更优雅的方式是将算子的版本号也作为 hash 的一部分,这样可以在算子本身升级后,自动触发 cache 失效。这种方式确实更优化,不需要用户人工再维护一个字段了。写到这里我突然想起来了,好像确实之前有说这种情况把相关算子重新发个版就可以了。
也就是说在设计 cache 是否过期时,需要将所有的因素都作为输入的数据,不论是计算相关的执行时间、输入数据、依赖服务的状态,还是其它的因素,例如更高层的拓扑状态,或者是更底层的硬件信息。反思了一下,一开始看到用输入字段来影响 cache 失效时,应该将 输入字段 作为 cache 失效的因素之一来思考,进而考虑是否还有其它因素,而不能局限于 输入字段 这个底层的概念中。这会就是所谓的抽象思维?
search-index :
允许用户根据关键词对数据进行搜索,以及通过多种方式对数据进行过滤。
说起这个让我想起了 Elasticsearch 一个神奇的特性。当前建库这边的 trace 数据,是通过 Elasticsearch 来构建查询索引的。最近有一个新的需求,是需要给一路高优业务做延时报表,往 Elasticsearch 中存了一个额外的字段,并且在构建报表数据的时候依赖基于这个字段对数据进行查询。当时想说肯定需要显示给这个字段设置索引,否则查询的时候走全表扫描就把服务干挂了。结果最后发现 Elasticsearch 会自动给相关字段添加索引,配置为
dynamic: true,不用人工配置索引。对比 MongoDB ,有时候会有业务私连集群查询不带索引导致把集群部分 shard 延时飙升,还会出现部分 shard 的表没有部分索引字段的问题。
stream processing : 当数据发生变更或者有事件时粒度进行处理
batch processing : 定期对一批量级较大的数据进行处理
书中强调,各项技术之间没有孰优孰劣,每一项技术都有优点和缺点,关键需要结合具体的场景进行分析,选取较为合适的技术。
技术对比
Operational System VS Analytical System
接下来是一些技术的对比。首先是 “Operational System VS Analytical System” 。简单来说,Operational System 重点强调对于数据的读写,也就是一般的后端系统 CRUD 做的事情,涉及到的数据量较少。对于 Analytical System 来说,重点为对于大批量的数据进行分析,一般会从前者中 Copy 一份镜像数据,然后将其处理为适用于数据分析的格式,最后执行相关的分类策略。
书中提到了 ”online transcation processing“,也即 OLTP ,和 ”online analytical processing“ ,也即 OLAP 。对于前者来说,应该就是 Operational System 做的事情,而后者是 Analytical System 来干的。一般来说一个系统只能属于其中的一类,比如说如果 scan 大量的数据,对于普通的数据库来说,可能会影响 cache 。
我想起了之前看的《Database Internals》提到,大部分数据库都对 cache 有过优化,并非简单的 LRU 。因为对于类似于 SELECT COUNT(*) FROM orders WHERE ... 这种仅用于分析的语句,扫描了几百万行,这几百万页数据全部被加载进 Buffer Pool,把之前所有热数据全部都挤出去了,但是等分析语句跑完,这个数据就再也不用了。一种方式是 MySQL 的实现,将 cache 进行了分类,分别为 Young 和 Old ,新的数据(页)先到 Old 区,之后被访问才会到 Young 区里。还有一种有意思的 FLU ,用到了一种叫 Count-Min Sketch 的概率数据结构,用一个固定大小的二维数组来近似估算任意 key 出现的频次,对于每一个 key ,每行 hash 一个位置,对这个位置 +1 ,在查询的时候,取所有行对应位置值的最小值即可,将哈希碰撞的影响降到最小。另外还引入的半衰期机制对相关的位置的值进行除 2 计算。
不过 LRU 并非一无是处,对于不涉及上面提到的极端的情况还是可以用的,这里又有一个很有意思的实现。一般 LRU 都是通过链表来实现,每次需要将 hit 到的节点移动到头部,并发的场景下操作链表会成为瓶颈。Clock 可以做到近似模拟 LRU ,但是成本会低非常多。它的数据结构是一个环形数组,每次 hit 的时候,将对应位置的 bit 置为 1 ,当需要驱逐的时候,扫描这个数组,将 1 的置 0,然后驱逐为 0 的。类似于是给为 1 的位置一次机会。扩展版本可以持续加到 3 ,也即 2 bit ,或者更多 bit 也行。为了对性能进行优化,前人想出的这些精妙的数据结构和算法还是挺有意思的。
扯远了,现在收回来。对于 Operational System ,主要的就是各种数据库,我算是相对比较熟悉的。而对于 Analytical System ,书中提到了 Data Warehousing 这个概念,也即数据仓库。前面提到了,Operational System 不太适合用做大批量的数据分析,并且很多的分析需求是涉及到多种类似的数据库的。那么此时数据仓库就派上用场了,根据书中的介绍,数据仓库汇聚了各个数据库中的一份全量的镜像数据,通过 ETL ,也即 extract-transform-load 这三步,将数据加工后,对外提供数据分析的功能。另外书中提到,部分数据库也提供了数据分析的功能,名字叫 hybrid transcational/analytical processing(HTAP)。
对于数据仓库来说,它对接的是各种的数据库,而对于 data lake (数据湖)来说,它包含的是各种类型的原始数据,比如说文本文件、图片、视频等,主要用于科学计算,例如说基于原始数据算向量或者 NLP 。我 21 年的时候在字节跳动的对象存储部门实习的时候,然后他们就准备搞数据湖。对象存储天然就可以存储各式各样的原始数据,反正存在磁盘上就是一个二进制的数据。
基于上述两种类型系统的对比,这里有提到了两种类型的数据,分别是 Systems of Record 和 Derived Data 。对于前者来说,是 source of truth ,也即是数据最原始的版本,所有的写操作都是针对这类数据的,而读操作读到的这类数据的值是最可信的,并且丢了就没了。对于后者来说,是前者的衍生物,类似于是一个 copy / 再处理,如果丢失了可以再基于前者生成,比如说 cache 。从两类数据的定义就可以看出,一般 Operational System 维护的是 Systems of Record ,而相对的 Analytical System 里使用的数据就是 Derived Data 。书中提到 ”By being clear about which data is derived from which other data, you can bring clarity to an otherwise confusing system architecture“ (通过理清数据的衍生关系,就可以减少对系统架构的困惑)。对于 Systems of Record 类型的数据来说,我们需要做的就是如何保证这个数据不丢,并且永远是准确的。对于 Derived Data ,需要做的就是维护和 System of Record 的一致性,以及根据业务需求实现一些额外的加工逻辑。
Cloud VS Self-Hosting
接下来是 ”Cloud VS Self-Hosting“ ,也即是使用云服务还是自己的服务。在这里分为了两个维度:”who builds the software and who deploys it“ ,也即服务是谁开发的以及服务是谁部署和维护的。这里分为了两个极端,一端是服务是自研的,并且是部署在自己的服务器上的,另一端是完全用 SaaS 服务。书中提到,其实大部分的场景下,会使用开源的软件,然后部署到自己的服务器上。一来使用开源的方案,可以基于具体的业务需求进行定制化改造,二来部署在自己的服务器上,数据安全可控。不过这种方式会引入一定的维护成本,也即需要有一个团队,或者至少一个人来对服务进行运维。
我所在组的大部分中间件服务以及一部分业务托管的服务都是采用这种方式进行部署的。例如 ZooKeeper、Kafka 和 MongoDB 是直接用开源的源码进行编译,然后对其进行了上云改造后,部署在了公司的 PaaS 平台上。其中对于 KRaft 版本的 Kafka ,我基于 3.9.1 的源码对其进行了一些改造,包括支持通过 Raft Log 感知 Voter 的变化和添加了 Voter 数量少于 n/2 +1 时强制选主的配置。对于 Grafana 和 Prometheus 则是直接部署在组内自运维的一批物理机上,按照模块的优先级进行了拆分,不同机器上的 Prometheus 采集不同模块的服务。这种方式确实比较省成本,根据之前同事的说法,自己裸布的 MongoDB 比集团云的 MongoDB 便宜了 90% ,不过由于是裸布的,少了一些必要的企业级能力,例如多租户和备份。对于 ZooKeeper 服务,22 年我刚入职的时候由于公司机房相关机柜退租,需要进行不停机迁移,当时这个事情是交给我做的,想了好几个天,搭了线下环境验证了多次,想了一个迁移的方案,也即每次人工先扩一部分实例,完成扩容后再缩,实时保证 voter 的数量为奇数。目前我主要负责维护上面提到的这些中间件和数据库,通过运维这些服务,结合实际的问题处理,确实可以学到不少东西,也手撸了不少运维工具辅助问题排查,并且可以将学到的理论知识进行实践,但代价是需要 24h 待命,随时处理问题。当然还有其他的一些服务会使用到公司的 SaaS 服务,例如 Redis、MySQL 和 K8S 。
Distributed VS Single-Node Systems
然后是 “Distributed VS Single-Node Systems” ,也即分布式服务和单机服务。尽管分布式计算带来了很多好多,例如高可用,但是也会带来一些问题,例如高可用的副本,需要有数据一致性机制和数据同步多分片无限扩展,带来的是多节点之间协同的网络成本(延时+出错)在一些情况下单机服务比多机服务的性能要好,最好看适用的场景。这让我想起了之前给大搜建库服务各个模块做的 Prometheus Exporter 。当时这部分监控的存储用的是集团云的 Prometheus 。因为大搜建库服务的实例数较多,使用部署在自运维的物理机上的 Prometheus 出现过几次爆内存的问题,为了服务的稳定性,后面就直接用了集团云的 Prometheus 。但是对于 Exporter ,由于物理机的性能足够,当时就没有部署至 PaaS 平台上,还用的是物理机,1 -2 个 Exporter 基本就可以把一台逻辑核为 40 的物理机的 CPU 打满。值得一提的是,集团云的 Prometheus 对单次采集有数据量的限制,针对这个问题,我也没搞分布式的 Exporter ,而是在一个 Exporter 内部开了多个端口,每个端口采集一部分实例的数据,以达到减少单次采集数据量的目的。关于实例分配的问题,我评估了一下,需要采集的实例数最多 10w ,都放到内存里的话,对于内存为 190 GB 的物理机来说,完全没有 oom 的风险,所以就直接存数组里进行切分了。
Data Systems, Law and Society
在本章的最后,讲的是 “Data Systems, Law and Society” ,其中提到了一个很有意思的点,我们存储数据是因为我们认为数据的价值比存储成本高,但是不能忽略额外的成本,例如数据泄漏的风险以及法律法规。一些法律要求个人数据需要支持删除,但是由于存储架构的限制,有些是不可改的 log ,有些被分发到 derived data 中了(cache、模型训练),没法保证可以删除。
总结
整体来说,本章仅为一个概览,在阅读的时候,看到熟悉的名词,进而联想到自己之前的一些工作经历和学到的知识还是挺有意思的。