PostgreSQL学徒

从实际案例聊聊逻辑解码

1前言

昨晚一位同事打电话找到我说:"有个库的逻辑复制貌似卡住不动了,调整了内存参数 logical_decoding_work_mem 也没有实质性的改善",于是我上去看了下,发现对应的复制槽所保留的 WAL 还是1月1号上午8点的,相当于堆积了两天的日志,但是复制槽的状态却又是 active 的,给人一个还在正常消费的错觉。

让我们一起分析下这个怪异的案例!

2分析

既然是从前天上午就开始保留WAL,遂翻一下前天的日志

Image

这一看就发现了端倪,日志中从凌晨开始就不断地在打印如下日内容

"oldest xmin is far in the past",,"Close open transactions soon to avoid wraparound problems. You might also need to commit or roll back old prepared transactions, or drop stale replication slots"

很好理解,这条告警表示当前数据库中存在一个十分老的事务ID,阻碍年龄的回收,提示你要赶紧检查一下预备事务或者复制槽了

 /*
  * If oldestXmin is very far back (in practice, more than
  * autovacuum_freeze_max_age / 2 XIDs old), complain and force a minimum
  * freeze age of zero.
  */

 safeLimit = ReadNextTransactionId() - autovacuum_freeze_max_age;
 if (!TransactionIdIsNormal(safeLimit))
  safeLimit = FirstNormalTransactionId;

 if (TransactionIdPrecedes(limit, safeLimit))
 {
  ereport(WARNING,
    (errmsg("oldest xmin is far in the past"),
     errhint("Close open transactions soon to avoid wraparound problems.\n"
       "You might also need to commit or roll back old prepared transactions, or drop stale replication slots.")));
  limit = *oldestXmin;
 }

由于这个库上有多个复制槽的存在,而且又没有配置max_slot_wal_keep_size,因此WAL不断堆积,磁盘很快就告警了,我去看了一下历史操作,发现值班DA亦或是主管DA做了不规范的操作,手动操作了WAL👇

Image

1月1号17:25的时候DA移除了某个 WAL ,直到10分钟之后才改了回来。真是个危险的操作。数据库自身是有一套完整的WAL清理机制的,因此手动操作 WAL 是很危险的事情,这也是为什么自10版本以后要把 pg_xlog 目录重命名为 pg_wal 的原因,就是怕产生误导,以为是"日志"可以清理。

另外在我们自研数据库中新增了一个特性,在原生 PostgreSQL 中,假如 WAL 缺失了,流复制会断开,提示 WAL segment has been removed,遇到此类错误需要手动在备库上配置从归档中拷贝缺失的 WAL ,假如归档没有那么很不幸只能重建,对于大库尤为头疼。因此我们加了一个特性,自动从归档中找需要的 WAL,可惜看历史操作,还做过 pg_archivecleanup 的操作,清除了部分 WAL ,因此日志中就在开始恶性循环了:启动新的逻辑订阅进程想要继续 → 发现需要的WAL被移除了,然后去归档目录找 → 发现归档中也被清理 → 于是又启动新的逻辑订阅进程,不断往复

Image

3复现

经过前面的简单分析,WAL被移除导致walsender缺失了解码所需的信息。

那让我们复现一下,新建一个发布订阅

postgres=# create table t1(id int primary key);
CREATE TABLE 
postgres=# create publication pub1 for table t1;
CREATE PUBLICATION

订阅端

postgres=# create table t1(id int primary key);
CREATE TABLE
postgres=# create subscription sub1 connection 'host=localhost port=5432' publication pub1;
NOTICE:  created replication slot "sub1" on publisher
CREATE SUBSCRIPTION

此时WAL从00000001000000000000009A开始保留

postgres=# select slot_name,active,active_pid,pg_walfile_name(restart_lsn) as wal_recv from pg_replication_slots ;
 slot_name | active | active_pid |         wal_recv         
-----------+--------+------------+--------------------------
 sub1      | t      |      22967 | 00000001000000000000009A
(1 row)

postgres=# select pg_walfile_name(pg_current_wal_lsn());
     pg_walfile_name      
--------------------------
 00000001000000000000009A
(1 row)

[postgres@xiongcc pg_wal]$ ls -lrth
total 17M
drwx------ 2 postgres postgres 4.0K Dec 25 19:18 archive_status
-rw------- 1 postgres postgres  16M Jan  3 11:14 00000001000000000000009A  ---👈🏻保留的WAL从这里开始

现在发布端开启一个事务但不提交(此时的订阅端当然是看不到这条数据的)发布端还是从00000001000000000000009A这个日志开始保留

postgres=# begin;
BEGIN
postgres=*# insert into t1 values(1);
INSERT 0 1
postgres=*# select slot_name,active,active_pid,pg_walfile_name(restart_lsn) as wal_recv from pg_replication_slots ;
 slot_name | active | active_pid |         wal_recv         
-----------+--------+------------+--------------------------
 sub1      | t      |      22967 | 00000001000000000000009A
(1 row)

接着用pgbench模拟一些写入以产生一些 WAL ,不一会儿就产生了许多个日志

[postgres@xiongcc pg_wal]$ ls -lrth
total 641M
drwx------ 2 postgres postgres 4.0K Dec 25 19:18 archive_status
-rw------- 1 postgres postgres  16M Jan  3 11:15 00000001000000000000009A  ---👈🏻最开始保留的WAL从这里开始
-rw------- 1 postgres postgres  16M Jan  3 11:15 00000001000000000000009B
-rw------- 1 postgres postgres  16M Jan  3 11:15 00000001000000000000009C
-rw------- 1 postgres postgres  16M Jan  3 11:15 00000001000000000000009D
-rw------- 1 postgres postgres  16M Jan  3 11:15 00000001000000000000009E
-rw------- 1 postgres postgres  16M Jan  3 11:15 00000001000000000000009F
-rw------- 1 postgres postgres  16M Jan  3 11:15 0000000100000000000000A0
-rw------- 1 postgres postgres  16M Jan  3 11:15 0000000100000000000000A1
-rw------- 1 postgres postgres  16M Jan  3 11:15 0000000100000000000000A2
-rw------- 1 postgres postgres  16M Jan  3 11:15 0000000100000000000000A3
-rw------- 1 postgres postgres  16M Jan  3 11:15 0000000100000000000000A4
-rw------- 1 postgres postgres  16M Jan  3 11:15 0000000100000000000000A5
-rw------- 1 postgres postgres  16M Jan  3 11:15 0000000100000000000000A6
-rw------- 1 postgres postgres  16M Jan  3 11:15 0000000100000000000000A7
-rw------- 1 postgres postgres  16M Jan  3 11:15 0000000100000000000000A8
-rw------- 1 postgres postgres  16M Jan  3 11:15 0000000100000000000000A9
-rw------- 1 postgres postgres  16M Jan  3 11:15 0000000100000000000000AA
-rw------- 1 postgres postgres  16M Jan  3 11:15 0000000100000000000000AB
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000AC
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000AD
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000AE
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000AF
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000B0
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000B1
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000B2
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000B3
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000B4
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000B5
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000B6
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000B7
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000B8
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000B9
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000BA
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000BB
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000BC
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000BD
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000BE
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000BF
-rw------- 1 postgres postgres  16M Jan  3 11:16 0000000100000000000000C0
-rw------- 1 postgres postgres  16M Jan  3 11:17 0000000100000000000000C1  ---👈🏻现在新插入的数据位于这个WAL中

至此,原先开启的事务现在已经横跨了很多个WAL,第一条数据位于00000001000000000000009A中,后面两条数据位于0000000100000000000000C1中

postgres=*# select pg_walfile_name(pg_current_wal_lsn());
     pg_walfile_name      
--------------------------
 0000000100000000000000C1
(1 row)

postgres=*# insert into t1 values(9999);
INSERT 0 1
postgres=*# insert into t1 values(10000);
INSERT 0 1
postgres=*# select relfilenode from pg_class where relname = 't1';
 relfilenode 
-------------
       16578
(1 row)

使用pg_waldump解析也可以看到这两条数据

[postgres@xiongcc pg_wal]$ pg_waldump 0000000100000000000000C1
...
rmgr: Heap        len (rec/tot):     64/   160, tx:        915, lsn: 0/C120DFC0, prev 0/C120DF88, desc: INSERT off 2 flags 0x08, blkref #0: rel 1663/5/16578 blk 0 FPW
rmgr: Btree       len (rec/tot):     53/   133, tx:        915, lsn: 0/C120E078, prev 0/C120DFC0, desc: INSERT_LEAF off 2, blkref #0: rel 1663/5/16579 blk 1 FPW
rmgr: Heap        len (rec/tot):     59/    59, tx:        915, lsn: 0/C120E100, prev 0/C120E078, desc: INSERT off 3 flags 0x08, blkref #0: rel 1663/5/16578 blk 0
rmgr: Btree       len (rec/tot):     64/    64, tx:        915, lsn: 0/C120E140, prev 0/C120E100, desc: INSERT_LEAF off 3, blkref #0: rel 1663/5/16579 blk 1
...

现在让我们模拟一下生产中的操作,比如将最新的0000000100000000000000C1手动移除一下,再去提交事务

[postgres@xiongcc pg_wal]$ mv 0000000100000000000000C1 0000000100000000000000C1.bak

发布端没有问题,数据可见

postgres=*# commit ;
COMMIT
postgres=# select * from t1;
  id   
-------
     1
  9999
 10000
(3 rows)

但是订阅端此时就不对了,数据一条也没有

postgres=# select * from t1;
 id 
----
(0 rows)

然后容易让人误解的地方来了,发布端查看复制槽的状态是active的,给人一种仍在消费的错觉

postgres=# select slot_name,active,active_pid,pg_walfile_name(restart_lsn) as wal_recv from pg_replication_slots ;
 slot_name | active | active_pid |         wal_recv         
-----------+--------+------------+--------------------------
 sub1      | t      |      23442 | 00000001000000000000009A
(1 row)

其实是因为发布订阅端在不断尝试重启中,pid在不断变化也可以证实这一点 👇

postgres=# select slot_name,active,active_pid,pg_walfile_name(restart_lsn) as wal_recv from pg_replication_slots ;
 slot_name | active | active_pid |         wal_recv         
-----------+--------+------------+--------------------------
 sub1      | t      |      23442 | 00000001000000000000009A
(1 row)

postgres=# select slot_name,active,active_pid,pg_walfile_name(restart_lsn) as wal_recv from pg_replication_slots ;
 slot_name | active | active_pid |         wal_recv         
-----------+--------+------------+--------------------------
 sub1      | t      |      23447 | 00000001000000000000009A
(1 row)

postgres=# select slot_name,active,active_pid,pg_walfile_name(restart_lsn) as wal_recv from pg_replication_slots ;
 slot_name | active | active_pid |         wal_recv         
-----------+--------+------------+--------------------------
 sub1      | t      |      24693 | 00000001000000000000009A
(1 row)

订阅端日志不停启动新的 logical replication worker

2023-01-03 12:22:59.470 CST [24895] LOG:  logical replication apply worker for subscription "sub1" has started
2023-01-03 12:23:10.893 CST [24895] ERROR:  could not receive data from WAL stream: ERROR:  could not find record while sending logically-decoded data: invalid magic number 0000 in log segment 0000000100000000000000C1, offset 0
2023-01-03 12:23:10.895 CST [23420] LOG:  background worker "logical replication worker" (PID 24895) exited with exit code 1
2023-01-03 12:23:10.902 CST [24900] LOG:  logical replication apply worker for subscription "sub1" has started
2023-01-03 12:23:24.161 CST [24900] ERROR:  could not receive data from WAL stream: ERROR:  could not find record while sending logically-decoded data: invalid magic number 0000 in log segment 0000000100000000000000C1, offset 0
2023-01-03 12:23:24.162 CST [23420] LOG:  background worker "logical replication worker" (PID 24900) exited with exit code 1
2023-01-03 12:23:24.168 CST [24905] LOG:  logical replication apply worker for subscription "sub1" has started
2023-01-03 12:23:37.071 CST [24905] ERROR:  could not receive data from WAL stream: ERROR:  could not find record while sending logically-decoded data: invalid magic number 0000 in log segment 0000000100000000000000C1, offset 0
2023-01-03 12:23:37.071 CST [23420] LOG:  background worker "logical replication worker" (PID 24905) exited with exit code 1
2023-01-03 12:23:37.077 CST [24911] LOG:  logical replication apply worker for subscription "sub1" has started
2023-01-03 12:23:48.859 CST [24911] ERROR:  could not receive data from WAL stream: ERROR:  could not find record while sending logically-decoded data: invalid magic number 0000 in log segment 0000000100000000000000C1, offset 0
2023-01-03 12:23:48.860 CST [23420] LOG:  background worker "logical replication worker" (PID 24911) exited with exit code 1
2023-01-03 12:23:48.866 CST [24916] LOG:  logical replication apply worker for subscription "sub1" has started

发布端也在不断重启

2023-01-03 12:24:27.135 CST [24928] ERROR:  XX000: could not find record while sending logically-decoded data: invalid magic number 0000 in log segment 0000000100000000000000C1, offset 0
2023-01-03 12:24:27.135 CST [24928] LOCATION:  XLogSendLogical, walsender.c:3063
2023-01-03 12:24:27.135 CST [24928] STATEMENT:  START_REPLICATION SLOT "sub1" LOGICAL 0/0 (proto_version '3', publication_names '"pub1"')
2023-01-03 12:24:27.163 CST [24961] LOG:  00000: starting logical decoding for slot "sub1"
2023-01-03 12:24:27.163 CST [24961] DETAIL:  Streaming transactions committing after 0/C120E2A0, reading WAL from 0/9A00FC60.
2023-01-03 12:24:27.163 CST [24961] LOCATION:  CreateDecodingContext, logical.c:572
2023-01-03 12:24:27.163 CST [24961] STATEMENT:  START_REPLICATION SLOT "sub1" LOGICAL 0/0 (proto_version '3', publication_names '"pub1"')
2023-01-03 12:24:27.164 CST [24961] LOG:  00000: logical decoding found consistent point at 0/9A00FC60
2023-01-03 12:24:27.164 CST [24961] DETAIL:  There are no running transactions.
2023-01-03 12:24:27.164 CST [24961] LOCATION:  SnapBuildFindSnapshot, snapbuild.c:1362
2023-01-03 12:24:27.164 CST [24961] STATEMENT:  START_REPLICATION SLOT "sub1" LOGICAL 0/0 (proto_version '3', publication_names '"pub1"')

现在让我们回想一下,之前我为什么会在《PostgreSQL DBA Daily 2.0》大图中写到导致WAL堆积的原因之一便是logical decoding下面的大事务 👇🏻

Image

逻辑复制也通过walsender实现,walsender不停地读取WAL日志,会对每一条WAL日志记录都进行解析,将解析出的元组按照事务进行分组(保存进ReorderBuffer),并按照事务的开始时间进行排序。在事务提交时,会对提交事务在ReorderBuffer中的所有信息使用相应的plugin进行解码,并将解码后的逻辑日志发送给备机或者接收工具(例如pg_recvlogical)。

对应到前面的例子,一个事务横跨了多个WAL,所以需要从事务最开始保留到事务提交所在的WAL才能确保解析是正确的。

另外我在复现的过程中对于 pg_logical/snapshots 目录下的东西来了兴趣

[postgres@xiongcc snapshots]$ ls -lrth
total 12K
-rw------- 1 postgres postgres 128 Jan  3 09:34 0-7000AA40.snap
-rw------- 1 postgres postgres 128 Jan  3 09:34 0-7000AA78.snap
-rw------- 1 postgres postgres 128 Jan  3 09:35 0-7000AB28.snap

这么多个snap是什么玩意?看名字貌似和快照有关,那这个文件是什么时候生成的呢?看看代码,根据路径来搜就行

/*
 * Serialize the snapshot 'builder' at the location 'lsn' if it hasn't already
 * been done by another decoding process.
 */

static void
SnapBuildSerialize(SnapBuild *builder, XLogRecPtr lsn)
{
 Size  needed_length;
 SnapBuildOnDisk *ondisk = NULL;
 char    *ondisk_c;
 int   fd;
 char  tmppath[MAXPGPATH];
 char  path[MAXPGPATH];
 int   ret;
 struct stat stat_buf;
 Size  sz;

 Assert(lsn != InvalidXLogRecPtr);
 Assert(builder->last_serialized_snapshot == InvalidXLogRecPtr ||
     builder->last_serialized_snapshot <= lsn);

 /*
  * no point in serializing if we cannot continue to work immediately after
  * restoring the snapshot
  */

 if (builder->state < SNAPBUILD_CONSISTENT)
  return;

 /* consistent snapshots have no next phase */
 Assert(builder->next_phase_at == InvalidTransactionId);

 /*
  * We identify snapshots by the LSN they are valid for. We don't need to
  * include timelines in the name as each LSN maps to exactly one timeline
  * unless the user used pg_resetwal or similar. If a user did so, there's
  * no hope continuing to decode anyway.
  */

 sprintf(path, "pg_logical/snapshots/%X-%X.snap",
   LSN_FORMAT_ARGS(lsn));
  ...
  ...

可以看到该函数根据LSN将快照定时刷到存储上,其中快照里面记录了一些必要的信息,比如

  • xmin/xmax/xip_list
  • 快照开始时达到一致性状态的LSN,start_decoding_at
  • 子事务的修改信息
  • 系统表元组xmin,因为逻辑解码依赖系统表元组的可见性,某个事务的修改要对其他事务可见
  • ...

在逻辑解码的过程中,下面这个SnapBuild就会一直不断被更新,因此也会涉及到脏数据

/*
 * This struct contains the current state of the snapshot building
 * machinery. Besides a forward declaration in the header, it is not exposed
 * to the public, so we can easily change its contents.
 */

struct SnapBuild
{

 /* how far are we along building our first full snapshot */
 SnapBuildState state;

 /* private memory context used to allocate memory for this module. */
 MemoryContext context;

 /* all transactions < than this have committed/aborted */
 TransactionId xmin;

 /* all transactions >= than this are uncommitted */
 TransactionId xmax;

 /*
  * Don't replay commits from an LSN < this LSN. This can be set externally
  * but it will also be advanced (never retreat) from within snapbuild.c.
  */

 XLogRecPtr start_decoding_at;

 /*
  * LSN at which we found a consistent point at the time of slot creation.
  * This is also the point where we have exported a snapshot for the
  * initial copy.
  *
  * The prepared transactions that are not covered by initial snapshot
  * needs to be sent later along with commit prepared and they must be
  * before this point.
  */

 XLogRecPtr initial_consistent_point;

 /*
  * Don't start decoding WAL until the "xl_running_xacts" information
  * indicates there are no running xids with an xid smaller than this.
  */

 TransactionId initial_xmin_horizon;

 /* Indicates if we are building full snapshot or just catalog one. */
 bool  building_full_snapshot;

 /*
  * Snapshot that's valid to see the catalog state seen at this moment.
  */

 Snapshot snapshot;

 /*
  * LSN of the last location we are sure a snapshot has been serialized to.
  */

 XLogRecPtr last_serialized_snapshot;

 /*
  * The reorderbuffer we need to update with usable snapshots et al.
  */

 ReorderBuffer *reorder;

 /*
  * TransactionId at which the next phase of initial snapshot building will
  * happen. InvalidTransactionId if not known (i.e. SNAPBUILD_START), or
  * when no next phase necessary (SNAPBUILD_CONSISTENT).
  */

 TransactionId next_phase_at;

 /*
  * Array of transactions which could have catalog changes that committed
  * between xmin and xmax.
  */

 struct
 {

  /* number of committed transactions */
  size_t  xcnt;

  /* available space for committed transactions */
  size_t  xcnt_space;

  /*
   * Until we reach a CONSISTENT state, we record commits of all
   * transactions, not just the catalog changing ones. Record when that
   * changes so we know we cannot export a snapshot safely anymore.
   */

  bool  includes_all_transactions;

  /*
   * Array of committed transactions that have modified the catalog.
   *
   * As this array is frequently modified we do *not* keep it in
   * xidComparator order. Instead we sort the array when building &
   * distributing a snapshot.
   *
   * TODO: It's unclear whether that reasoning has much merit. Every
   * time we add something here after becoming consistent will also
   * require distributing a snapshot. Storing them sorted would
   * potentially also make it easier to purge (but more complicated wrt
   * wraparound?). Should be improved if sorting while building the
   * snapshot shows up in profiles.
   */

  TransactionId *xip;
 }   committed;
};

而文件的命名则根据传入的 lsn 计算,很简单就是一个宏

/*
 * Handy macro for printing XLogRecPtr in conventional format, e.g.,
 *
 * printf("%X/%X", LSN_FORMAT_ARGS(lsn));
 */

#define LSN_FORMAT_ARGS(lsn) (AssertVariableIsOfTypeMacro((lsn), XLogRecPtr), (uint32) ((lsn) >> 32)), ((uint32) (lsn))

由于lsn是一个64位的数字,右移32位就是斜杠左边的数字,左边32位,右边32位,加起来64位,能够支持4GB*4GB的日志,所以右移32位就是斜杠左边的数字,第二个数字就是斜杠右边的数字,比如目前存在这样一个复制槽 👇

postgres=# select restart_lsn,confirmed_flush_lsn,slot_name,active,active_pid from pg_replication_slots ;
 restart_lsn | confirmed_flush_lsn | slot_name | active | active_pid 
-------------+---------------------+-----------+--------+------------
 0/E9006D48  | 0/E9006D80          | sub1      | t      |      26786
(1 row)

postgres=# SELECT * FROM pg_ls_logicalsnapdir() order by modification;
      name       | size |      modification      
-----------------+------+------------------------
 0-E9006C60.snap |  128 | 2023-01-03 14:04:20+08
 0-E9006C98.snap |  128 | 2023-01-03 14:05:58+08
 0-E9006D48.snap |  128 | 2023-01-03 14:06:04+08
(3 rows)

那根据前面的分析,这个文件名应该是0-E9006D48.snap,确认一下

[postgres@xiongcc snapshots]$ ls -lrth
total 12K
-rw------- 1 postgres postgres 128 Jan  3 14:04 0-E9006C60.snap
-rw------- 1 postgres postgres 128 Jan  3 14:05 0-E9006C98.snap
-rw------- 1 postgres postgres 128 Jan  3 14:06 0-E9006D48.snap  ---👈🏻文件在这里

果然如我们所料,当然有存储就有还原,SnapBuildRestore函数用于将存储上的快照恢复

/*
 * Restore a snapshot into 'builder' if previously one has been stored at the
 * location indicated by 'lsn'. Returns true if successful, false otherwise.
 */

static bool
SnapBuildRestore(SnapBuild *builder, XLogRecPtr lsn)
{
 SnapBuildOnDisk ondisk;
 int   fd;
 char  path[MAXPGPATH];
 Size  sz;
 int   readBytes;
 pg_crc32c checksum;

 /* no point in loading a snapshot if we're already there */
 if (builder->state == SNAPBUILD_CONSISTENT)
  return false;

 sprintf(path, "pg_logical/snapshots/%X-%X.snap",
   LSN_FORMAT_ARGS(lsn));

 fd = OpenTransientFile(path, O_RDONLY | PG_BINARY);
  ...

假如恢复过,在日志中就会有类似的打印,比如重启之后,数据库要找到一个一致性位点

"logical decoding found consistent point at xxx","Logical decoding will begin using saved snapshot."

其实要简单理解的话,就把这些数据理解成逻辑复制相关的脏数据,所以也会按照我们熟知的逻辑,在checkpoint的时候进行回收

/*
 * Remove all serialized snapshots that are not required anymore because no
 * slot can need them. This doesn't actually have to run during a checkpoint,
 * but it's a convenient point to schedule this.
 *
 * NB: We run this during checkpoints even if logical decoding is disabled so
 * we cleanup old slots at some point after it got disabled.
 */

void
CheckPointSnapBuild(void)
{
 XLogRecPtr cutoff;
 XLogRecPtr redo;
 DIR     *snap_dir;
 struct dirent *snap_de;
 char  path[MAXPGPATH + 21];
  ...

验证一下

postgres=# insert into t1 values(100);
INSERT 0 1
postgres=# SELECT * FROM pg_ls_logicalsnapdir() order by modification;
      name       | size |      modification      
-----------------+------+------------------------
 0-E9007230.snap |  128 | 2023-01-03 14:39:53+08
 0-E9007430.snap |  128 | 2023-01-03 14:40:07+08
 0-E90074E0.snap |  128 | 2023-01-03 14:40:13+08
(3 rows)

postgres=# insert into t1 values(101);
INSERT 0 1
postgres=# SELECT * FROM pg_ls_logicalsnapdir() order by modification;
      name       | size |      modification      
-----------------+------+------------------------
 0-E9007230.snap |  128 | 2023-01-03 14:39:53+08
 0-E9007430.snap |  128 | 2023-01-03 14:40:07+08
 0-E90074E0.snap |  128 | 2023-01-03 14:40:13+08
 0-E90077D0.snap |  128 | 2023-01-03 14:44:18+08
(4 rows)

postgres=# select restart_lsn,confirmed_flush_lsn,slot_name,active,active_pid from pg_replication_slots ;
 restart_lsn | confirmed_flush_lsn | slot_name | active | active_pid 
-------------+---------------------+-----------+--------+------------
 0/E90077D0  | 0/E9007808          | sub1      | t      |      26786
(1 row)

postgres=# checkpoint ;   ---checkpoint的时候回收相关脏数据
CHECKPOINT
postgres=# select restart_lsn,confirmed_flush_lsn,slot_name,active,active_pid from pg_replication_slots ;
 restart_lsn | confirmed_flush_lsn | slot_name | active | active_pid 
-------------+---------------------+-----------+--------+------------
 0/E9007808  | 0/E90078B8          | sub1      | t      |      26786
(1 row)

postgres=# SELECT * FROM pg_ls_logicalsnapdir() order by modification;  ---snap文件被回收
      name       | size |      modification      
-----------------+------+------------------------
 0-E90077D0.snap |  128 | 2023-01-03 14:44:18+08
 0-E9007808.snap |  128 | 2023-01-03 14:44:26+08
(2 rows)

postgres=# insert into t1 values(102);
INSERT 0 1
postgres=# SELECT * FROM pg_ls_logicalsnapdir() order by modification;
      name       | size |      modification      
-----------------+------+------------------------
 0-E90077D0.snap |  128 | 2023-01-03 14:44:18+08
 0-E9007808.snap |  128 | 2023-01-03 14:44:26+08
 0-E90078B8.snap |  128 | 2023-01-03 14:44:38+08
(3 rows)

根据代码,将日志设为DEBUG1就可以看到相关信息,elog(DEBUG1, "serializing snapshot to %s", path);

2023-01-03 15:30:24.443 CST [28399] DEBUG:  00000: serializing snapshot to pg_logical/snapshots/0-EB008870.snap  ---同步中
2023-01-03 15:30:24.443 CST [28399] LOCATION:  SnapBuildSerialize, snapbuild.c:1659
2023-01-03 15:30:24.449 CST [28399] DEBUG:  00000: "sub1" has now caught up with upstream server
2023-01-03 15:30:24.449 CST [28399] LOCATION:  WalSndLoop, walsender.c:2525

4小结

小结一下,由于同一个事务的WAL日志是不连续的,并且可能横跨多个WAL(如本文中的案例),而逻辑解码要求同一个事务的日志按顺序相邻,因此会借助ReorderBuffer结构体来保存这些信息,同时用ReorderBufferTXN标识每一个事务(xid和TXN进行映射),在事务提交时,借助decode plugin对元组进行逻辑解码,比如自带的test_decoding、decoder_raw等等,按plugin自己定制解析后输出逻辑解析成品,同一个ReorderBufferTXN中的操作会全部发送给订阅端。

因此假如同一个事务的日志不连续或者缺失了,就会使逻辑解码这个过程出现异常,所以切忌手动操作WAL,十分危险。

目前两个群聊都满200人了,需要加群只能我手动拉了。