You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

将OLEDB行集转换为Queue对象以实现多线程与锁定

在SSIS的C# Script Task中实现线程安全的结果集队列访问

我之前在SSIS并行处理数据的场景里解决过几乎一样的问题——核心就是把Execute SQL Task返回的Full Result Set转换成线程安全的队列,确保并行执行的子任务不会重复处理同一条数据。下面是具体的实现步骤和代码:

步骤1:准备包级变量

首先,你需要新增一个包级变量来存储线程安全队列:

  • 变量名:User::CoreTablesQueue
  • 类型:Object

这个变量用来共享队列实例,让所有并行子任务都能访问到同一个队列。

步骤2:初始化线程安全队列(前置Script Task)

在所有并行子任务执行之前,先执行一个单独的Script Task,把User::CoreTables中的DataTable转换成线程安全的ConcurrentQueue<DataRow>:

using System;
using System.Data;
using System.Collections.Concurrent;
using Microsoft.SqlServer.Dts.Runtime;

public class ScriptMain
{
    public void Main()
    {
        // 从包变量中获取Execute SQL Task返回的结果集
        DataTable coreTablesDT = Dts.Variables["User::CoreTables"].Value as DataTable;
        
        // 校验数据有效性
        if (coreTablesDT == null || coreTablesDT.Rows.Count == 0)
        {
            Dts.TaskResult = (int)ScriptResults.Failure;
            return;
        }
        
        // 初始化线程安全队列
        ConcurrentQueue<DataRow> tablesQueue = new ConcurrentQueue<DataRow>();
        
        // 将DataTable的所有行加入队列
        foreach (DataRow row in coreTablesDT.Rows)
        {
            tablesQueue.Enqueue(row);
        }
        
        // 将队列存入包变量,供并行任务共享
        Dts.Variables["User::CoreTablesQueue"].Value = tablesQueue;
        
        Dts.TaskResult = (int)ScriptResults.Success;
    }
    
    enum ScriptResults
    {
        Success = 0,
        Failure = 1
    }
}

步骤3:并行子任务中获取队列首行数据

每个并行执行的子Script Task中,通过ConcurrentQueue的TryDequeue方法安全获取首行数据——这个方法是线程安全的,内部已经实现了锁机制,不会出现多个任务同时取到同一条数据的情况:

using System;
using System.Data;
using System.Collections.Concurrent;
using Microsoft.SqlServer.Dts.Runtime;

public class ScriptMain
{
    public void Main()
    {
        // 获取共享的线程安全队列
        ConcurrentQueue<DataRow> tablesQueue = Dts.Variables["User::CoreTablesQueue"].Value as ConcurrentQueue<DataRow>;
        
        if (tablesQueue == null)
        {
            Dts.TaskResult = (int)ScriptResults.Failure;
            return;
        }
        
        // 尝试取出队列首行数据(线程安全)
        if (tablesQueue.TryDequeue(out DataRow currentRow))
        {
            // 这里处理当前行的数据,比如获取列值
            string targetTableName = currentRow["你的列名"].ToString();
            
            // 执行你的业务逻辑:调用其他子进程、执行SQL操作等
            
            Dts.TaskResult = (int)ScriptResults.Success;
        }
        else
        {
            // 队列已空,没有剩余数据需要处理
            Dts.TaskResult = (int)ScriptResults.Success;
        }
    }
    
    enum ScriptResults
    {
        Success = 0,
        Failure = 1
    }
}

关键说明

  • 为什么用ConcurrentQueue? 它是.NET专门为多线程场景设计的线程安全集合,内部已经封装了高效的锁机制,比手动给普通Queue加lock更可靠,也避免了锁竞争带来的性能问题。
  • 初始化时机 一定要确保初始化队列的Script Task在所有并行子任务之前执行,否则子任务会拿到空队列。
  • 结果集校验 别忘了校验User::CoreTables是否有效(比如是否为null、是否有数据),避免空指针异常。

如果你非要用普通Queue配合手动锁实现(不推荐),可以参考下面的代码片段,但要注意锁的范围必须覆盖整个Dequeue操作:

// 初始化普通Queue
Queue<DataRow> tablesQueue = new Queue<DataRow>();
foreach (DataRow row in coreTablesDT.Rows)
{
    tablesQueue.Enqueue(row);
}
Dts.Variables["User::CoreTablesQueue"].Value = tablesQueue;

// 并行任务中取数据
Queue<DataRow> tablesQueue = Dts.Variables["User::CoreTablesQueue"].Value as Queue<DataRow>;
DataRow currentRow = null;

// 手动加锁确保线程安全
lock (tablesQueue)
{
    if (tablesQueue.Count > 0)
    {
        currentRow = tablesQueue.Dequeue();
    }
}

if (currentRow != null)
{
    // 处理数据逻辑
}

内容的提问来源于stack exchange,提问作者TomNash

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.20 10:06:30