public override bool ExeModule()
        {
            Stopwatch.Restart();
            try
            {
                if (string.IsNullOrEmpty(QueueKey))
                {
                    Logger.AddLog($"[{ModuleParam.ModuleName}] 队列Key不能为空!", eMsgType.Warn);
                    ChangeModuleRunStatus(eRunStatus.NG);
                    return false;
                }

                // 确保队列存在并定义槽位
                if (!Solution.Ins.QueueDic.ContainsKey(QueueKey))
                    Solution.Ins.QueueDic[QueueKey] = new JGTechVision.Services.DataOut(QueueKey);

                // 核心1:获取对应的全局缓存队列数据
                var outQueue = Solution.Ins.QueueDic[QueueKey];

                if (!_slotsDefined)
                {
                    DefineAllSlots(outQueue);
                    _slotsDefined = true;
                }

                lock (outQueue)
                {
                    foreach (var slot in QueueSlots.Where(s => s.IsEnable))
                    {
                        // 核心2:通过变量链接获取上游模块的值 (LinkVar.Text 是 "&模块.变量" 格式)
                        object val = GetLinkValue(slot.LinkVar.Text);
                        if (val == null) continue;

                        int idx = slot.SlotIndex;
                        switch (slot.DataType)
                        {
                            case "double":
                                var dList = (List<double>)outQueue.GetDataQueue(idx);
                                if (dList.Count >= outQueue.LimitLength) dList.RemoveAt(0);
                                // 核心:添加数据
                                dList.Add(Convert.ToDouble(val)); break;
                            case "int":
                                var iList = (List<int>)outQueue.GetDataQueue(idx);
                                if (iList.Count >= outQueue.LimitLength) iList.RemoveAt(0);
                                iList.Add(Convert.ToInt32(val)); break;
                            case "string":
                                var sList = (List<string>)outQueue.GetDataQueue(idx);
                                if (sList.Count >= outQueue.LimitLength) sList.RemoveAt(0);
                                sList.Add(val.ToString()); break;
                            case "bool":
                                var bList = (List<bool>)outQueue.GetDataQueue(idx);
                                if (bList.Count >= outQueue.LimitLength) bList.RemoveAt(0);
                                bList.Add(Convert.ToBoolean(val)); break;
                        }
                    }
                }

                // 唤醒等待的 DataOut 读取方
                if (Solution.Ins.QueueSignDic.ContainsKey(QueueKey))
                    // 核心:唤醒Dataout
                    Solution.Ins.QueueSignDic[QueueKey].Set();

                ChangeModuleRunStatus(eRunStatus.OK);
                return true;
            }
            catch (Exception ex)
            {
                Logger.AddLog($"[{ModuleParam.ModuleName}] {ex.Message}", eMsgType.Error);
                ChangeModuleRunStatus(eRunStatus.NG);
                return false;
            }
        }

QueueSignDic 是一个跨模块信号唤醒机制,本质是一个 生产者-消费者 的通知器。

这个项目里的数据流是这样的:

  1. 生产者(DataIn 模块) 往某个命名队列(由 QueueKey 标识)里写数据
  2. 消费者(DataOut 模块) 从同一个队列里取数据

QueueSignDic 的作用就是当消费者在等待数据时,生产者写完数据后"喊醒"它。

具体来说就是三个场景:

初始化时 — 创建 AutoResetEvent(false)(初始为"未信号"状态)

[DataOut.cs:24](/D:\JGTechVision\GJTechVisionV1.0.0\01_Sourse - 0604-2\VM\01Main\VM.Start\Services\DataOut.cs:24) 在 DataOut 构造函数里为每个队列创建一个信号量

生产者写完数据 — 调用 Set() 打开信号

[DataIn.cs:102](/D:\JGTechVision\GJTechVisionV1.0.0\01_Sourse - 0604-2\VM\01Main\VM.Start\Services\DataIn.cs:102) 中的 WakeAll() 和 [DataInViewModel.cs:96](/D:\JGTechVision\GJTechVisionV1.0.0\01_Sourse - 0604-2\VM\02Plugins\010文件通讯\Plugin.DataIn\ViewModels\DataInViewModel.cs:96) 都是在数据入队后 Set()

消费者等待数据 — 调用 WaitOne(TimeOut) 阻塞等待

[DataOutViewModel.cs:62-73](/D:\JGTechVision\GJTechVisionV1.0.0\01_Sourse - 0604-2\VM\02Plugins\010文件通讯\Plugin.DataOut\ViewModels\DataOutViewModel.cs:62-73) 里当 IsWait=true 时,先检查队列是否为空,如果为空就阻塞等待,直到被生产者 Set() 唤醒或超时

AutoResetEvent 为什么合适? 因为它是"自动重置"的:一个线程被唤醒后,它自动变回"未信号"状态,不需要手动重置。这非常适合"来一波数据唤醒一次消费者消费一波"的模式。用 Dictionary<string, AutoResetEvent> 是因为系统里可能有多条命名队列并行工作,每条队列独立唤醒自己的消费者。



核心内容:

两阶段入队模型 — “暂存 → 批量提交”:插件/脚本先调 AddDoubleQueueIn(3.14) 暂存到 m_DataQueueInListExeModule 执行时统一推入 QueueDic[key] 对应的 DataOut 队列,最后 finally 清空缓冲并 WakeAll() 唤醒等待者

ExeModule() 逐行拆解 — 8 步:找 DataOut → lock(dataOut) 加锁 → 边界校验(QueueIndex 越界/负数)→ 逐项类型字符串比对 → 10 种类型分发 → 异常静默吞掉 → finally 清空 + WakeAll

10 种类型分发机制 — 不是泛型,而是字符串类型名("double"/"int"/"string"/"bool" 及其数组 + "HImage")比对后强转为对应的 List<T>,优点是灵活,缺点是拼写错误只能在运行时发现

QueueIndex 偏移 — 多个 DataIn 可以向同一个 DataOut 的不同槽位区域写入,不一定要从槽位 0 开始

完整数据流 — 从变量链接获取值 → AddXXXQueueIn 暂存 → ExeModule 写入 → WakeAll 唤醒 → DataOut 读取的端到端链路

5 个隐患 — 字符串类型匹配脆弱、异常被注释掉静默吞掉、lock(dataOut) 粒度过大阻塞读取端、无 QueueKey 不存在时的容错、10 个 AddXXX 方法重复代码

DataIn.cs 完整解析 — 跨模块数据入队引擎

文件: Services\DataIn.cs (230行)
继承: ModuleBase — 本身可作为流程节点使用
角色: 将上游模块的输出数据按类型推入共享队列 (DataOut), 实现跨模块/跨流程的数据传递
配套: DataOut.cs (出队端) / Plugin.DataIn (我们创建的 UI 插件) / Plugin.DataOut (UI 插件)


1. 它解决什么问题

在视觉流程中, 模块 A 的输出需要传给模块 B 使用。普通变量链接只能在同一 Project 内传递。DataIn + DataOut 通过 QueueKey 命名的共享队列 实现了跨 Project 的数据传递

Project A (采集流程):                   Project B (汇总流程):
  ┌──────────┐                             ┌──────────┐
  │ DataIn   │──→ QueueDic["Q1"] ──→       │ DataOut  │
  │ QueueKey │     (全局共享队列)           │ QueueKey │
  │ ="Q1"    │                             │ ="Q1"    │
  └──────────┘                             └──────────┘

2. 源码结构 (230行)

DataIn : ModuleBase
│
├── 属性
│   ├── QueueKey           ← 队列标识 (与 DataOut 的 Key 对应)
│   └── QueueIndex         ← 数据写入的起始槽位偏移
│
├── 内部缓冲 (用完即清)
│   ├── m_DataQueueInList  ← List<object> 暂存待入队的数据
│   └── m_DataTypeInList   ← List<string> 暂存对应的类型名
│
├── ★ 10 个 AddXXXQueueIn 方法 (数据暂存)
│   ├── AddIntQueueIn(int)
│   ├── AddDoubleQueueIn(double)
│   ├── AddStringQueueIn(string)
│   ├── AddBoolQueueIn(bool)
│   ├── AddIntListQueueIn(List<int>)
│   ├── AddDoubleListQueueIn(List<double>)
│   ├── AddStringListQueueIn(List<string>)
│   ├── AddBoolListQueueIn(List<bool>)
│   ├── AddHImageQueueIn(HImage)
│   └── AddHImageListQueueIn(List<HImage>)
│
├── WakeAll()              ← 唤醒所有等待此队列的 DataOut
└── ExeModule()            ← 将缓冲数据按类型推入 DataOut 队列

3. 两阶段入队模型

DataIn 采用 “暂存 → 批量提交” 的两阶段模型:

阶段1: 暂存 (由插件/脚本调用 AddXXXQueueIn)
  │
  ├→ AddDoubleQueueIn(3.14)    → m_DataQueueInList=[3.14]    m_DataTypeInList=["double"]
  ├→ AddStringQueueIn("ABC")   → m_DataQueueInList=[3.14,"ABC"] m_DataTypeInList=["double","string"]
  └→ AddBoolQueueIn(true)      → 同上...
  │
  ▼
阶段2: 批量入队 (ExeModule 执行)
  │
  ├→ 根据 QueueKey 找到 DataOut
  ├→ lock(dataOut) 线程安全
  ├→ 逐个按类型推入 dataOut 的对应槽位
  └→ finally: 清空 m_DataQueueInList + m_DataTypeInList

4. ExeModule() — 核心入队逻辑 (109→220行)

public override bool ExeModule()
{
    bool flag = true;
    try
    {
        // ① ★ 找到对应的 DataOut 队列 (跨流程共享)
        if (!Solution.Ins.QueueDic.ContainsKey(QueueKey))
        {
            Logger.AddLog($"没有找到对应的队列 [{QueueKey}]");
            return false;
        }
        DataOut dataOut = Solution.Ins.QueueDic[QueueKey];

        // ② ★ 加锁: 保证并发安全
        lock (dataOut)
        {
            int dataOutLength = dataOut.GetQueueCount();

            // ③ 边界校验
            if (dataOutLength < (QueueIndex + m_DataQueueInList.Count))
            {
                Logger.AddLog("入队变量的长度超过数据出队的变量的长度");
                return false;  // 槽位索引越界
            }
            if (QueueIndex < 0)
            {
                Logger.AddLog("入队变量的索引为负值");
                return false;
            }

            // ④ ★ 逐项按类型推入队列
            for (int i = 0; i < m_DataQueueInList.Count; i++)
            {
                // 类型匹配检查
                if (m_DataTypeInList[i] != dataOut.GetDataType(i + QueueIndex))
                {
                    Logger.AddLog("数据入队类型与对应的数据出队类型不匹配");
                    WakeAll();  // 唤醒等待者 (避免死锁)
                    return false;
                }

                // ★ 按 10 种类型分发
                switch (m_DataTypeInList[i])
                {
                    case "int":
                        List<int> list1 = (List<int>)dataOut.GetDataQueue(i + QueueIndex);
                        list1.Add((int)m_DataQueueInList[i]);
                        break;
                    case "double":
                        List<double> list2 = (List<double>)dataOut.GetDataQueue(i + QueueIndex);
                        list2.Add((double)m_DataQueueInList[i]);
                        break;
                    case "string":
                        List<string> list3 = (List<string>)dataOut.GetDataQueue(i + QueueIndex);
                        list3.Add((string)m_DataQueueInList[i]);
                        break;
                    case "bool":
                        List<bool> list4 = (List<bool>)dataOut.GetDataQueue(i + QueueIndex);
                        list4.Add((bool)m_DataQueueInList[i]);
                        break;
                    case "HImage":
                        List<HImage> list5 = (List<HImage>)dataOut.GetDataQueue(i + QueueIndex);
                        list5.Add((HImage)m_DataQueueInList[i]);
                        break;
                    // ... 数组类型同理 (List<int[]> 等)
                }
            }
        }
    }
    catch (Exception ex)
    {
        // ★ 异常被静默吞掉 — 隐患
        /// MessageBox.Show(ex);
    }
    finally
    {
        // ⑤ ★ 入队完成 → 清空缓冲 → 唤醒等待者
        m_DataQueueInList.Clear();
        m_DataTypeInList.Clear();
        WakeAll();  // → QueueSignDic[QueueKey].Set() 唤醒 DataOut.GetStr()
    }
    return flag;
}

5. 类型分发机制

DataIn 和 DataOut 之间通过字符串类型名做类型匹配, 而非泛型:

// DataIn 端声明类型
m_DataTypeInList[i] = "double"

// DataOut 端注册类型 (通过 DefineDoubleQueue)
m_DataTypeList[idx] = "double"

// 入队时比对
if ("double" == "double")  → 强转为 List<double> → Add

支持的 10 种类型:

类型字符串 运行时容器类型 AddXXXQueueIn 参数
"int" List<int> int
"double" List<double> double
"string" List<string> string
"bool" List<bool> bool
"int[]" List<List<int>> List<int>
"double[]" List<List<double>> List<double>
"string[]" List<List<string>> List<string>
"bool[]" List<List<bool>> List<bool>
"HImage" List<HImage> HImage
"HImage[]" List<List<HImage>> List<HImage>

6. WakeAll() — 信号唤醒

private void WakeAll()
{
    if (!Solution.Ins.QueueSignDic.ContainsKey(QueueKey))
        Solution.Ins.QueueSignDic.Add(QueueKey, new AutoResetEvent(false));

    Solution.Ins.QueueSignDic[QueueKey].Set();
    // → 唤醒所有在 GetStr() 中 WaitOne 等待的 DataOut
}

调用时机:

  • 入队成功 → 通知 DataOut “有新数据了”
  • 类型不匹配 → 也唤醒 (避免 DataOut 永久阻塞)
  • finally 块 → 无论如何都唤醒

7. 完整数据流

外部调用 (如 Plugin.DataIn.ExeModule)
  │
  ├→ 通过变量链接获取上游模块的值
  │    GetLinkValue(slot.LinkVar.Text) → 3.14
  │
  ├→ 暂存到 DataIn 内部缓冲
  │    dataIn.AddDoubleQueueIn(3.14)     → m_DataQueueInList=[3.14]
  │
  └→ 触发入队
       dataIn.ExeModule()
         │
         ├→ QueueDic["Q1"] → 找到 DataOut 实例
         ├→ lock(dataOut)
         ├→ 类型校验: "double" == dataOut.GetDataType(0)
         ├→ List<double> list = dataOut.GetDataQueue(0)
         ├→ list.Add(3.14)                     ← ★ 数据写入了!
         ├→ WakeAll() → QueueSignDic["Q1"].Set()
         └→ m_DataQueueInList.Clear()          ← 清空缓冲

  ════════ 另一边, DataOut 正在等待 ════════

  DataOut.GetStr() 或 Plugin.DataOut.ExeModule
    │
    ├→ lock(dataOut)
    ├→ List<double> list = dataOut.GetDataQueue(0)
    ├→ double val = list.Last()               ← ★ 数据读取!
    └→ 返回给下游

8. QueueIndex 偏移机制

QueueIndex 允许不从槽位 0 开始写入, 而是从指定偏移开始:

DataOut 定义了 6 个槽位:
  [0]=List<double>  [1]=List<double>  [2]=List<int>
  [3]=List<string>  [4]=List<bool>    [5]=List<bool>

QueueIndex = 2 时:
  m_DataQueueInList = [100, "hello", true]
  → 实际写入: 槽位[2]=100, 槽位[3]="hello", 槽位[4]=true

这意味着多个 DataIn 可以向同一个 DataOut 的不同槽位区域写入。


9. 设计分析

优点

  1. 队列共享: 通过 Solution.Ins.QueueDic 全局字典, 任意 Project 的 DataIn 都能写入
  2. 类型安全: 入队时显式校验 DataType 字符串, 防止类型错误
  3. 信号驱动: WakeAll() 确保 DataOut 不会被永久阻塞
  4. 缓冲清理: finally 保证缓冲总是被清空, 不会残留数据

隐患

问题 说明
字符串类型匹配脆弱 类型靠字符串比对, 拼写错误只能在运行时发现
异常静默吞掉 catch (Exception ex) { /// MessageBox.Show(ex); } — 异常被注释掉
锁粒度过大 lock(dataOut) 锁住了整个入队过程, DataOut 读取也被阻塞
无 QueueKey 不存在时的容错 直接返回 false, 但 WakeAll 仍在 finally 中执行
10 个 AddXXX 方法重复代码 除了类型不同, 逻辑完全一样, 可以用泛型方法消除

文档说明: 基于 DataIn.cs (230行) 源码静态分析生成。与 DataOut.cs 配对使用, 共同构成跨模块数据队列系统。当前版本 2026-06-10。

Logo

AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。

更多推荐