DataIn.cs 完整解析 — 跨模块数据入队引擎
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 是一个跨模块信号唤醒机制,本质是一个 生产者-消费者 的通知器。
这个项目里的数据流是这样的:
- 生产者(DataIn 模块) 往某个命名队列(由
QueueKey标识)里写数据 - 消费者(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_DataQueueInList,ExeModule 执行时统一推入 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. 设计分析
优点
- 队列共享: 通过
Solution.Ins.QueueDic全局字典, 任意 Project 的 DataIn 都能写入 - 类型安全: 入队时显式校验
DataType字符串, 防止类型错误 - 信号驱动:
WakeAll()确保 DataOut 不会被永久阻塞 - 缓冲清理:
finally保证缓冲总是被清空, 不会残留数据
隐患
| 问题 | 说明 |
|---|---|
| 字符串类型匹配脆弱 | 类型靠字符串比对, 拼写错误只能在运行时发现 |
| 异常静默吞掉 | catch (Exception ex) { /// MessageBox.Show(ex); } — 异常被注释掉 |
| 锁粒度过大 | lock(dataOut) 锁住了整个入队过程, DataOut 读取也被阻塞 |
| 无 QueueKey 不存在时的容错 | 直接返回 false, 但 WakeAll 仍在 finally 中执行 |
| 10 个 AddXXX 方法重复代码 | 除了类型不同, 逻辑完全一样, 可以用泛型方法消除 |
文档说明: 基于 DataIn.cs (230行) 源码静态分析生成。与 DataOut.cs 配对使用, 共同构成跨模块数据队列系统。当前版本 2026-06-10。
AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。
更多推荐




所有评论(0)