当订单状态发生变更,库存系统、通知服务、数据仓库都希望第一时间拿到最新数据。传统做法是在业务代码里手动发送事件,但这种方式高度依赖开发者是否记得埋点,也容易在历史遗留项目中产生遗漏。SQL触发器把捕获逻辑下沉到数据库层,只要目标表出现INSERT、UPDATE或DELETE,就能自动记录变更,不再依赖应用层是否配合。配合消息队列之后,触发器先将变更写入本地outbox表,生产者异步投递到Kafka或RabbitMQ,不同技术栈的消费者各自订阅并同步到目标存储。这套方案既降低了源库与目标系统的耦合,也具备较好的扩展性。接下来从触发器设计、队列集成和可靠性保障三方面展开。

SQL触发器如何捕获数据变更
触发器是与表绑定的数据库对象,可以在INSERT、UPDATE、DELETE执行之后自动运行。MySQL、PostgreSQL、SQL Server都支持AFTER触发器,SQLite也支持类似机制。触发器的最大价值在于它运行在事务内部,因此把变更写入outbox表时,如果原事务回滚,outbox记录也会被回滚,避免产生孤儿消息。这一点是很多外部轮询方案难以保证的。
以MySQL为例,下面的代码为orders表创建一个AFTER INSERT触发器,将新插入行的主键、操作类型和关键字段打包成JSON,写入cdc_outbox表。outbox表的created_at可用于后续投递顺序判断。
DELIMITER $$
CREATE TRIGGER trg_orders_after_insert
AFTER INSERT ON orders
FOR EACH ROW
BEGIN
INSERT INTO cdc_outbox (table_name, row_id, operation, payload, created_at)
VALUES ('orders', NEW.id, 'INSERT', JSON_OBJECT('id', NEW.id, 'status', NEW.status), NOW());
END$$
DELIMITER ;
UPDATE和DELETE也需要类似的触发器,但它们访问旧值和删除行的方式不同。PostgreSQL中UPDATE触发器可以同时读取OLD和NEW,SQL Server中则通过inserted和deleted两张临时表获取数据。无论哪种数据库,触发器中建议只做轻量级操作,例如拼装JSON并插入outbox表,避免在触发器中直接调用HTTP接口或发送消息。直接发送网络请求会让业务事务被外部延迟拖累,数据库连接也会被长时间占用。
还需要注意,触发器的读写开销与表的写入频率成正比。对于每秒数千行的批量插入,行级触发器的高频执行可能成为瓶颈。因此该方案更适合中低写入量的核心业务表,或作为变更捕获的补充手段。高吞吐场景可以考虑基于binlog或WAL的CDC工具。
触发器与消息队列的集成链路
触发器只负责把变更写入outbox表,真正的消息发布由独立的投递进程完成。这个进程可以是一个常驻的Python、Go或Java服务,也可以使用数据库自身的异步通知机制,例如PostgreSQL的LISTEN/NOTIFY。常驻进程通常每隔几百毫秒查询一次未发布的outbox记录,读取一定批量后投递到消息队列,再更新这些记录的已发布状态。
下面的Python示例展示如何从PostgreSQL读取outbox记录并发送到Kafka。为了降低数据库压力,代码每次只读取最近未发布的一批记录,并且发送完成后再统一标记。
import psycopg2
from kafka import KafkaProducer
conn = psycopg2.connect(dbname='shop', user='sync', password='secret')
cur = conn.cursor()
producer = KafkaProducer(bootstrap_servers='localhost:9092')
cur.execute("""
SELECT id, table_name, row_id, operation, payload
FROM cdc_outbox
WHERE published = false
ORDER BY id
LIMIT 100
""")
rows = cur.fetchall()
for row in rows:
msg_id, table, row_id, op, payload = row
topic = 'dbsync.' + table
key = str(row_id).encode('utf-8')
value = str(payload).encode('utf-8')
producer.send(topic, key=key, value=value)
for row in rows:
cur.execute('UPDATE cdc_outbox SET published = true WHERE id = %s', (row[0],))
conn.commit()
投递过程中需要重点关注顺序。如果同一个主键的多次变更被分发到不同分区,消费者就可能先处理新版本再处理旧版本,导致目标端数据倒退。解决方式有两种:对于Kafka,可以把key设为主键,确保同一主键进入同一分区;对于RabbitMQ,可以使用一致性哈希交换机,或者把变更路由到同一个队列。消费端还需要保存每条消息的版本号,只有当收到的版本高于本地版本时才覆盖写入。
目标平台侧可以是MySQL、PostgreSQL、Redis、Elasticsearch或对象存储。消费者在收到消息后完成反序列化,判断操作类型,再执行相应的INSERT、UPDATE或DELETE逻辑。不同平台的写入API差异很大,但核心结构基本一致:先幂等校验,再落库,最后更新偏移量或确认消息。
可靠性与一致性如何保障
触发器加outbox方案的最大优势是捕获阶段的事务一致性:业务写入orders表与outbox写入发生在同一个数据库事务中,不会出现业务已提交但变更未记录的尴尬情况。真正容易出问题的是投递阶段,因为读取outbox、发送消息、标记published这三个动作并不是原子的。如果发送消息成功,但程序在标记published之前崩溃,下次重启后同一条消息会被再次发送,因此消费端必须实现幂等。
幂等键通常使用源表名、主键和版本号拼接而成。例如目标表可以增加一个source_version列,或者使用独立的同步状态表记录每个源行已应用的最大版本。消费者收到消息后先根据主键查询当前版本,如果消息版本小于等于已应用版本,直接跳过。这样即使重复投递、重试或消息重复,也不会造成重复写入或数据回退。
CREATE TABLE sync_state (
source_table VARCHAR(64) NOT NULL,
source_id BIGINT NOT NULL,
applied_version BIGINT NOT NULL,
updated_at DATETIME NOT NULL,
PRIMARY KEY (source_table, source_id)
);
INSERT INTO sync_state (source_table, source_id, applied_version, updated_at)
VALUES ('orders', 1001, 7, NOW())
ON DUPLICATE KEY UPDATE
applied_version = GREATEST(applied_version, VALUES(applied_version)),
updated_at = NOW();
除了幂等,失败处理同样关键。生产者如果暂时无法连接Kafka或RabbitMQ,可以保留outbox记录并持续重试,避免数据丢失。消费者处理失败时,不应该无限阻塞当前队列。可以将消息投递到死信队列,并记录失败原因,由运维或定时任务修复后重新入队。监控上建议关注outbox积压数量、投递延迟、消费失败率三个指标。
顺序性在前文已经提到,还需要注意时钟问题。如果源库和应用服务器分布在多个时区,created_at和版本号应尽量使用数据库单调递增的序列或自增主键,而不是应用服务器时间。分布式环境下,单纯依赖时间戳判断先后顺序很容易因时钟偏移产生误判。
常见误区与替代方案对比
一个常见误区是认为触发器可以直接在数据库里发布消息。某些数据库提供了外部插件,例如MySQL的sys_exec或PostgreSQL的扩展语言,可以在触发器中调用外部命令。但这种方式把网络IO引入事务上下文,不仅拖慢业务写入,还会在消息系统不可用时让事务失败,产生业务不可用的连锁反应。正确做法始终是先写outbox表,再异步投递。
另一个误区是过度依赖触发器而忽略表结构维护。当源表新增字段、重命名列或调整索引时,触发器里的JSON拼装逻辑需要同步修改。如果维护不当,消息中可能缺少关键字段,消费者侧的校验逻辑也会产生误报。建议在CI流程中加入触发器脚本的版本管理,每个源表的触发器与建表脚本一起评审和发布。
对于写入量极大或表结构频繁变化的系统,基于binlog或WAL的CDC工具可能是更好的选择。Debezium、Flink CDC、Maxwell等组件可以直接解析数据库日志,不需要在每张业务表上手工创建触发器。它们对源库的影响更小,也能捕获所有表的变化。但代价是部署复杂度更高,需要管理连接器、偏移量和Schema变更。SQL触发器方案更适合中小型项目、少量核心表或者需要快速落地的场景。
综合来看,SQL触发器与消息队列的组合在低侵入、事务一致性和实现简单之间取得了平衡。只要把outbox投递、幂等消费和死信处理设计清楚,就能搭建一条稳定可靠的数据同步通道,为跨平台实时分析、缓存刷新、索引更新等场景提供统一的变更事件流。