
简介《数据库进程间通信解决方案之MQ_2.docx》针对数据库在数据变更时难以主动通知其他进程的痛点提供了一套基于消息队列的实践方案适合需要在不动业务代码的前提下完成外部联动的开发者、DBA与架构师。文档从传统定时轮询的缺陷讲起进而介绍FIFO管道、ZeroMQ等IPC/RPC机制并重点阐述MySQL插件的开发与使用包括UDF函数zmq_client的编译装载、触发器的创建以及发送短信/邮件、图片处理、身份证校验、网站静态化等典型调用场景。通过文档可掌握如何将数据库事件实时推送至MQ服务端再交由下游程序异步处理或同步应答从而降低系统耦合、提升响应速度。资源包仅含1个docx文件约116KB内容紧凑目前已有97人浏览学习适合对数据库与消息队列集成感兴趣的技术人员参考。1. 数据库进程间通信轮询之外的另一个答案数据库进程间通信最常见的实现是轮询程序每隔几秒去 SELECT 一次看字段有没有变化。数据量小还能忍量一上来轮询的延迟、无效查询、数据库压力就叠在一起变成一笔沉默成本。这份资源给的是另一条路——把 MySQL 变成消息发送端通过 UDF 插件调用 ZeroMQ把数据变化实时推给外部进程。它解决的场景很具体数据库里数据变了外部程序能不能立刻知道并干活同时不用改现有业务代码、不引入 CDC 组件只需要在库里建几个触发器和函数就能把短信发送、图片处理、页面静态化这些活拆出去异步做。对手里有存量项目、外包已交付项目、想在不动代码的前提下加实时通知能力的人来说这条路线值得照着走一遍。2. MQ 方案怎么立住IPC、RPC 与三种协议前缀的选型2.1 为什么轮询不是好答案延迟与无效查询的双重浪费轮询的本质是“猜”你猜数据大概什么时候会变就把查询周期设成那个粒度。猜短了数据库频繁被无效 SELECT 打扰猜长了业务方看到的数据总比别人慢半拍。比如商品价格改了静态化页面要等下一个轮询周期才更新这段时间里用户看到的就是旧价格。这是轮询在实时性上的天花板不是把 sleep 时间调小就能解决的因为你永远不知道下一次变化什么时候来。事件通知则是反过来的思路变化发生时主动喊一声而不是让外部程序反复确认。MySQL 里能捕获变化的地方有三处触发器、EVENT 定时任务、事务内部的语句流。这份资源的核心做法就是把 UDF 函数放到这些位置里让数据库在数据变化的同时把变化内容作为消息推出去。数据库仍然是业务的数据库但多了一张“嘴”可以对外说话。这里的几个缩写需要先对齐IPC 是进程间通信Inter-Process Communication指同一台机器内两个进程交换数据ITC 是线程间通信InterThread Communication发生在同一个进程内部RPC 是远程过程调用Remote Procedure Calls通信双方可以跨操作系统、跨服务器。这份资源里的插件同时覆盖了这三层本机用 ipc 前缀同进程内用 inproc 前缀跨机器用 tcp 前缀后面会逐个讲。2.2 从 FIFO 到 ZeroMQ管道方案的边界在哪作者早先写过一版基于 FIFO命名管道的插件思路是在 MySQL 里写管道文件外部程序从管道里读数据。FIFO 的优势是零依赖系统自带不用装任何库但它的短板非常明显管道是单向的、无缓冲区的写端和读端必须同时在线否则写操作会阻塞消息也可能直接丢失。更关键的是它只适合同一台服务器内的进程通信数据库和应用分机部署后就玩不转了。ZeroMQ 解决的是这层问题它把 socket 抽象成消息队列自带缓冲、自动重连、多协议支持、多语言绑定。数据库这边只要把消息交给 ZeroMQ剩下的路由、排队、传输都由库完成不关心对端是 C 程序、Python 程序还是另一个数据库。相比裸 socketZeroMQ 不需要你手动处理粘包拆包相比 FIFO它可以跨机器相比 RabbitMQ 这类重型消息中间件它不需要独立的 broker 进程嵌入在插件里就能跑。两者的适用范围可以这样看通信载体通信范围缓冲区对端必须在线依赖FIFO 管道单机进程间无是否则阻塞无ZeroMQ ipc单机进程间有否自动排队libzmqZeroMQ tcp跨机器有否可重连libzmq插件最终以 UDF 形式提供编译产物是 libzeromq.so里面封装了多个函数。文档里明确列出的有 zmq_client 和 zmq_publish前者用于“请求-应答”式的发送后者面向发布-订阅场景剩余函数以源码里的注释和头文件为准。对大多数业务场景先用 zmq_client 就够了。2.3 协议前缀怎么选inproc、ipc、tcp 的适用边界插件支持三种协议前缀对应三类通信距离地址写法分别是inproc://my_publisher线程间通信发生在同一个进程内部适合在程序里同时加载发送端和接收端的场景。ipc:///tmp/feeds/0本机进程间通信底层走 UNIX socket 文件没有 TCP 协议栈的开销延迟最低。同机部署时优先用它。tcp://server001:5555跨服务器通信数据库在 A 机、消费程序在 B 机或者数据库是多实例共享一个 MQ 集群时使用。这里有一个方向性问题容易搞反插件是客户端负责发起连接外部处理程序是服务端负责 bind 端口等待连接。文件里测试服务器的启动方式是先跑./server再在 MySQL 里执行select zmq_client(...)顺序不能颠倒否则客户端会连接失败并返回 false。地址里的 5555 是服务端 bind 的端口客户端只写地址不写端口选择连接由 ZeroMQ 内部完成。注意tcp 模式下服务端 bind 在 0.0.0.0 还是某个具体 IP决定了客户端能连哪里。生产环境建议 bind 内网 IP不要裸奔到公网。3. 编译装载 UDF依赖安装、cmake 与 create function 全流程3.1 依赖包装全才能不出怪错编译 UDF 需要三样东西编译器、MySQL 客户端开发库、构建工具。在 Ubuntu/Debian 系系统上对应这几包sudo apt-get install pkg-config sudo apt-get install libmysqlclient-dev sudo apt-get install gcc g make cmakepkg-config 用来探测库的头文件和链接路径libmysqlclient-dev 提供 MySQL 客户端 API 的头文件与动态库gcc/g 负责编译 C/C 代码make 和 cmake 负责构建流程。少了 pkg-config 或者 dev 包缺失cmake 阶段经常报“找不到 MySQL library”而不是在编译阶段报错排查起来反而费时间。注意libmysqlclient-dev 的版本必须和正在运行的 MySQL 服务端大版本一致。比如服务端是 MySQL 8.0编译机却装了 5.7 的 dev 包编出来的 so 在装载时容易出符号错误。源码来自 GitHub先把它拿到本地git clone https://github.com/netkiller/mysql-zmq-plugin.git cd mysql-zmq-plugin进入目录后先扫一眼 README 和 CMakeLists.txt确认这个版本对应的 MySQL 版本区间。原作者文档写得比较早如果你用的是 MySQL 8.0 之后的版本建议先确认 CMakeLists.txt 里有没有针对新版本 ABI 的兼容处理。3.2 cmake 编译与 so 文件落位编译命令很直接三步走cmake . make sudo make installcmake . 表示在当前目录读取 CMakeLists.txt 并生成 Makefile默认不传参数时会自动探测 MySQL 的头文件和库文件位置。make 开始实际编译产物是动态库文件名字一般是 libzeromq.so。make install 会把 so 文件安装到 MySQL 的插件目录但如果你用源码方式安装过 MySQL或者系统里存在多个 MySQL 实例安装位置可能不是你当前实例的 plugin_dir这一步需要人工确认。确认插件目录的方法是查 MySQL 变量SHOW VARIABLES LIKE plugin_dir;注意路径不对的话就算 create function 成功实际调用函数时也会报 “cant open shared library”。这时候手动把 so 复制过去然后重启一次 MySQL 让它重新加载依赖。复制命令示例sudo cp libzeromq.so /usr/lib/mysql/plugin/如果 make install 之后 MySQL 仍然加载不了 libzeromq.so先检查动态库依赖是否完整ldd /usr/lib/mysql/plugin/libzeromq.soldd 输出里如果出现 libzmq.so 找不到的提示说明系统里缺 libzmq需要先安装 ZeroMQ 的基础库否则插件连装载都过不去。这一步是排查“函数创建成功但 select 报错”的最高频根因。3.3 create function 装载与安装验证so 文件就位后在 MySQL 里注册函数。需要注意UDF 的装载是把函数名和 so 文件绑定写入系统的 func 表这一步需要具备 mysql.func 的插入权限通常用 root 或专门的管理账号执行CREATE FUNCTION zmq_client RETURNS STRING SONAME libzeromq.so; CREATE FUNCTION zmq_publish RETURNS STRING SONAME libzeromq.so;RETURNS STRING 声明函数的返回类型是字符串SONAME 后面跟的是动态库文件名不需要写完整路径MySQL 会到 plugin_dir 下找。同一文件名可以被多个函数共用所以这两个函数用的是同一个 so。如果装错了想重来卸载命令是对称的DROP FUNCTION zmq_client; DROP FUNCTION zmq_publish;验证是否装载成功查 func 表SELECT * FROM mysql.func WHERE name LIKE zmq%;查询结果里能看到函数名、返回类型、关联的 so 文件名type 为 function。这一步能过说明插件已经被 MySQL 认识接下来才是真正的通信测试。4. 把 MQ 接进业务触发器、批量 SELECT 与同步异步取舍4.1 触发器三件套INSERT、UPDATE、DELETE 都要通知以文档里的静态化场景为例。假设你接手了一个外包电商项目外包团队早走了现在要给商品页做静态化又不敢在十几处商品操作代码里逐个加逻辑。最保险的做法是用触发器商品表只要有变化就把记录的主键推给 MQ静态化程序收到后重新下载生成页面。DELIMITER // CREATE TRIGGER demo_after_insert AFTER INSERT ON demo FOR EACH ROW BEGIN SELECT zmq_client(tcp://localhost:5555, NEW.id); END// CREATE TRIGGER demo_after_update AFTER UPDATE ON demo FOR EACH ROW BEGIN SELECT zmq_client(tcp://localhost:5555, NEW.id); END// CREATE TRIGGER demo_after_delete AFTER DELETE ON demo FOR EACH ROW BEGIN SELECT zmq_client(tcp://localhost:5555, NEW.id); END// DELIMITER ;逻辑说明三个触发器分别覆盖插入、更新、删除NEW.id 是当前操作行的主键值MySQL UDF 的调用必须用 SELECT 形式所以触发器主体里写的是SELECT zmq_client(...)而不是直接函数调用。服务端程序拿到 id 后再按需回查数据库就能区分这是新建、覆盖还是删除动作。这里有个设计要点消息体里只传主键 id不传整行数据。原因有两层一是消息体越短网络和队列压力越小二是业务字段随时可能变化传过去的数据可能到消费端时已经过期拿到 id 再回查才是准的。静态化程序的行为对应关系是插入时生成 html更新时覆盖 html删除时删掉 html 文件。4.2 批量外呼场景一条 SELECT 把整个结果集送进队列触发器的局限是它只对单行操作有意义如果有一批历史数据要补发或者每个月要给所有订阅用户推一次短信触发器就不适用了。这时候可以直接跑查询把结果集一行一行打给 MQSELECT zmq_client(tcp://localhost:5555, mobile) FROM demo WHERE subscribed Y;这条 SQL 的逻辑是先从 demo 表筛出所有 subscribed 为 Y 的用户手机号然后对结果集逐行调用 zmq_client把手机号推给 MQ 服务端。执行完成后服务端的消息队列里就有这批待发送的号码消费程序可以从队列里多线程领取并发短信。如果你要传递的不止一个字段可以用 concat 把多个字段拼进一个消息体分隔符自己定SELECT zmq_client(tcp://localhost:5555, CONCAT(name, ,, mobile, ,news)) FROM demo; SELECT zmq_client(tcp://localhost:5555, CONCAT(name, |, mobile, |news)) FROM demo;参数说明第一个逗号前的 tcp 地址是 MQ 服务端地址第二个参数是消息体。消息体的格式完全由服务端解析程序决定用逗号还是竖线分隔两端约定好就行。我建议固定一种分隔符并在服务端做容错否则后来改字段顺序时容易踩坑。4.3 拼装 JSON 与多字段消息体concat 的写法边界把消息体组装成 JSON是另一种常见做法因为服务端解析 JSON 比手工 split 分隔符更不容易出错SELECT zmq_client(tcp://localhost:5555, CONCAT({name:, name, , tel:, mobile, , template:news})) FROM demo;执行后客户端这边能看到类似{name:neo, tel:13113668891} OK的结果OK 是服务端收到消息后的确认回执。需要清楚的是这里的 JSON 是用 CONCAT 手工拼出来的字符串不是结构化数据。它只是在文本层面长成了 JSON 的样子服务端拿到后按 JSON 解析仅此而已。这个写法有两个边界要留意。第一字段值里如果本身带逗号、冒号或花括号拼出来的 JSON 就是非法的比如姓名叫 “neo, jr” 这类带逗号的输入服务端解析会断到中间。第二手工拼 JSON 不做任何转义字段里的双引号和反斜杠会原样进入消息体。如果你的字段内容是用户可输入的更稳妥的做法是触发器里只传主键 id服务端拿到 id 后再回查数据库取字段这样绕开了所有转义问题。提示不要指望在 SQL 里用 JSON_OBJECT 来解决 UDF 的消息体格式问题。插件是按字符串收发的JSON_OBJECT 拼出来的还是字符串特殊字符该转义还是要转义。4.4 同步与异步两条路线怎么选MQ 通知不一定都是异步的文档里明确区分了两种模式。发短信、发邮件、处理图片这类任务是盲目发送发出去之后能不能达、用户看不看数据库这边根本无法确认所以 MQ 端收到任务后立即回一个“成功”确认消息进队列就算完事剩下的交给消费端慢慢做。这属于异步。身份证号校验这类任务恰恰相反它要求数据库这边拿到校验结果才能决定下一步走哪个分支。MQ 端的处理程序几乎不耗时能立即返回结果所以可以用同步方式客户端 select 的返回值里直接带着校验结论。同步等待的代价是当 MQ 服务端繁忙时 select 会被拖住所以要确认服务端的处理链路足够快再启用同步模式。判断维度建议走异步建议走同步任务结果是否需要回库不需要需要后续决策依赖结果对端处理耗时秒级以上毫秒级对端是否可能积压可能需要队列缓冲不积压典型场景短信、邮件、图片处理身份证校验、黑名单查询我的习惯是默认全部异步只有业务上确实需要拿回执做后续判断的任务才改同步。异步方案在 MQ 端进程重启时还能靠队列缓冲扛一扛同步方案一旦 MQ 服务端抖动数据库这边的写入链路会跟着遭殃。5. 部署避坑插件不加载、端口握不上、消息丢失的五条记录5.1 插件装上却调不到Function does not exist现象CREATE FUNCTION 执行成功mysql.func 表里也能查到记录但真正执行SELECT zmq_client(...)时报错提示函数不存在。原因最常见的是 so 文件不在当前实例的 plugin_dir 下。MySQL 加载 UDF 时只认 plugin_dir 目录里的文件你把 so 复制到了别的路径或者系统里有多个 MySQL 实例、make install 把文件装到了另一个实例的目录里都会出现“注册成功但调用失败”。解决先用SHOW VARIABLES LIKE plugin_dir;确认当前实例的插件目录再ls -l查看该目录下有没有 libzeromq.so。没有就复制过去然后执行DROP FUNCTION zmq_client;再重新 CREATE FUNCTION。如果 so 文件在但调用仍失败执行ldd /usr/lib/mysql/plugin/libzeromq.so检查动态库依赖libzmq.so 缺失时插件在调用瞬间会崩掉或报错。5.2 消息发不出去服务端一直没收到现象SELECT 正常执行返回的却是 false服务端程序没有任何输出或者返回看起来正常但服务端那边要等很久才收到。原因插件是 ZeroMQ 客户端外部程序是服务端二者通过 REQ/REP 模式匹配。服务端没启动或者启动时 bind 的地址和客户端 connect 的地址不一致连接就握不上。另外服务端如果 bind 在 localhost而客户端用的是别的机器地址同样连不上。解决先在数据库所在机器上执行netstat -lntp | grep 5555确认服务端确实在监听再用文档里自带的测试程序跑一遍最简单链路——先启动./server再执行SELECT zmq_client(tcp://localhost:5555,Hello world!);服务端能打印出 Received 说明链路通。跨机部署时服务端 bind 地址要改成内网 IP客户端也连对应 IP不要两边都写 localhost。同机部署我建议直接用ipc:///tmp/feeds/0少一层 TCP 协议栈也能少碰一类防火墙误拦问题。5.3 触发器拖慢写入同步阻塞与从库双写现象给表加上触发器后INSERT/UPDATE 的执行时间明显变长业务方反馈写库变慢更严重时从库也报触发器相关错误。原因如果 MQ 服务端处理慢或者网络抖动触发器里的 SELECT zmq_client 会一直等到服务端返回才结束这个等待时间直接算进 DML 事务里把正常的写操作拖住。另一个问题是从库别重复建触发器否则主从两边各发一次消息消费端会重复消费。解决触发器里只发消息不做任何重量级操作确保单次调用在毫秒级返回MQ 服务端采用异步处理收到消息立即回执再慢慢干活。排查时可以用SHOW PROCESSLIST;看有没有长时间的 SELECT zmq_client 卡在 State 列。从库上不要建触发器靠主库复制过去的变更已经能保证数据一致性不需要再从库发一遍消息。5.4 JSON 拼接翻车字段里的逗号和引号现象服务端收到消息后 JSON 解析失败或者一条消息被错误地解析成了两条排查发现消息体里字段值带了逗号、引号等特殊字符。原因CONCAT 手工拼 JSON 不转义。字段内容里有逗号时消息体变成了{name:neo, jr, tel:13113668891}服务端按逗号分列就错位了。这是手工拼接的固有问题不是插件的问题。解决三个办法按优先级选。第一触发器消息里只传主键 id服务端拿到 id 后自己回查数据库取字段彻底避开转义第二对可能含特殊字符的字段先做 REPLACE把逗号替换成全角或去掉再拼进去第三如果机器上 MySQL 版本足够新在服务端回查时用 JSON_OBJECT 直接生成合法 JSON而不是依赖 SQL 端手工拼。5.5 消息进了队列却丢了先写库后发 MQ 的补偿现象数据库里商品价格已经更新静态化页面却还是旧版查看消费端日志发现消息根本没到或者到了但消费端崩了没处理完。原因ZeroMQ 的消息队列默认是内存态的不持久化。MQ 服务端进程重启队列里还没消费的消息就丢了。数据库这边不会重发因为触发器只在数据变更瞬间触发一次。这就是“先写数据库、后发 MQ”模型的老问题——库落库了消息却没送达或没消费。解决给这类对一致性有要求的场景加一个补偿对账任务。利用表里的更新时间字段定期扫描最近十分钟内更新过的记录重新发一次 MQ 消息消费端按主键幂等处理重复消息不重复影响结果。我一般用 EVENT 定时任务来做CREATE EVENT resend_staticize ON SCHEDULE EVERY 5 MINUTE DO SELECT zmq_client(tcp://localhost:5555, id) FROM demo WHERE updated_at NOW() - INTERVAL 10 MINUTE;参数说明EVENT 每五分钟扫一次把最近十分钟内更新过的记录重新发给 MQ静态化程序收到重复消息时先判断页面文件是否已最新最新的直接跳过。这套补偿不能替代主链路它只兜底能容忍分钟级延迟的场景够用了。要求更高的场景需要考虑把对端换成带持久化的消息系统。6. 进阶技巧让 MQ 端返回结构化结果校验型请求走同步 RPCzmq_client 的返回值完全由服务端那边决定。文档里的 server 测试程序回的是 “xxx OK”你就可以在生产服务端里回结构化内容让客户端一个 SELECT 拿到校验结果不用落表、不用二次查询。这个技巧特别适合身份证校验、会员状态检查这类需要回执的任务。服务端用 Python 写一个 REP 模式的处理程序常见的做法是这样import zmq import json context zmq.Context() socket context.socket(zmq.REP) socket.bind(ipc:///tmp/feeds/0) while True: # 接收客户端发来的身份证号 id_number socket.recv().decode(utf-8) # 这里放真正的校验逻辑 valid len(id_number) 18 result json.dumps({valid: valid, received: id_number}) # 把结构化结果返回给数据库端 socket.send(result.encode(utf-8))逻辑说明socket.recv 拿到的是数据库那边传入的一行消息处理完把 JSON 字符串通过 socket.send 返回。zmq_client 会把这个返回值变成 SELECT 结果显示出来。客户端 SQL 对应写法SELECT zmq_client(ipc:///tmp/feeds/0, id_number) FROM demo;返回结果类似{valid: true, received: 110101...} OK。有了结构化返回你可以在上层程序里对这个字符串做 LIKE 判断也可以原样交给应用解析 JSON。注意返回值末尾的 OK 是服务端追加的确认标记解析前先按约定的分隔符切开。这个技巧的价值是让数据库具备了双向通信能力不再是单方面往外推消息。从那以后我每在一个新库上装完这套插件都会先跑一遍 test/server 验证链路通不通再专门测一次“服务端故意不启动时客户端返回什么”确认失败分支是可控的最后才把触发器挂到业务表上。顺序反了排查时就会把插件问题、网络问题、业务问题混在一起浪费半天时间。希望帮到你。本文还有配套的精品资源点击获取