基本概念

逻辑复制是数据同步的一种方式,逻辑二字主要指的是传输数据的格式,与之相对的是物理复制。

逻辑数据格式

逻辑数据格式是指,能够由一定规则描述的,符合SQL直觉的一种数据格式,例如:

[MsgType:Insert,Relation:test,New Tuple:[a, int, 100]]

从这种数据格式上,能够通过解析等方式,将数据还原为SQL语句,这也同时表现了逻辑复制的第一个特点:无视物理结构上的差异,可以在异构、跨版本的数据库之间进行数据同步。

接上例,通过一定的解析步骤,我们总能够得到:

Insert into test(a) values 100::int;

对于不同内核来说,可能在语法上稍有差异,但只需要修改一些解析方式,也能够得到自身看的懂、可以运行的命令。

与物理复制的区别

逻辑复制的另外一个特点,是相对于另一种主流的数据同步方式:物理复制,来作为比较的。

几乎所有的数据库内核,在进行数据存储的时候,都是采用随机写的方式,最大化利用性能优势,将数据存储在磁盘上。那么从物理上看,一个relation的所有数据,可能并不是连续的:他们可能分布在不同的page上,由指针进行连接。那么从磁盘角度看,一个page的分布,可能是这样的:

--------------
Page Header
--------------
Tuple 1 for r1
--------------
Tuple 2 for r1
--------------
Tuple 1 for r2
--------------
Tuple 1 for r3
--------------
Tuple 2 for r2
--------------

如此一来,想要在物理层面进行数据传输时,以page为单位,是无法将不同relation的tuple分开的,也就是细粒度不够。

逻辑复制的另一个特点,是可以以表为单位,进行数据同步,由一个Reader完成。

Reader逐行的去读取物理中的tuple,筛选那些符合要求的tuple,进行下一步工作。

逻辑复制提高的数据同步的细粒度,可以在库、表、甚至行上进行数据过滤,显然是一种更灵活的同步方式,在某些场景下,其效果更好,常见的有:

应用场景

  1. 某数据库使用厂商需要进行版本升级,但升级前后版本由于迭代问题,在物理存储上存在差异,导致无法进行直接替换;
  2. 某数据库使用厂商需要将原有的数据库数据迁移到一款新的数据库产品上,而这两个数据库的内核架构完全不同,物理存储存在本质差异;
  3. 某厂商部门之间使用的数据库产品不同,但部门之间存在协作,需要进行数据传递; …

架构

alt text

如图所示,逻辑复制功能,是两个实例(或生产者与消费者)之间完成的数据同步,通常我们把作为数据源端的实例,称为发布端(Publisher),而准备接收数据并进行回放的实例称为订阅端(Subscriber)。

发布端

发布端的组成部分,共分为以下几个部分:WALReader、Decoder、Sender,通常是由三个独立的线程,按照其角色完成不同的工作:

  • WALReader:主要完成对Xlog日志的读取,即从磁盘上将当前已罗盘的Xlog record读取到内存,并发送给Decoder;
  • Decoder:解码线程,由于WALReader发送来的数据,是物理层面的记录,Decoder会按照特定的规则,将物理信息转换为逻辑信息,从内存上看,使用一个结构体,保存Xlog record中的xid、LSN、DMLType等信息。同时,Decoder还需要将这些信息,按照约定好的协议(Protocol)将数据进行打包,其中应用最广泛的端到端协议是pgoutput(后面会介绍他的协议格式以及细节实现)。Decoder会将这些按照协议进行拼接的逻辑数据存储到私有内存中,当达到一定条件时,将数据发送给Sender线程执行发送。
  • Sender:将Decoder发送来的数据通过网络发送给订阅端。

发布端如何保证事务一致性

对于数据同步来说,最重要的是事务一致性,如果不能保证一致性,那么逻辑复制将是毫无意义的。而要保证事务一致性,通常主要责任在于发布端,发布端需要确保自己发送逻辑信息的顺序,是严格的,即使不能保证发送顺序一定与WAL日志写入顺序完全一致,但必须保证事务的最终一致性。为了保证事务一致性,在发布端,通常会有以下几个部分保证:

  • 严格按照commit顺序进行事务发送:事务的提交顺序非常重要,Decoder解码到一条逻辑数据时,并不会立刻发送给Sender,而是将条变更存储起来,直到解码到该事务对应的commit,才将整个事务完发送;
  • 错误重发:当订阅端接收到一个事务的commit时,即代表这个事务经传输完毕,订阅端会发送一个确认LSN给发布端,发布端需要确认LSN是在当前处理LSN的前面,如果此时LSN发生错乱,发布端会立即这条返回的LSN开始,重新发送这些变更;
  • commit时间确认:在Decoder将Xlog信息解码为逻辑协议信息时对于某些类型的消息,会携带commit timestampz,这是该事务提时的瞬时时间戳;

发布端对于内存的管理

在业务系统中,有时不可避免的会出现大事务,通常包含几M甚至百M的事务变更数据,在逻辑复制模块中,上文提到,通常是将所有事务存储于内存,当遇到该事务的commit时,才进行发送。对于这种大事务来说,如果多几条,可能很容易造成内存压力,引起CPU抖动,有影响系统整体性能的可能性。

为了保证系统整体性能,当遇到这种大事务的时候,通常会设计一个内存阈值,当在内存中已经存储了多条事务变更时,一旦当前内存占用达到阈值,发布端会选择当前内存中一个最大的事务(不一定是完整的,只是transaction size最大)spill to disk;诚然,这会带来一定的IO开销,因为当我们读到被spill的事务的commit时,需要把它从磁盘上读出来,但相比于其带来的内存负荷,这是合理的取舍。

发布端的性能瓶颈

通过这个架构,我们不难看出一个性能瓶颈,那就是整个流程,都是串行完成的:

  • WALReader在串行的读取Xlog;
  • Decoder在串行的对Xlog record进行解码;
  • Sender在串行的发送解码好的逻辑信息;

根据测试,在每秒大约100MB/s的Xlog写入负载下,单线程解码的速度只有~9MB/s,这对于一些使用逻辑复制的场景来说,是不可接受的,因为差异如此大的性能,如果在数据源端负载不停的情况下,数据同步端将永远没有追赶上的可能性,Lag会累积越来越大,最终将导致不得不停止业务负载,等待同步,这对于一些实时性要求高的业务来说是完全不可接受的。

所以单线程解码的瓶颈,立即成为制约整个逻辑复制模块的关键点。

发布端并行解码

由于单线程解码的性能无法满足使用要求,所以多线程并行解码,就成为一个解决问题的途径。

并行解码的本质,是依靠Reader本身读取的快速性能(仅有Read操作,性能几乎完全取决于物理磁盘性能),快速的读取Xlog record,并将其传递给并行的解码线程。

但是这里有一个关键的问题:我们应该如何进行分发?线程之间相互独立,彼此不可能知道对方的状态(除非使用线程间通信,然而这样并行的优势会被大大缩小),如果将事务变更随机发送,由于线程的调度问题,有的线程快,有的线程慢,那么保证事务一致性是一个非常困难的问题。或许增加一个收集(Collector)线程,来将所有的解码信息进行整合,并按照事务顺序进行重排序是一个方案,然而重排序的代价非常大,会十分影响高并发的性能。

对于这个问题,OG使用的方式是对HASH XID,将所有相同XID事务变更都交给同一个线程处理,由此来保证事务内的相对顺序是严格的。而事务间的相对顺序,则交由单线程的Sender保证:当Deocder完成本事务的处理,发送给Sender后,Sender内部会按照commit的顺序,对事务进行排序,严格按照提交顺序发送,确保事务一致性。

订阅端

在数据同步端,订阅端,基本是由两个部分组成:recieve与apply。

  • receive:负责接收发布端传输的逻辑信息(或心跳信息),发送相应的确认LSN与心跳reply,通常发布端发送的逻辑信息是一个batch,其中包含了一个事务的所有变更操作,receive负责将这写batch拆解为一个个独立的变更操作,并将这些独立的变更发送给apply。
  • apply:对receive发来的每条独立的变更进行回放,通常是Insert、Update、Delete操作,apply负责解析这些逻辑数据中的关键部分,例如:Insert信息中通常包含relation信息,new tuple信息等等,apply将这些信息重组为对应的内存变量,发送给heap_insert。

订阅端的性能瓶颈

从功能实现上看,事务一致性大多数是由发布端保证的,只要发布端