sql server触发器必须用service broker实现事务一致的实时解耦,因其毫秒级本地消息传递不阻塞主事务;而http/kafka直调会引发超时、死锁或安全风险。

SQL Server 触发器不能直接推送消息到 Kafka、RabbitMQ 或 HTTP 接口——任何尝试在触发器里调用 sp_OACreate、xp_cmdshell、sp_send_dbmail 或外部 HTTP 客户端的操作,都会导致事务阻塞、超时甚至死锁。唯一安全、原生支持、事务一致的实时解耦方案,是使用 Service Broker。
为什么必须用 Service Broker 而不是直接发 HTTP?
触发器运行在主 DML 事务上下文中,所有操作必须轻量、快速、无外部依赖。一旦触发器里发起网络请求:
-
sp_send_dbmail会同步等待 SMTP 响应,DML 事务卡住数秒甚至分钟 -
xp_cmdshell调用curl属于高危禁用操作,多数生产环境已关闭,且违反最小权限原则 - 自定义 CLR 函数调用 Socket 或 HttpClient 同样破坏事务原子性,且需 UNSAFE ASSEMBLY 权限,审计难、升级易崩
- 哪怕只是
INSERT INTO remote_table(跨服务器链接),也会引入分布式事务(MSDTC),大幅增加锁争用和失败率
而 Service Broker 是 SQL Server 内置的消息系统,所有操作都在本地数据库完成:BEGIN DIALOG → SEND → COMMIT,毫秒级,不跨事务边界,失败不影响主表写入。
Service Broker 必须启用且配置四要素
缺一不可,否则 SEND 会报错 Service does not exist 或 Invalid contract:
- 启用数据库 Broker:
ALTER DATABASE YourDB SET ENABLE_BROKER WITH ROLLBACK IMMEDIATE(注意:不是NEW_BROKER,除非你确定要丢弃旧会话) - 定义消息类型:
CREATE MESSAGE TYPE [//YourApp/ChangeEvent] VALIDATION = WELL_FORMED_XML(推荐 XML,兼容性好;JSON 需 SQL Server 2016+ 且验证麻烦) - 创建约定:
CREATE CONTRACT [//YourApp/ChangeContract] ([//YourApp/ChangeEvent] SENT BY INITIATOR)(只允许发送方发,不需回复) - 建队列和服务:
CREATE QUEUE ChangeQueue; CREATE SERVICE [//YourApp/ChangeService] ON QUEUE ChangeQueue ([//YourApp/ChangeContract])
队列表名不要带 schema 前缀(如 dbo.ChangeQueue),CREATE SERVICE 语句中只写队列裸名。
触发器里只能做三件事:BEGIN DIALOG、SEND、INSERT(可选日志)
下面是一个典型变更捕获触发器示例,监听 Orders 表的 INSERT/UPDATE/DELETE:
CREATE TRIGGER tr_Orders_Change ON Orders
AFTER INSERT, UPDATE, DELETE
AS
BEGIN
DECLARE @h UNIQUEIDENTIFIER;
DECLARE @msg NVARCHAR(MAX);
<p>-- 只取关键字段,避免大对象或 LOB 拖慢触发器
SELECT @msg = (
SELECT
CASE
WHEN EXISTS(SELECT <em> FROM inserted) AND EXISTS(SELECT </em> FROM deleted) THEN 'UPDATE'
WHEN EXISTS(SELECT * FROM inserted) THEN 'INSERT'
ELSE 'DELETE'
END AS op,
ISNULL(i.OrderId, d.OrderId) AS id,
i.Status AS new_status,
d.Status AS old_status
FROM inserted i
FULL JOIN deleted d ON i.OrderId = d.OrderId
FOR JSON PATH, WITHOUT_ARRAY_WRAPPER
);</p><p>BEGIN DIALOG @h
FROM SERVICE [//YourApp/ChangeService]
TO SERVICE N'//YourApp/ChangeService'
ON CONTRACT [//YourApp/ChangeContract]
WITH ENCRYPTION = OFF;</p><p>SEND ON CONVERSATION @h MESSAGE TYPE [//YourApp/ChangeEvent] (@msg);</p><p>-- 可选:记录发送时间,用于监控积压
INSERT INTO ChangeLog (table_name, event_type, sent_time)
VALUES ('Orders', 'trigger_send', GETDATE());
END</p>
注意点:
-
BEGIN DIALOG的 TO SERVICE 必须和 FROM SERVICE 相同(单向队列场景),否则报错Target service does not exist - 不要在触发器里
WAITFOR (RECEIVE ...)——那是激活存储过程的事,放这里会阻塞 - 避免在 SELECT … FOR JSON 中引用大文本字段(如
description),JSON 序列化开销大,易拖慢主业务 - 如果触发器涉及多表 JOIN 或子查询,务必加
NOLOCK提示,防止锁升级
后续消费必须用激活存储过程,不能轮询
队列消息不会自动消失,必须由激活存储过程持续接收并转发。手动轮询(如定时作业查 SELECT TOP 1 FROM sys.transmission_queue)延迟高、资源浪费、易漏消息。
激活过程核心结构:
CREATE PROCEDURE usp_ConsumeChangeQueue
AS
BEGIN
DECLARE @h UNIQUEIDENTIFIER;
DECLARE @msgType SYSNAME;
DECLARE @msgBody VARBINARY(MAX);
<p>WHILE (1=1)
BEGIN
BEGIN TRY
WAITFOR (
RECEIVE TOP (1)
@h = conversation_handle,
@msgType = message_type_name,
@msgBody = message_body
FROM ChangeQueue
), TIMEOUT 5000;</p><pre class="brush:php;toolbar:false;"> IF (@@ROWCOUNT = 0) BREAK;
IF (@msgType = N'http://schemas.microsoft.com/SQL/ServiceBroker/EndDialog')
BEGIN
END CONVERSATION @h;
END
ELSE IF (@msgType = N'//YourApp/ChangeEvent')
BEGIN
-- 解析 @msgBody(XML 或 JSON),调用外部客户端(Kafka Producer / HttpClient)
-- 注意:此处才做网络调用,与主事务完全隔离
EXEC usp_SendToKafka @msgBody;
END CONVERSATION @h;
END
END TRY
BEGIN CATCH
-- 记录错误但不中断循环,避免队列堵塞
INSERT INTO ErrorLog VALUES (ERROR_MESSAGE(), GETDATE());
END CONVERSATION @h WITH ERROR = 1 DESCRIPTION = N'Processing failed';
END CATCHEND END
关键细节:
- 必须设为队列的 ACTIVATION 存储过程:
ALTER QUEUE ChangeQueue WITH ACTIVATION (STATUS = ON, PROCEDURE_NAME = usp_ConsumeChangeQueue, MAX_QUEUE_READERS = 2) - 每个消息处理完必须
END CONVERSATION,否则会话堆积,最终耗尽会话句柄 - 外部调用失败时,不能简单
ROLLBACK—— 那会让消息重回队列无限重试;应END CONVERSATION WITH ERROR并落库告警 - 激活过程本身不能有事务嵌套,
usp_SendToKafka若失败,需幂等重试或转入死信表
真正容易被忽略的是:Service Broker 的会话生命周期管理比消息体内容更关键。90% 的积压问题,根源不在 JSON 大小或网络慢,而在未正确 END CONVERSATION 导致会话堆积,最终触发 sys.transmission_queue 溢出,新消息无法入队。










