pg_walreceiver


难度 中等

启动流程

启动receiver进程

graph TB main-->PostmasterMain-->StartChildProcess-->fork_process-->AuxiliaryProcessMain-->WalReceiverMain

连接主库

graph TB WalReceiverMain-->walrcv_connect-->libpqrcv_connect

wal_receiver整体流程

wal_receiver整体流程由三层循环组成,第一层循环使用launch机制,第二层和第三层则使用for(;;)实现。

  • 第一层循环是wal_receiver进程的控制循环,负责挂起和监听socket是否有数据到来
  • 第二层循环是流复制流程,从一次流复制开始到流复制结束或流复制出现错误
  • 第三层循环是wal数据和心跳数据接收,不断的从socket的buffer中收取数据,直到buffer中没有数据

wal数据接收流程

graph TB loop-->walrcv_receive-->XLogWalRcvProcessMsg-->ProcessWalSndrMessage-->|收到wal数据|XLogWalRcvWrite-->pg_pwrite-->loop ProcessWalSndrMessage-->|收到心跳|XLogWalRcvSendReply-->loop
  1. wal数据接收在walreceiver的第三层循环,通过循环调用walrcv_receive收取数据,当返回的数据长度小于等于0时退出;其中等于0时表示没有数据,小于0表示收包遇到错误;

    for (;;)
    {
        if (len > 0)
        {
            /*
             * Something was received from primary, so reset
             * timeout
             */
            last_recv_timestamp = GetCurrentTimestamp();
            ping_sent = false;
            XLogWalRcvProcessMsg(buf[0], &buf[1], len - 1,
                                 startpointTLI);
        }
        else if (len == 0)
            break;
        else if (len < 0)
        {
            ereport(LOG,
                    (errmsg("replication terminated by primary server"),
                     errdetail("End of WAL reached on timeline %u at %X/%X.",
                               startpointTLI,
                               LSN_FORMAT_ARGS(LogstreamResult.Write))));
            endofwal = true;
            break;
        }
        len = walrcv_receive(wrconn, &buf, &wait_fd);
    }
  2. 当收到数据后,会根据报文的第一个字节区分wal数据和心跳,并进行区别处理;

    switch (type)
    {
        case 'w':                /* WAL records */
            {
                /* copy message to StringInfo */
                hdrlen = sizeof(int64) + sizeof(int64) + sizeof(int64);
                if (len < hdrlen)
                    ereport(ERROR,
                            (errcode(ERRCODE_PROTOCOL_VIOLATION),
                             errmsg_internal("invalid WAL message received from primary")));
                appendBinaryStringInfo(&incoming_message, buf, hdrlen);
       
                /* read the fields */
                dataStart = pq_getmsgint64(&incoming_message);
                walEnd = pq_getmsgint64(&incoming_message);
                sendTime = pq_getmsgint64(&incoming_message);
                ProcessWalSndrMessage(walEnd, sendTime);
       
                buf += hdrlen;
                len -= hdrlen;
                XLogWalRcvWrite(buf, len, dataStart, tli);
                break;
            }
        case 'k':                /* Keepalive */
            {
                /* copy message to StringInfo */
                hdrlen = sizeof(int64) + sizeof(int64) + sizeof(char);
                if (len != hdrlen)
                    ereport(ERROR,
                            (errcode(ERRCODE_PROTOCOL_VIOLATION),
                             errmsg_internal("invalid keepalive message received from primary")));
                appendBinaryStringInfo(&incoming_message, buf, hdrlen);
       
                /* read the fields */
                walEnd = pq_getmsgint64(&incoming_message);
                sendTime = pq_getmsgint64(&incoming_message);
                replyRequested = pq_getmsgbyte(&incoming_message);
       
                ProcessWalSndrMessage(walEnd, sendTime);
       
                /* If the primary requested a reply, send one immediately */
                if (replyRequested)
                    XLogWalRcvSendReply(true, false);
                break;
            }
        default:
            ereport(ERROR,
                    (errcode(ERRCODE_PROTOCOL_VIOLATION),
                     errmsg_internal("invalid replication message type %d",
                                     type)));
    }

    如果是wal报文,则进行写入;如果是心跳报文需要进一步读取报文的最后一个字节是否为1,为1表示需要回复心跳报文。

  3. 当一次数据收取结束后(socket的buffer中没有数据),此时walreceiver也需要进行回复,通知walsender收到了数据;

    /* Let the primary know that we received some data. */
    XLogWalRcvSendReply(false, false);
  4. 同时如果写入了数据,需要将数据flush到磁盘

    /*
     * If we've written some records, flush them to disk and
     * let the startup process and primary server know about
     * them.
     */
    XLogWalRcvFlush(false, startpointTLI);

文章作者: growdu
版权声明: 本博客所有文章除特別声明外,均采用 CC BY 4.0 许可协议。转载请注明来源 growdu !
  目录
分类导航
随笔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