MySQL触发器无法直接发消息到Kafka或RabbitMQ,必须通过审计缓冲表中转,由外部消费者读取后投递;该表需用InnoDB、无主键无索引、结构极简,并妥善处理NULL、大字段和时间格式。

触发器里不能直接发消息到Kafka或RabbitMQ
MySQL触发器本身不支持网络调用,无法直接把变更写入外部消息队列。常见误区是想在 BEFORE UPDATE 或 AFTER INSERT 里调用 curl 或驱动发消息——这会报错或被禁用(即使启用 sys_exec 也极不稳定,且多数生产环境禁用)。真正可行的路径是:触发器只负责把变更记录落库到一张轻量级审计缓冲表,再由外部消费者(如Python脚本、Debezium、Canal)轮询或监听该表,转投队列。
审计缓冲表设计要避开主键和索引陷阱
这张表不是用来查的,是给消费者“搬数据”的中转站,所以结构必须极简:
-
id用BIGINT UNSIGNED AUTO_INCREMENT,但别设主键——主键会引发插入竞争和自增锁,高并发下拖慢原表DML - 只保留必要字段:
table_name(VARCHAR(64))、operation(ENUM('INSERT','UPDATE','DELETE'))、pk_value(存主键值,TEXT或VARCHAR(255),兼容复合主键拼接)、old_data和new_data(JSON类型,MySQL 5.7+ 支持,比TEXT更安全) - 不建任何二级索引——消费者按
id顺序读取后立即DELETE或用pt-archiver归档,索引纯属冗余开销 - 引擎必须用
BLACKHOLE?不,那是误解;实际要用InnoDB,但autocommit=1+ 短事务,避免阻塞
AFTER 触发器里 JSON_OBJECT 要小心 NULL 和二进制字段
用 JSON_OBJECT() 构造 old_data/new_data 最方便,但有三个坑:
- 如果某列是
NULL,JSON_OBJECT('col', col)会把col键整个丢掉,导致审计字段缺失——得先用IFNULL(col, '##NULL##')填充占位符 - 含
BLOB、TEXT大字段的表,直接塞进JSON_OBJECT可能超max_allowed_packet,建议只审计关键业务字段,用显式列名拼装,而非SELECT JSON_OBJECT(...)全字段 - 时间字段(
DATETIME)在OLD/NEW中是字符串格式,但有时带微秒或有时不带,统一用DATE_FORMAT(OLD.updated_at, '%Y-%m-%d %H:%i:%s')格式化,避免下游解析失败
示例片段(对 users 表):
CREATE TRIGGER tr_users_audit_after_insert
AFTER INSERT ON users
FOR EACH ROW
INSERT INTO audit_queue (table_name, operation, pk_value, new_data)
VALUES (
'users',
'INSERT',
NEW.id,
JSON_OBJECT(
'id', NEW.id,
'name', IFNULL(NEW.name, '##NULL##'),
'email', IFNULL(NEW.email, '##NULL##'),
'created_at', DATE_FORMAT(NEW.created_at, '%Y-%m-%d %H:%i:%s')
)
);消费者必须处理重复和顺序问题
哪怕用 AFTER 触发器,也不能保证“先写审计表、再提交事务”绝对原子——MySQL崩溃时可能留下未消费的脏记录。所以消费者逻辑必须:
- 每次只拉固定条数(如
SELECT * FROM audit_queue ORDER BY id LIMIT 100),处理完再DELETE WHERE id <= ?,避免长事务锁表 - 发送到Kafka前,用
id做幂等键(key=id),靠Kafka端去重,而不是在MySQL里加唯一约束(会冲突) - 不依赖
id严格递增顺序:MySQL自增在主从切换、批量插入时可能跳变,下游应接受乱序,靠业务时间戳(如created_at)做最终排序 - 必须监控
audit_queue表的DATA_LENGTH增长速率,超过阈值(如 1GB)就告警——说明消费者挂了或队列积压
真正难的从来不是写触发器,而是让审计链路在主库压力大、网络抖动、消费者重启时仍不丢不重。缓冲表只是起点,后面每一步都得按分布式系统来想。


















