这是本节的多页打印视图。 点击此处打印.

返回本页常规视图.

通信层

Quatm 分布式通信与数据管理系统:基于 TCP/IPC 的进程间通信框架

概述

通信层是 Quatm 框架的通信骨干。所有进程间通信(设备驱动、分析脚本、GUI 面板之间的数据交互)都经 TCP/IPC 通道以 发布/订阅(PUB/SUB) 模式进行,命令与控制则通过 RPC 实现。这确保了:

  • 故障隔离:某个相机驱动崩溃不会导致整个实验中断
  • 并行处理:图像处理与实验时序可同时运行
  • 网络透明:各组件可以运行在不同的机器上
实验脚本
     │
     ▼
[with realtime():]  ──→  FPGA  ──→  DAC/TTL 输出
     │
     ▼
分析脚本  ◄──  图像/数据流 (TCP/IPC PUB/SUB)
     │
     ▼
结果发布  ──→  DataManager  ──→  HDF5 存储

客户端类型

quatm.servers 提供了四种标准客户端,用于处理不同类型的进程间通信:

客户端用途适用场景
CommandClient控制驱动和进程向设备发送指令、启停驱动
DataClient传输小数据1D 曲线、标量值、元数据
ImageClient传输大数据块相机图像(2D 数组)
MessageClient日志消息传递错误、警告、信息、调试消息

基础客户端 — GenericClient

所有客户端的基类,封装了 TCP/IPC 连接与数据收发逻辑。

from quatm.servers.clients import GenericClient

client = GenericClient("my_channel")
client.subscribe("data_stream_name")
# 发送数据
client.send({"temperature": 25.0, "timestamp": 12345})
# 接收数据
if client.has_new_data():
    data = client.recv()
方法说明
subscribe(name)订阅指定的数据流
unsubscribe(name)取消订阅
send(data)发送数据
recv()接收数据
has_new_data()检查是否有新数据到达

数据客户端 — DataClient

专为小数据集设计的客户端,数据以 JSON 格式序列化,支持 NumPy 数组附件。

from quatm.servers.clients import DataClient

data_client = DataClient("analysis_result")
data_client.send({
    "fit_params": {"A0": 1.5, "sigma": 0.3, "pos": 10.2},
    "fit_curve": numpy_array
})

图像客户端 — ImageClient

优化用于传输相机图像等大型 2D 数据块。

from quatm.servers.clients import ImageClient

img_client = ImageClient("camera_output")
img_client.send(image_array)  # 发送 NumPy 图像数组

命令客户端 — CommandClient

用于向驱动程序发送控制命令,使用独立通道以确保命令不被数据流阻塞。

消息客户端 — MessageClient

标准化的消息传递接口,所有消息带有时间戳和来源信息。

from quatm.servers import send_error, send_warning, send_info, send_debug

send_info("实验开始执行")
send_warning("激光功率偏低,请检查")
send_error("相机连接失败")

分布式属性系统

Quatm 通过分布式属性树管理所有运行时配置参数,支持跨进程实时同步。

配置层级

层级说明示例
Configuration深层、基础系统属性运行哪些硬件、可用驱动列表
Properties对象特定参数ROI 位置、校准值、拟合参数
Preferences不影响数据的次要选择鼠标指针形状、窗口位置

PropertyAttribute

将分布式属性映射为 Python 属性,读写操作自动同步到属性树。

from quatm.servers.properties import PropertyAttribute, Properties

class MyComponent:
    # 声明分布式属性,默认值为 42.0
    _my_param = PropertyAttribute('/MyComponent/param', 42.0)

    def __init__(self):
        self._props = Properties('MyComponent')

    def do_something(self):
        # 读取属性值
        x = self._my_param.value
        # 写入属性值,自动同步到分布式树
        self._my_param.value = 99.0

Properties

管理组件内部属性的同步副本,后台守护线程经发布/订阅(TCP/IPC)连接到中央属性中心。

from quatm.servers.properties import Properties

props = Properties('MyProcess')
props.set('/path/to/property', value)
current_value = props.get('/path/to/property')

配置读取器

读取静态配置文件 configuration/configfile.json,提供系统级别的运行时配置。

from quatm.servers.configreader import ConfigReader

reader = ConfigReader()
config = reader.getConfiguration()

常用路径常量

quatm.servers 导出以下路径常量,方便定位工作目录:

常量说明
workpath工作根目录
driverpath驱动文件目录
iconpath图标资源目录
datapath数据存储目录
experiment_path实验脚本目录
configpath配置文件目录

1 - 客户端

Data/Image/Command 客户端:发布/订阅通信(TCP/IPC),连接目标由中央配置。

概述

quatm/servers/clients.py 提供通信客户端:连接目标由中央配置决定(Data / Image / Command Hub)。GenericClient 提供标准收发接口,适用于小数据(<10k);大 尺寸数据(图像)请用 ImageClient

客户端类

连接说明
GenericClient由配置决定通用收发(含 numpy 数组),subscribe/unsubscribe
DataClientData Hub普通数据的发布与订阅
ImageClientImage Hub图像级大数据发送(send 直接收 numpy 数组)
CommandClientCommand Hub向驱动发送控制命令(命令名拼入主题)
NpEncodernumpy 类型 JSON 序列化

关键方法(GenericClient)

方法说明
send(datadict, arr=None)发布数据(可附带一个 numpy 数组)
recv()接收一帧:(主题, 字典[, 数组]);超时返回 None
subscribe / unsubscribe订阅 / 取消订阅通道
has_new_data非阻塞检查是否有新数据
close关闭 PUB/SUB socket(幂等)

用法

from quatm.servers.clients import DataClient

dc = DataClient(name="pmt")
dc.send({"pmt": 1234}, arr=numpy_array)
topic, data = dc.recv()

2 - 属性树

分布式分层属性树:本地副本经 Property Hub 与 Propertylogger 双向同步。

概述

quatm/servers/properties.py 提供属性树(Properties):管理控制中心的配置与 属性,维护属性字典的同步副本。本地 set / get 会发布到 Property Hub,中央的变 更也推送回本地副本。属性寻址类似文件系统(如 Drivers/xxx 或以 / 开头的绝对路 径)。

关键方法

方法说明
set(key, value)设置属性并发布变更到 Property Hub
get(key, default=...)读取值(总是返回深拷贝;缺失时创建默认值)
delete(key)删除条目(含全部子条目)并广播
changes()返回自上次调用以来的变更键
close / close_all关闭 socket 与后台线程

PropertyAttribute 把属性树中的一个属性当作普通 Python attribute 读写(宿主需有 _props)。

用法

from quatm.servers.properties import Properties

props = Properties()
props.set("laser/power", 50e-3)
power = props.get("laser/power")

3 - 数据管理

多流数据汇总:按实验运行对齐,缺失超时补零后写 HDF5。

概述

quatm/servers/datamgr.py 提供数据管理(DataSummary):汇总多个输入通道的数 据流,整理成每个 run 的完整字典。收到数据先累积在 incompleteData,当所有订阅通 道的数据都到达时整条记录移入 completeData;若 run 结束超时仍有数据未到,缺失部 分以零补齐。

关键方法

方法说明
initDataq(重新)订阅属性里配置的通道与实验开始/结束事件
recvData / processIncomingData接收数据并归并到最近的 run 记录
savetoFile把收集的 run 字典写入 HDF5 结构化数据集
loadFromFile从 HDF5 读回一组已存的 run 字典
clear清空已收集的全部数据
run_forever阻塞式循环持续接收数据

用法

from quatm.servers.datamgr import DataSummary

summary = DataSummary(name="sum")
summary.initDataq()     # 订阅属性中配置的通道
summary.run_forever()   # 持续收数;savetoFile 把 run 字典写入 .h5

4 - 消息系统

标准化消息流客户端:把带时间戳/来源的消息发到 Message Hub。

概述

quatm/servers/messageclient.py 提供消息系统MessageClient 是标准化的消 息流客户端,把带时间戳与来源文件的消息发到 Message Hub(PUB)。日常代码一般不要 直接构造,而是使用模块级发送函数。

发送函数

函数说明
send_info(msg)Info 级:进入消息流作为普通通知
send_warning(msg)Warning 级:进入消息流并触发警告通知
send_error(msg)Error 级:进入消息流并触发错误通知
send_debug(msg)Debug 级:仅在调试通道显示

用法

from quatm.servers.messageclient import send_info, send_error

send_info("实验启动完成")
send_error("相机连接超时")

5 - 配置与属性持久化

configreader + propertylogger:JSON 配置读取与属性树落盘、活跃流监视。

模块

模块功能
configreader.ConfigReader只读封装:读取 configfile.jsongetConfiguration()
configreader.Config()读取 configfile.json 并返回配置字典
configreader.Properties()读取属性盘文件 properties.json 并返回字典
propertylogger.run_propertylogger()属性树落盘 + INIT 应答(Property Logger 服务)
datalogger.DataStreamLogger监视活跃数据流并写入属性树
imagelogger.ImageStreamLogger监视活跃图像流并写入属性树

用法

from quatm.servers.configreader import ConfigReader

cfg = ConfigReader()
config = cfg.getConfiguration()   # 完整配置字典

6 - 外围服务模块

InfluxDB、通知推送、图像网络服务、Web Hub 与 XSUB/XPUB 代理等外围服务。

服务列表

模块功能
influxdbInfluxDB 客户端封装:硬件状态写入(data2db)与历史查询(query_data),可作 RPC 服务
notifications通知推送:配置服务、push_message 推送、post_to_mattermost 兼容入口
image_network_server图像网络服务:把各图像流最新帧刷到 Web figure
webhubWeb 远程控制平台的服务器入口(quatm.web,自 configfile Servers 段取 host/port)
xsub_xpubXSUB/XPUB 转发服务器(run_server
experiment_status实验状态监控 / 预警通知
propertylogger / datalogger / imagelogger属性树落盘与活跃数据/图像流监视(见配置系统

用法

from quatm.servers.notifications import push_message

push_message("实验完成")   # 经已启用的推送服务发出(绝不抛异常)