Windows环境下如何在PostgreSQL函数与触发器中向RabbitMQ发送消息
Windows环境下PostgreSQL函数/触发器向RabbitMQ发送消息实操方案
需提前安装的额外工具
你可以根据自己的需求二选一方案,对应准备以下工具:
方案1(推荐,轻量无额外依赖)
- 与你PostgreSQL版本、系统架构完全匹配的
pg_amqp预编译插件(PG官方的AMQP协议对接插件,无需额外 runtime) - 可正常访问的RabbitMQ服务(提前配置好vhost、账号密码、交换机/队列的访问权限)
方案2(灵活性高,适合自定义逻辑)
- PostgreSQL对应版本的PL/Python3语言扩展
- 与PostgreSQL版本适配的Python运行环境 + Python AMQP客户端库
pika - 可正常访问的RabbitMQ服务(提前配置好vhost、账号密码、交换机/队列的访问权限)
操作步骤
方案1:pg_amqp插件方案
- 下载和你PostgreSQL大版本、系统位数完全匹配的
pg_amqpWindows预编译包,禁止跨版本使用 - 解压插件包,将后缀为
.dll的文件放入PostgreSQL安装目录下的lib文件夹,将.control、.sql后缀的文件放入share/extension文件夹 - 打开psql命令行工具连接到目标数据库,执行命令启用插件:
CREATE EXTENSION amqp;
- 配置RabbitMQ连接信息:
-- 插入RabbitMQ服务配置,返回的id为后续调用的broker标识,示例返回id为1 INSERT INTO amqp.broker (host, port, vhost, username, password) VALUES ('127.0.0.1', 5672, '/', 'guest', 'guest');
- 编写发送消息的触发器函数并绑定触发器:
-- 编写消息发送函数,示例为将新增行转成JSON发送 CREATE OR REPLACE FUNCTION send_rabbitmq_trigger_func() RETURNS trigger AS $$ DECLARE msg_content text; BEGIN msg_content = row_to_json(NEW)::text; -- 参数依次为:broker id、交换机名称、路由键、消息内容 PERFORM amqp.publish(1, 'test_exchange', 'test_routing_key', msg_content); RETURN NEW; END; $$ LANGUAGE plpgsql; -- 绑定触发器,示例为test_table表插入数据后触发消息发送 CREATE TRIGGER after_insert_send_msg AFTER INSERT ON test_table FOR EACH ROW EXECUTE FUNCTION send_rabbitmq_trigger_func();
方案2:PL/Python方案
如果找不到对应版本的pg_amqp编译包,或者需要自定义消息序列化、重试等逻辑,可以用该方案:
- 安装和你PostgreSQL版本匹配的Python运行环境,比如PostgreSQL 14对应Python 3.9,版本不匹配会导致PL/Python加载失败
- 打开psql执行命令启用PL/Python扩展:
CREATE EXTENSION plpython3u;
- 给PG对应的Python环境安装pika库:
pip install pika
- 编写Python消息发送函数:
CREATE OR REPLACE FUNCTION send_rabbitmq_msg(msg text) RETURNS void AS $$ import pika # 配置RabbitMQ连接 credentials = pika.PlainCredentials('guest', 'guest') connection = pika.BlockingConnection(pika.ConnectionParameters('127.0.0.1', 5672, '/', credentials)) channel = connection.channel() # 发送消息,交换机、路由键根据你的实际配置修改 channel.basic_publish(exchange='test_exchange', routing_key='test_routing_key', body=msg) connection.close() $$ LANGUAGE plpython3u;
- 编写触发器函数调用上述Python函数即可,逻辑和方案1的触发器函数一致。
注意事项
- 生产环境禁止在函数内硬编码RabbitMQ账号密码,可以存储在加密配置表或者系统环境变量中读取
- 建议添加异常捕获逻辑,避免RabbitMQ故障导致PG事务回滚,影响正常业务写入
pg_amqp默认会将消息发送纳入当前PG事务,事务回滚时消息也不会发出,如果不需要该特性可以用PL/Python方案自行控制事务边界
内容的提问来源于stack exchange,提问作者ZedZip
相关产品推荐
相关产品推荐

