如何用SSIS读取MSMQ公共队列消息并写入SQL Server表?
用SSIS实现MSMQ消息读取并写入SQL Server表
一、配置SSIS连接管理器
- 新建SSIS包,添加两个连接管理器:
- MSMQ连接管理器:选择「公共队列」,输入队列的完整路径(比如
FormatName:DIRECT=OS:.\public$\YourQueueName) - SQL Server连接管理器:指向目标数据库,配置好登录凭据
- MSMQ连接管理器:选择「公共队列」,输入队列的完整路径(比如
二、将MSMQ消息存入变量
- 创建一个字符串类型的变量(比如
var_MSMQMessage),用于暂存读取到的消息内容 - 添加「Message Queue Task」,配置如下:
- 操作类型选择「Receive」
- 连接选择之前创建的MSMQ连接管理器
- 在「接收设置」里,勾选「将消息内容存储在变量中」,选择
var_MSMQMessage变量 - 根据需求设置消息处理方式(比如接收后删除消息,或保留)
三、将消息写入SQL Server表
根据消息格式的不同,分两种处理方式:
1. 简单文本消息
- 添加「Execute SQL Task」,配置如下:
- 连接选择SQL Server连接管理器
- SQL语句写插入逻辑,比如:
INSERT INTO TargetTable (MessageContent, ReceivedTime) VALUES (?, GETDATE()) - 切换到「参数映射」标签,将
var_MSMQMessage变量映射到SQL语句中的?占位符,数据类型选择NVARCHAR
2. 结构化消息(如XML)
- 添加「Script Task」,配置如下:
- 在「脚本」标签里,将
var_MSMQMessage设为只读变量 - 选择C#作为脚本语言,编辑脚本解析消息并插入数据库,示例代码:
using System; using System.Data; using Microsoft.SqlServer.Dts.Runtime; using System.Data.SqlClient; using System.Xml.Linq; public void Main() { string rawMessage = Dts.Variables["var_MSMQMessage"].Value.ToString(); XDocument messageDoc = XDocument.Parse(rawMessage); // 提取结构化字段 string orderId = messageDoc.Descendants("OrderId").First().Value; decimal amount = decimal.Parse(messageDoc.Descendants("Amount").First().Value); // 写入SQL Server string connStr = Dts.Connections["SQLServerConn"].ConnectionString; using (SqlConnection conn = new SqlConnection(connStr)) { conn.Open(); string insertSql = @"INSERT INTO OrderMessages (OrderId, Amount, ReceivedTime) VALUES (@OrderId, @Amount, GETDATE())"; using (SqlCommand cmd = new SqlCommand(insertSql, conn)) { cmd.Parameters.AddWithValue("@OrderId", orderId); cmd.Parameters.AddWithValue("@Amount", amount); cmd.ExecuteNonQuery(); } } Dts.TaskResult = (int)ScriptResults.Success; }
- 在「脚本」标签里,将
四、批量处理多条消息
如果队列中存在多条消息,用「Foreach Loop Container」实现循环处理:
- 将「Message Queue Task」和后续的写入任务拖入容器内
- 配置Foreach Loop的枚举器为「Foreach MSMQ Enumerator」,指向目标公共队列
- 设置枚举模式为「Enumerate messages in the queue」,每次循环自动读取下一条消息并赋值给变量
五、配置SQL Server代理作业
- 将SSIS包部署到SSIS目录(或保存到文件系统,根据SQL Server版本)
- 新建SQL Server代理作业,添加步骤:
- 类型选择「SQL Server Integration Services Package」
- 选择部署好的SSIS包,配置执行账号
- 设置作业计划为每隔数小时执行一次
注意事项
- 确保SQL Server代理的运行账号拥有MSMQ公共队列的读取权限,以及目标SQL Server表的写入权限
- 建议添加错误处理:比如将处理失败的消息移至死信队列,或写入错误日志表
- 如果消息是二进制格式,需将变量类型改为
Object,在Script Task中做相应的类型转换
内容的提问来源于stack exchange,提问作者pranay_gaurav123
相关产品推荐
相关产品推荐

