Postgresql源码(117)libpq的两套实现(socket/shm_mq)

2023-12-19 12:52

本文主要是介绍Postgresql源码(117)libpq的两套实现(socket/shm_mq),希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

libpq的通信方式

libpq提供了两套通信方式

  • socket
  • shm_mq

分别实现在下面两个文件中

  • pqcomm.c
  • pqmq.c

什么时候用socket通信?

除了下述并行场景,其他场景全部使用socket通信。

static const PQcommMethods PqCommSocketMethods = {.comm_reset = socket_comm_reset,.flush = socket_flush,.flush_if_writable = socket_flush_if_writable,.is_send_pending = socket_is_send_pending,.putmessage = socket_putmessage,.putmessage_noblock = socket_putmessage_noblock
};

什么时候使用mq通信?

并行框架中会将子进程的libpq的通信改成mq通信,用于子进程给父进程发送错误信息。

static const PQcommMethods PqCommMqMethods = {.comm_reset = mq_comm_reset,.flush = mq_flush,.flush_if_writable = mq_flush_if_writable,.is_send_pending = mq_is_send_pending,.putmessage = mq_putmessage,.putmessage_noblock = mq_putmessage_noblock
};

使用MQ通信需要用pq_redirect_to_shm_mq函数指定使用的dsm和mq。

注意这个pq_mq_handle是申请在dsm上的,专门用于并行框架。

void
pq_redirect_to_shm_mq(dsm_segment *seg, shm_mq_handle *mqh)
{PqCommMethods = &PqCommMqMethods;pq_mq_handle = mqh;whereToSendOutput = DestRemote;FrontendProtocol = PG_PROTOCOL_LATEST;on_dsm_detach(seg, pq_cleanup_redirect_to_shm_mq, (Datum) 0);
}

使用位置在并行框架子进程入口

void ParallelWorkerMain(...)
{......// 拿到父进程在共享内存中申请mq的内存其实地址error_queue_space = shm_toc_lookup(toc, PARALLEL_KEY_ERROR_QUEUE, false);// 用自己的woker num偏移得到自己的mqmq = (shm_mq *) (error_queue_space +ParallelWorkerNumber * PARALLEL_ERROR_QUEUE_SIZE);// 配置PGPROC到mq上shm_mq_set_sender(mq, MyProc);// mq是单纯的mq抽象,用的时候一般使用mq handle,在这里包装一层成为mqhmqh = shm_mq_attach(mq, seg, NULL);// 配置libpq的消息队列为mqhpq_redirect_to_shm_mq(seg, mqh);// 记录父进程的pid为leader pidpq_set_parallel_leader(fps->parallel_leader_pid,fps->parallel_leader_backend_id);
}

配置好后,子进程已经记录了父进程的pid,在子进程中需要发送消息时:

int
mq_putmessage(...)
{...for (;;){// 先把书库放入mq中,flush到共享内存result = shm_mq_sendv(pq_mq_handle, iov, 2, true, true);if (pq_mq_parallel_leader_pid != 0){... // 这里只有子进程能走进来,通知父进程读取SendProcSignal(pq_mq_parallel_leader_pid,PROCSIG_PARALLEL_MESSAGE,pq_mq_parallel_leader_backend_id);...}}if (result != SHM_MQ_WOULD_BLOCK)break;(void) WaitLatch(MyLatch, WL_LATCH_SET | WL_EXIT_ON_PM_DEATH, 0,WAIT_EVENT_MESSAGE_QUEUE_PUT_MESSAGE);ResetLatch(MyLatch);CHECK_FOR_INTERRUPTS();}...
}

子进程发完了,信息会留存在mq中。然后给父进程发信号。

父进程收到kill过来的信号,进入信号处理函数(函数已经绑定sigusr1了),标记ParallelMessagePending

void
procsignal_sigusr1_handler(SIGNAL_ARGS)
{...if (CheckProcSignal(PROCSIG_PARALLEL_MESSAGE))HandleParallelMessageInterrupt();...
}void
HandleParallelMessageInterrupt(void)
{InterruptPending = true;ParallelMessagePending = true;SetLatch(MyLatch);
}

等下次调用CHECK_FOR_INTERRUPTS宏,执行ProcessInterrupts时处理具体的消息。

在函数中会shm_mq_receive接受子进程发到mq中的消息。

void
HandleParallelMessages(void)
{...HOLD_INTERRUPTS();......res = shm_mq_receive(pcxt->worker[i].error_mqh, &nbytes,&data, true);if (res == SHM_MQ_WOULD_BLOCK)break;else if (res == SHM_MQ_SUCCESS){StringInfoData msg;initStringInfo(&msg);appendBinaryStringInfo(&msg, data, nbytes);HandleParallelMessage(pcxt, i, &msg);pfree(msg.data);}elseereport(ERROR,(errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE),errmsg("lost connection to parallel worker")));}}}...RESUME_INTERRUPTS();
}

elog如何发送错误日志?

无论是并发时的子进程,还是普通进程,调用elog发送日志都会经过两步:

errstart(int elevel, const char *domain)    做一些初始化和配置errfinish(const char *filename, int lineno, const char *funcname)调用EmitErrorReport发送错误

EmitErrorReport负责将日志发送到client和server log

EmitErrorReport/* Send to server log, if enabled */if (edata->output_to_server)send_message_to_server_log(edata);/* Send to client, if enabled */if (edata->output_to_client)send_message_to_frontend(edata);

这里发送到client的日志send_message_to_frontend中,会走libpq的逻辑:

  1. 普通场景libpq使用PqCommSocketMethods的实现,将日志发送给客户端。
  2. 并发场景子进程中libpq使用PqCommMqMethods的实现,将日志发送给父进程。

这篇关于Postgresql源码(117)libpq的两套实现(socket/shm_mq)的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/512243

相关文章

hdu1043(八数码问题,广搜 + hash(实现状态压缩) )

利用康拓展开将一个排列映射成一个自然数,然后就变成了普通的广搜题。 #include<iostream>#include<algorithm>#include<string>#include<stack>#include<queue>#include<map>#include<stdio.h>#include<stdlib.h>#include<ctype.h>#inclu

JAVA智听未来一站式有声阅读平台听书系统小程序源码

智听未来,一站式有声阅读平台听书系统 🌟&nbsp;开篇:遇见未来,从“智听”开始 在这个快节奏的时代,你是否渴望在忙碌的间隙,找到一片属于自己的宁静角落?是否梦想着能随时随地,沉浸在知识的海洋,或是故事的奇幻世界里?今天,就让我带你一起探索“智听未来”——这一站式有声阅读平台听书系统,它正悄悄改变着我们的阅读方式,让未来触手可及! 📚&nbsp;第一站:海量资源,应有尽有 走进“智听

【C++】_list常用方法解析及模拟实现

相信自己的力量,只要对自己始终保持信心,尽自己最大努力去完成任何事,就算事情最终结果是失败了,努力了也不留遗憾。💓💓💓 目录   ✨说在前面 🍋知识点一:什么是list? •🌰1.list的定义 •🌰2.list的基本特性 •🌰3.常用接口介绍 🍋知识点二:list常用接口 •🌰1.默认成员函数 🔥构造函数(⭐) 🔥析构函数 •🌰2.list对象

【Prometheus】PromQL向量匹配实现不同标签的向量数据进行运算

✨✨ 欢迎大家来到景天科技苑✨✨ 🎈🎈 养成好习惯,先赞后看哦~🎈🎈 🏆 作者简介:景天科技苑 🏆《头衔》:大厂架构师,华为云开发者社区专家博主,阿里云开发者社区专家博主,CSDN全栈领域优质创作者,掘金优秀博主,51CTO博客专家等。 🏆《博客》:Python全栈,前后端开发,小程序开发,人工智能,js逆向,App逆向,网络系统安全,数据分析,Django,fastapi

让树莓派智能语音助手实现定时提醒功能

最初的时候是想直接在rasa 的chatbot上实现,因为rasa本身是带有remindschedule模块的。不过经过一番折腾后,忽然发现,chatbot上实现的定时,语音助手不一定会有响应。因为,我目前语音助手的代码设置了长时间无应答会结束对话,这样一来,chatbot定时提醒的触发就不会被语音助手获悉。那怎么让语音助手也具有定时提醒功能呢? 我最后选择的方法是用threading.Time

Android实现任意版本设置默认的锁屏壁纸和桌面壁纸(两张壁纸可不一致)

客户有些需求需要设置默认壁纸和锁屏壁纸  在默认情况下 这两个壁纸是相同的  如果需要默认的锁屏壁纸和桌面壁纸不一样 需要额外修改 Android13实现 替换默认桌面壁纸: 将图片文件替换frameworks/base/core/res/res/drawable-nodpi/default_wallpaper.*  (注意不能是bmp格式) 替换默认锁屏壁纸: 将图片资源放入vendo

C#实战|大乐透选号器[6]:实现实时显示已选择的红蓝球数量

哈喽,你好啊,我是雷工。 关于大乐透选号器在前面已经记录了5篇笔记,这是第6篇; 接下来实现实时显示当前选中红球数量,蓝球数量; 以下为练习笔记。 01 效果演示 当选择和取消选择红球或蓝球时,在对应的位置显示实时已选择的红球、蓝球的数量; 02 标签名称 分别设置Label标签名称为:lblRedCount、lblBlueCount

Kubernetes PodSecurityPolicy:PSP能实现的5种主要安全策略

Kubernetes PodSecurityPolicy:PSP能实现的5种主要安全策略 1. 特权模式限制2. 宿主机资源隔离3. 用户和组管理4. 权限提升控制5. SELinux配置 💖The Begin💖点点关注,收藏不迷路💖 Kubernetes的PodSecurityPolicy(PSP)是一个关键的安全特性,它在Pod创建之前实施安全策略,确保P

Java ArrayList扩容机制 (源码解读)

结论:初始长度为10,若所需长度小于1.5倍原长度,则按照1.5倍扩容。若不够用则按照所需长度扩容。 一. 明确类内部重要变量含义         1:数组默认长度         2:这是一个共享的空数组实例,用于明确创建长度为0时的ArrayList ,比如通过 new ArrayList<>(0),ArrayList 内部的数组 elementData 会指向这个 EMPTY_EL

如何在Visual Studio中调试.NET源码

今天偶然在看别人代码时,发现在他的代码里使用了Any判断List<T>是否为空。 我一般的做法是先判断是否为null,再判断Count。 看了一下Count的源码如下: 1 [__DynamicallyInvokable]2 public int Count3 {4 [__DynamicallyInvokable]5 get