PostgreSQL数据库复制——后台一等公民进程WalReceiver&startup交互_postgressql walreceive线程_肥叔菌的博客-CSDN博客


难度 中等

在这里插入图片描述
在PostgreSQL流复制过程中,有三个进程协同工作:walsender进程,walreceiver进程和startup进程。其中walsender进程属于主节点的进程,主要用来向备节点发送wal record; walreceiver和startup进程属于备节点进程,wal receiver主要用来接收主端发送来的wal record并写入磁盘上的XLOG文件中,之后startup进程就会对这些wal数据进行replay。三个进程共同协作,完成主备的整个流复制过程。本篇博客主要关注于WalReceiver进程和startup进程之间的交互逻辑。首先看一下这三个进程的调用堆栈,可以更加方便定位所需阅读的代码细节:

walsender进程是用来发送WAL日志记录的,执行顺序如下:

PostgresMain()->exec_replication_command()->StartReplication()->WalSndLoop()->XLogSendPhysical()

graph TB PostgresMain-->exec_replication_command-->StartReplication -->WalSndLoop-->XLogSendPhysical

walreceiver进程是用来接收WAL日志记录的,执行顺序如下:sigusr1_handler()->StartWalReceiver()->AuxiliaryProcessMain()->WalReceiverMain()->walrcv_receive()

graph TB sigusr1_handler-->StartWalReceiver-->AuxiliaryProcessMain -->WalReceiverMain-->walrcv_receive

startup进程是用来apply日志的,执行顺序如下:PostmasterMain()->StartupDataBase()->AuxiliaryProcessMain()->StartupProcessMain()->StartupXLOG()

graph TB PostmasterMain-->StartupDataBase-->AuxiliaryProcessMain -->StartupProcessMain-->StartupXLOG

在流复制启动过程中,三个进程的启动顺序是从备库到主库,即:startup —> walreceiver —> walsender。但是值得注意的是Startup启动后,不会马上发送信号给postmaster来启动wal receiver进程,它先会进行一系列条件的判断然后决定是否通知postmaster启动wal receiver进程。 我们知道Startup进程回放日志所需要的WAL文件有3个来源:归档中获取、pg_wal文件夹下获取、从primary 节点以流复制方式获取。在实际流复制过程中, 如果是非归档,则先会从pg_wal中获取;否则优先从archive归档中获取(Archive Mode); 如果两者都没有,startup要恢复的wal,只能从primary 节点以流复制方式获取,这时startup会发送信号(通过函数SendPostmasterSignal(PMSIGNAL_START_WALRECEIVER) ,可以查看之前的博文了解这段过程)给postmaster进程,请求其启动wal receiver进程从Primary节点来获取wal数据。

在流复制运行中,WAL数据的流向则是walsender进程占据主动位置:walsender —> walreceiver —> startup。从主库backend执行业务操作所产生的XLOG会顺着上述流程从主库walsender进程网络发送到walreceiver网络接收并落盘,最终备库startup进程会对XLOG进行应用。

startup进程主要流程

startup进程进入standby模式和apply日志主要过程:

  1. 读取pg_control文件,找到redo位点;读取recovery.conf,如果配置standby_mode=on则进入standby模式。
  2. 如果是Hot Standby需要初始化clog、subtrans、事务环境等。初始化redo资源管理器,比如Heap、Heap2、Database、XLOG等。
  3. 读取WAL record,如果record不存在需要调用XLogPageRead->WaitForWALToBecomeAvailable->RequestXLogStreaming唤醒walreceiver从walsender获取WAL record。
  4. 对读取的WAL record进行redo,通过record->xl_rmid信息,调用相应的redo资源管理器进行redo操作。比如heap_redo的XLOG_HEAP_INSERT操作,就是通过record的信息在buffer page中增加一个record。还有部分redo操作(vacuum产生的record)需要检查在Hot Standby模式下的查询冲突,比如某些tuples需要remove,而存在正在执行的query可能读到这些tuples,这样就会破坏事务隔离级别。通过函数ResolveRecoveryConflictWithSnapshot检测冲突,如果发生冲突,那么就把这个query所在的进程kill掉。
  5. 检查一致性,如果一致了,Hot Standby模式可以接受用户只读查询;更新共享内存中XLogCtlData的apply位点和时间线;如果恢复到时间点,时间线或者事务id需要检查是否恢复到当前目标;
  6. 回到步骤3,读取next WAL record。
    在这里插入图片描述

WalReceiver进程主要流程

WalReceiver进程进入工作状态后主要执行流程如下所示:

  1. WalReceiver进程进入流复制之前,startup进程已经recovery.conf文件中的primary_conninfo参数信息解析后填充walrcv->conninfowalrcv->slotnamewalrcv->receiveStartwalrcv->receiveStartTLI等共享内存WalRcvData中。 WalReceiver进程通过walrcv_connect连接主库。
  2. 进入第一层死循环,执行identify_system命令,获取主库systemid/timeline/xlogpos等信息,执行TIMELINE_HISTORY命令拉取history文件,如果需要,则创建临时复制槽。
  3. (3.1)执行wal_startstreaming开始启动流复制
    (3.1.1)通过walrcv_receive获取WAL日志,并向wal segment文件中写入xlog,期间也会回应主库发过来的心跳信息(发送write位点、flush位点、apply位点)
    (3.1.2)如果walrcv_receive获取不到数据,向主库发送write位点、flush位点、apply位点,发送feedback信息(xmin、xmin_epoch、catalog_xmin、catalog_xmin_epoch),避免vacuum删掉备库正在使用的记录,如果flush了XLOG逻辑位置唤醒startup进程,跳转到4
    (3.1.3)如果walrcv_receive获取EOF,向主库发送write位点、flush位点、apply位点,发送feedback信息(xmin、xmin_epoch、catalog_xmin、catalog_xmin_epoch),避免vacuum删掉备库正在使用的记录,如果flush了XLOG逻辑位置唤醒startup进程,跳转到5
    (3.2)如果wal_startstreaming返回false,说明主库在该时间线上已经没有wal可以发送了,跳转到6
  4. WaitLatchOrSocket等待超时/网络可读/latch被触发:
    (4.1)等待超时,向主库发送接收位点、flush位点、apply位点,发送feedback信息(xmin、xmin_epoch、catalog_xmin、catalog_xmin_epoch),跳转到3.1.x
    (4.2)startup进程通过latch唤醒WalReceiver进程,向主库发送write位点、flush位点、apply位点,跳转到3.1.x
  5. 执行walrcv_endstreaming结束流复制,接收wal sender发送过来的时间线文件,进入步骤6
  6. 关闭目前的xlog segment文件描述符,发送feedback信息(xmin、xmin_epoch、catalog_xmin、catalog_xmin_epoch),如果flush了XLOG逻辑位置唤醒startup进程,进入步骤7
  7. WalRcvWaitForStartPosition函数等待startup进程更新receiveStart和receiveStartTLI,一旦更新,进入步骤2。
    在这里插入图片描述

WalReceiver&startup交互

graph TB XLogPageRead-->WaitForWALToBecomeAvailable-->RequestXLogStreaming -->|发送PMSIGNAL_START_WALRECEIVER信号启动walreceiver|SendPostmasterSignal -->sigusr1_handler-->StartWalReceiver-->AuxiliaryProcessMain -->WalReceiverMain-->walrcv_receive
  1. startup进程向postmaster请求启动WalReceiver进程
    startup进程通过如下WaitForWALToBecomeAvailable->RequestXLogStreaming流程设置receiveStart和receiveStartTLI,要求WalReceiver进行流复制。
void RequestXLogStreaming(TimeLineID tli, XLogRecPtr recptr, const char *conninfo, const char *slotname, bool create_temp_slot) {
WalRcvData *walrcv = WalRcv;
...
walrcv->receiveStart = recptr;
walrcv->receiveStartTLI = tli;
latch = walrcv->latch;
SpinLockRelease(&walrcv->mutex);
    // startup进程将WalReceiver状态从WALRCV_STOPPED转变为WALRCV_STARTING时,设置launch为true,向postmaster请求启动WalReceiver进程
    // startup进程通过latch唤醒WalReceiver进程
if (launch) SendPostmasterSignal(PMSIGNAL_START_WALRECEIVER);
else if (latch) SetLatch(latch);
}
  1. walReceiver进程调用XLogWalRcvFlush函数如果已经flush了XLOG逻辑位置则唤醒startup进程
static void XLogWalRcvFlush(bool dying) {
if (LogstreamResult.Flush < LogstreamResult.Write) {
WalRcvData *walrcv = WalRcv;
issue_xlog_fsync(recvFile, recvSegNo);
LogstreamResult.Flush = LogstreamResult.Write;
/* Update shared-memory status */
SpinLockAcquire(&walrcv->mutex);
if (walrcv->flushedUpto < LogstreamResult.Flush) {
walrcv->latestChunkStart = walrcv->flushedUpto;
walrcv->flushedUpto = LogstreamResult.Flush;
walrcv->receivedTLI = ThisTimeLineID;
}
SpinLockRelease(&walrcv->mutex);
WakeupRecovery(); /* Signal the startup process and walsender that new WAL has arrived */ // 唤醒startup进程
if (AllowCascadeReplication()) WalSndWakeup();
if (!dying) { /* Also let the primary know that we made some progress */
XLogWalRcvSendReply(false, false);
XLogWalRcvSendHSFeedback(false);
}
}
}
  1. startup进程要求WalReceiver进行现在发送应用反馈,每当应用interesting xlog records时,startup进程都会调用此方法,以便 walreceiver 可以检查它是否需要将apply通知(notification)发送回主节点,主库可能在应用了synchronous_commit = remote_apply参数的COMMIT中等待walreceiver的反馈。从这里可以看到WalReceiver除了发送write位点、flush位点、apply位点消息之外,针对remote_apply从库应用之后才commit这种情况,增加了startup进程强制WalReceiver唤醒发送反馈这一特性,加速主库进行commit操作。
void WalRcvForceReply(void) {
Latch   *latch;
WalRcv->force_reply = true;  // 设置强制回复标志
SpinLockAcquire(&WalRcv->mutex); /* fetching the latch pointer might not be atomic, so use spinlock */
latch = WalRcv->latch;
SpinLockRelease(&WalRcv->mutex);
if (latch) SetLatch(latch);
}
  1. walReceiver进程WalRcvWaitForStartPosition函数等待startup进程设置receiveStart和receiveStartTLI
static void WalRcvWaitForStartPosition(XLogRecPtr *startpoint, TimeLineID *startpointTLI) {
WalRcvData *walrcv = WalRcv;
intstate;
SpinLockAcquire(&walrcv->mutex);
state = walrcv->walRcvState;
if (state != WALRCV_STREAMING) {
SpinLockRelease(&walrcv->mutex);
if (state == WALRCV_STOPPING) proc_exit(0);
else elog(FATAL, "unexpected walreceiver state");
}
walrcv->walRcvState = WALRCV_WAITING;
walrcv->receiveStart = InvalidXLogRecPtr;
walrcv->receiveStartTLI = 0;
SpinLockRelease(&walrcv->mutex);
/* nudge startup process to notice that we've stopped streaming and are now waiting for instructions. */
WakeupRecovery();
for (;;) {
ResetLatch(MyLatch);
ProcessWalRcvInterrupts();
SpinLockAcquire(&walrcv->mutex);
if (walrcv->walRcvState == WALRCV_RESTARTING) {
/* No need to handle changes in primary_conninfo or primary_slotname here. Startup process will signal us to terminate in case those change. */
*startpoint = walrcv->receiveStart;
*startpointTLI = walrcv->receiveStartTLI;
walrcv->walRcvState = WALRCV_STREAMING;
SpinLockRelease(&walrcv->mutex);
break;
}
if (walrcv->walRcvState == WALRCV_STOPPING) { /* We should've received SIGTERM if the startup process wants us to die, but might as well check it here too. */
SpinLockRelease(&walrcv->mutex);
exit(1);
}
SpinLockRelease(&walrcv->mutex);
(void) WaitLatch(MyLatch, WL_LATCH_SET | WL_EXIT_ON_PM_DEATH, 0, WAIT_EVENT_WAL_RECEIVER_WAIT_START);
}
}

startup进程则通过WaitForWALToBecomeAvailable->RequestXLogStreaming流程设置receiveStart和receiveStartTLI,并且是通过latch唤醒WalReceiver进程。

  1. walReceiver进程关闭
    (1) startup进程执行如下ShutdownWalRcv函数请求walreceiver进程关闭,并通过walRcvStoppedCV监视walReceiver进程关闭,等待 walreceiver 通过将状态设置为 WALRCV_STOPPED 来确认其死亡。
void ShutdownWalRcv(void) {
WalRcvData *walrcv = WalRcv;
pid_twalrcvpid = 0;
boolstopped = false;
/* Request walreceiver to stop. Walreceiver will switch to WALRCV_STOPPED mode once it's finished, and will also request postmaster to not restart itself.*/ // 请求walreceiver进程关闭
SpinLockAcquire(&walrcv->mutex);
switch (walrcv->walRcvState) {
case WALRCV_STOPPED: break;
case WALRCV_STARTING:
walrcv->walRcvState = WALRCV_STOPPED;
stopped = true;
break;
case WALRCV_STREAMING:
case WALRCV_WAITING:
case WALRCV_RESTARTING:
walrcv->walRcvState = WALRCV_STOPPING;
/* fall through */
case WALRCV_STOPPING:
walrcvpid = walrcv->pid;
break;
}
SpinLockRelease(&walrcv->mutex);
/* Unnecessary but consistent. */
if (stopped) ConditionVariableBroadcast(&walrcv->walRcvStoppedCV);

/* Signal walreceiver process if it was still running. */ // 向walreceiver进程发送SIGTERM信号
if (walrcvpid != 0) kill(walrcvpid, SIGTERM);

/* Wait for walreceiver to acknowledge its death by setting state to WALRCV_STOPPED. */ // 等待 walreceiver 通过将状态设置为 WALRCV_STOPPED 来确认其死亡
ConditionVariablePrepareToSleep(&walrcv->walRcvStoppedCV);
while (WalRcvRunning())
ConditionVariableSleep(&walrcv->walRcvStoppedCV, WAIT_EVENT_WAL_RECEIVER_EXIT);
ConditionVariableCancelSleep();
}

walReceiver进程信号处理函数WalRcvDie向walRcvStoppedCV条件变量进行广播,通知startup进程。

static void WalRcvDie(int code, Datum arg) {
WalRcvData *walrcv = WalRcv;
XLogWalRcvFlush(true); /* Ensure that all WAL records received are flushed to disk */
SpinLockAcquire(&walrcv->mutex); /* Mark ourselves inactive in shared memory */
walrcv->walRcvState = WALRCV_STOPPED;
walrcv->pid = 0;
walrcv->ready_to_display = false;
walrcv->latch = NULL;
SpinLockRelease(&walrcv->mutex);
ConditionVariableBroadcast(&walrcv->walRcvStoppedCV);
if (wrconn != NULL) walrcv_disconnect(wrconn); /* Terminate the connection gracefully. */
WakeupRecovery(); /* Wake up the startup process to notice promptly that we're gone */
}

walReceiver进程WalReceiverMain处理walRcvState状态向walRcvStoppedCV条件变量进行广播,通知startup进程。
在这里插入图片描述

(2) startup进程WaitForWALToBecomeAvailable会通过WalRcvStreaming从WALRCV_STARTING状态向Streaming转变或是否存活。
startup进程通过WalRcvRunning函数监控walReceiver进程从WALRCV_STARTING状态向Running转变超时或是否存活,会向walRcvStoppedCV条件变量进行广播,以保证退出ShutdownWalRcv函数的这个sleep循环。或者在启动进程重载配置文件时StartupRequestWalReceiverRestart函数会调用WalRcvRunning函数。
在这里插入图片描述

图片和部分内容参考自阿里内核月报 https://www.kancloud.cn/taobaomysql/monthly/81110
在这里衷心感谢阿里内核团队的分享


  目录
分类导航
随笔2 AI27 算法1 计算机基础13 博客搭建7 ChatGPT2 集群63 计算机通信1 数据库34 数据库深入80 DPDK26 Docker11 Elasticsearch4 编辑工具4 FAQ1 Go Web1 hometown2 编程语言16 网络9 OPC1 Linux38 openGauss4 页面12 PostgreSQL54 程序员自我修养1 协议11 成长之路1 stock1 存储5 工具20 VPP18 视频作品1 Vue13 Web1 代码示例11 数据库15 BenchmarkSQL1 PostgreSQL 源码修炼之路14
最热文章
1
13 逻辑复制深入
数据库深入🔥 1570
2
0 Postgresql存储、索引及系统优化、主备切换
PostgreSQL🔥 1495
3
一文读懂openguass dcf网络模块
集群🔥 1420
4
逻辑复制源码分析
数据库深入🔥 1327
5
PostgreSQL 分区表:从一行 `PARTITION BY` 到路由热路径的全链路拆解
数据库🔥 1094
6
applyparallelworker.c 之 LA 端源码深度解析:Leader Apply Worker 的指挥中枢
数据库深入🔥 1082
7
PostgreSQL Background Worker 全解:从 `RegisterBackgroundWorker` 到逻辑复制 4 类 worker 的全生命周期
数据库🔥 1078
8
PostgreSQL的后台进程walsender分析 - 关系型数据库 - 亿速云
PostgreSQL🔥 1033
9
PostgreSQL 逻辑复制的监控:六张视图 + 一组可执行 SQL,把 publisher/subscriber 的速率与健康度彻底看透
数据库🔥 1032
10
PostgreSQL 逻辑复制支持 DDL 之后:DDL 与 DML 的时序难题(重点:分区表)
数据库🔥 999
11
reorderbuffer.c 源码深度解析:PostgreSQL 逻辑复制的"事务重组引擎
数据库深入🔥 953
12
PostgreSQL 内核开发:读取一张表的 9 步标准流程与缓存全景
数据库🔥 938
13
从 `postgres` 二进制到生产级守护 —— PostgreSQL 最外层模块与启动全流程拆解
数据库🔥 936
14
支持逻辑复制同步 DDL 适配 SQL Server 方案
数据库深入🔥 934
15
PostgreSQL 逻辑复制的 ReorderBuffer 与事务机制:从一行 WAL 到一致性变更流的全链路绑定
数据库🔥 913
16
DDL同步架构(美化版)
数据库深入🔥 908
17
PostgreSQL Latch 机制详解:从一行 SetLatch 到 epoll 的内核之旅
数据库🔥 871
18
pgbench 源码全解:一个 C 文件如何撑起 PostgreSQL 官方压测工具
数据库🔥 860
19
PostgreSQL libpq 机制与缓冲区详解
数据库🔥 850
20
PostgreSQL 逻辑复制 spill 文件深度剖析:从 `xid-*.spill` 到 TPC-C 的增长方程
数据库🔥 845