流程 (Flow)
目录
入门
分屏幕指南
- 列表界面 · 编辑界面 (可视化画布) · 执行历史界面
消息·节点参考
- 消息结构 · 按触发器的负载示例 · 节点目录
- 触发器 · 过滤器 · 变换
- 动作 — 集成/存储 · 发布资产事件 · 领域 CRUD
- 边缘 · 外部集成 · 流程控制
- 节点选项详细参考 · 脚本节点编写 · JS 执行环境规范 · 内存安全编写模式 · 移动推送通知渠道 · 外部认证令牌自动刷新
示例·模式
运维
生产运维
- 常见问题 (FAQ) · 运维最佳实践 · 安全·敏感信息处理
- 流程指标与报警 · 消息处理语义与背压 · 执行历史中能看到什么
- 集群·HA 动作 · 端到端追踪 · 图模式目录 · 流程测试最佳实践
- 性能限制与调优 · 审计·历史追踪 · 新流程部署检查清单 · 部署策略 (Canary/Blue-Green/A·B) · 紧急应对流程
- 节点快速设置参考 · 流程 REST API · Webhook 触发器 · OPC/PLC 工业集成模式
- 外部系统集成 Cookbook · 数据变换 Cookbook · 可复用脚本集
问题排查
其他
概述
流程是以可视化图形方式定义与外部系统(MES/ERP/SCADA 等)之间的数据集成、以及领域对象(标签·资产·工单·作业人员等)的自动创建·更新的自动化菜单。运营人员无需编写代码即可直接设计·部署·运维自动化场景。
将100余种节点通过拖放的方式放置在画布上,并用连线连接,从而构建数据处理管道。
路径: 左侧菜单 > Automation > 流程
学习路线图 — 按角色推荐的入门顺序
本手册是超过4,000行的综合指南。根据自己的角色和目的从相应章节开始阅读会更高效。
👶 初次使用者 (1小时内创建第一个流程)
→ 完成首个流程部署。之后需要的节点可在节点目录中检索。
🧑🏭 现场运营人员 (编写自动化场景)
- 按触发器的负载示例 — 理解实际数据结构
- 节点选项详细参考 — 熟悉常用节点选项
- 数据变换 Cookbook — 复制使用常见变换模式
- 图模式目录 — 选择接线模式
- 示例流程 (17种) — 参考按场景的完整示例
🛠 系统管理员 (运维·调优·故障应对)
🔌 开发者·集成工程师 (外部系统集成)
- 外部系统集成 Cookbook — Slack/Teams/Jira/SAP 示例
- 流程 REST API — 通过程序操作流程
- Webhook 触发器 — 从外部触发流程
- OPC/PLC 工业集成模式 — 工业现场场景
- JS 执行环境规范 + 可复用脚本集
🔐 安全负责人 (审计·认证)
- 安全·敏感信息处理 — 凭证保管
- 外部认证令牌自动刷新 — OAuth2 令牌运维
- 权限 — 各角色可执行的操作
- 审计·历史追踪 — 保存变更/执行历史
📚 快速参考 (已经熟悉的用户)
| 查找信息 | 章节 |
|---|---|
| 节点 ID 与一句话说明 | 节点目录 |
| 节点选项默认值 | 节点选项详细参考 |
| 消息 JSON 示例 | 按触发器的负载示例 |
| 直接使用的脚本 | 可复用脚本集 · 数据变换 Cookbook |
| API 调用 curl | 流程 REST API |
| 解读执行历史 | 执行历史中能看到什么 |
| 问题排查指南 | 分步调试指南 · 常见问题 |
| 快速选项表 | 节点快速设置参考 |
所有章节链接均为同一文档内的锚点。
Ctrl+F关键词搜索也很有效。
核心概念
| 术语 | 说明 |
|---|---|
| 流程 (Flow) | 由节点和连线(关系)的集合构成的一个自动化工作流。是具有入口点的有向图,通过部署/解除开关来控制是否激活 |
| 节点 (Node) | 接收消息处理后传递给下一个节点的单位。分为7个类别(触发器·过滤器·变换·动作·外部集成·流程控制·边缘) |
| 关系 (Relation) | 从节点输出到下一个节点的连线的标签。节点直接指定SUCCESS/FAILURE/TRUE/FALSE/MATCH/NO_MATCH/DEFAULT/THROTTLED/EXHAUSTED等标签进行分支。所有标签都统一为大写 |
| 消息 (Message) | 在流程内部流动的负载。包含type(分类)、originator(主体实体)、data(正文)、metadata(上下文) |
| 触发器 (Trigger) | 作为流程起点的节点。分为领域事件(标签点·报警·资产事件等)、外部进入(Webhook·MQTT·外部DB)、时间(调度)三种 |
快速上手
通过以下4个步骤最快熟悉流程。
- 在列表界面点击
새 플로우按钮,创建空流程(仅输入名称和描述)。 - 在编辑界面依次拖拽左侧调色板中的触发器节点(如
flow_on_tag_alarm) → 过滤器 → 动作(如flow_send_email)并放置,用连线连接节点之间。 - 点击各节点在右侧检查器中输入选项,然后点击右上角的保存 → 测试运行,确认一次动作。
- 打开顶部部署开关后,每次有触发事件进入时流程都会自动执行,可在实时调试面板和执行历史界面中确认结果。
端到端教程
这是一个从头到尾创建可在运营环境中实际使用的自动化流程的分步教程。场景: 当电机温度超过80°C时自动发布紧急维护工单并通过邮件通知负责人。
步骤1 — 创建新流程
- 在左侧菜单中点击Automation > 流程 → 打开列表界面
- 点击右上角新建流程按钮
- 输入以下信息后点击确认
| 项目 | 输入值 |
|---|---|
| 名称 | 모터 과열 자동 정비 발행 |
| 说明 | [자동화] 모터 자산의 온도 80°C 초과 시 긴급 정비 작업지시 + 이메일. 담당: ops@example.com |
流程创建后会自动打开空画布的编辑界面。
步骤2 — 放置触发器节点
从资产遥测事件接收电机温度数据。
- 展开左侧调色板的触发器 (Trigger) 类别
- 将
flow_on_tag_point节点拖拽到画布上 - 点击节点 → 在右侧检查器中输入以下选项
| 选项 | 值 |
|---|---|
| 显示名 | 태그 포인트 인입 |
tag_id_pattern | MOTOR-*.TEMP |
通过
tag_id_pattern仅让电机温度标签通过,可大幅减少后续处理量。其他标签会被处理为SKIPPED,不计入计数器。
步骤3 — 阈值过滤器
添加脚本过滤器,只让超过80°C的消息通过。
- 从过滤器 (Filter) 类别拖拽
flow_script_filter - 从触发器节点的输出端口连线到新过滤器节点的输入端口
- 在检查器中输入以下内容
| 选项 | 值 |
|---|---|
| 显示名 | 임계값 필터 (80°C 초과) |
language | JS |
script | data.value > 80 |
步骤4 — 领域变更 (标签 → 资产)
由于报警·工单以资产为单位发布更自然,将消息的originator从标签转换为上级资产。
- 从变换 (Transform) 类别拖拽
flow_change_originator - 从过滤器节点的
TRUE输出连线 - 检查器:
| 选项 | 值 |
|---|---|
| 显示名 | Tag → Asset 변경 |
entity_type | Asset |
id_field | metadata.asset_id |
metadata.asset_id在标签点接入时会自动填充。如果消息中没有,也可以用脚本从标签ID的前缀部分(如MOTOR-001.TEMP→MOTOR-001)提取。
步骤5 — 发布工单
自动创建紧急维护工单。
- 从动作 (Action) — 领域 CRUD 类别拖拽
flow_create_work_order - 从变换节点的
SUCCESS输出连线 - 检查器:
| 选项 | 值 |
|---|---|
| 显示名 | 긴급 정비 작업지시 생성 |
asset_id_field | originator.id |
title_field | (静态值) 긴급 점검 — 모터 과열 |
master_id_field | (可选) data.master_id (未设置时自动生成) |
description_field | (静态值) 자동 발행: 임계 온도 초과로 긴급 점검 필요 |
default_priority | HIGH |
步骤6 — 邮件通知 (成功分支)
工单发布成功后通知负责人。
- 从外部集成 (External) 类别拖拽
flow_send_email - 从工单节点的
SUCCESS输出连线 - 检查器:
| 选项 | 值 |
|---|---|
| 显示名 | 정비 담당자 이메일 |
to | ops@example.com |
subject_template | [과열 정비] ${originator.id} 작업지시 ${data.work_order_id} |
body_template | 자산 ${originator.id} 의 온도가 ${data.value}°C 로 상승하여 자동으로 긴급 정비가 발행되었습니다.\n작업지시 ID: ${data.work_order_id} |
步骤7 — 失败分支处理
工单发布本身也可能失败(如资产已删除或权限问题)。立即向运营人员发送推送通知。
- 拖拽外部集成的
flow_send_push - 从工单节点的
FAILURE输出连线 - 检查器:
| 选项 | 值 |
|---|---|
title | 자동화 실패 |
body_template | ${originator.id} 정비 자동 발행 실패: ${data.error} |
步骤8 — 保存和测试运行
- 点击右上角保存按钮。会自动加载快照,之后可回滚。
- 点击右上角测试运行按钮 → 在JSON编辑器中输入以下内容 → 发布
{
"type": "POST_TELEMETRY",
"originator": { "entity_type": "Tag", "id": "MOTOR-001.TEMP" },
"data": { "value": 92.5 },
"metadata": { "ts": 1746247200000, "tag_id": "MOTOR-001.TEMP", "asset_id": "MOTOR-001" }
}
对话框下方显示✓ JSON OK (type=POST_TELEMETRY) → 点击发布。
- 在实时调试面板确认以下内容:
- 触发器节点亮绿灯 → 消息通过
- 过滤器节点: 因为
92.5 > 80所以走TRUE分支 - 变换节点: originator被更改为Asset/MOTOR-001
- 工单节点: SUCCESS,自动赋予
data.work_order_id - 邮件节点: 尝试发送
步骤9 — 验证和部署
- 在工单界面确认是否注册了新工单
- 确认邮件是否正常到达 (测试环境的邮箱)
- 使用低于80°C的消息再测试一次(应在过滤器中被阻断):
{ "type": "POST_TELEMETRY", "originator": {"entity_type":"Tag","id":"MOTOR-001.TEMP"},
"data": {"value": 70}, "metadata": {"asset_id":"MOTOR-001"} }
过滤器节点的分支为FALSE,后续节点显示灰色 — 正常。
- 验证所有分支后,用全部计数重置重置统计窗口
- 打开顶部部署开关
现在在实际运营环境中,一旦电机温度超过80°C,即会自动发布维修工单并向负责人发送邮件。
步骤10 — 运维监控
部署后5~10分钟内建议确认以下内容比较安全。
| 位置 | 确认项目 |
|---|---|
| 列表界面 | 对应流程行的执行次数是否在正常范围内增加(是否爆发) |
| 实时调试面板 | 是否没有节点失败 |
| 执行历史界面 | 如有失败消息,确认NODE_ERROR的原因 |
| 接收邮件 | 是否未以非预期频率发送通知 |
下一步
若要进一步完善此流程:
- 添加重试接线 — 邮件/推送发送失败时进行退避重试 (参见示例8)
- 自动调整波段 — 用6小时平均值 ± 3σ 自动修正阈值 (参见示例6)
- 累计停机时间 — 将过热历史累积到资产属性 (参见示例14)
- 隔离质量产线 — 连续过热时自动停止产线 (参见示例15)
界面组成
流程由以下三个界面组成。
| 界面 | 用途 |
|---|---|
| 列表 | 已注册流程一览·搜索·批量部署·导入 |
| 编辑 | 通过可视化画布放置·连接·设置节点 |
| 执行历史 | 查询节点单位的执行日志·时间线 |
列表界面
由顶部搜索·创建区域和流程一览表组成。
顶部工具
| 项目 | 说明 |
|---|---|
| 状态过滤器 | 전체 / 배포 / 해제 |
| 名称·说明搜索 | 按关键词过滤流程 |
| 新建流程 | 创建空流程(输入名称·说明) |
| 导入 | 上传导出得到的JSON恢复流程 |
| 全部重新部署 | 一次性重新加载所有已激活的流程 |
| 刷新 | 重新加载列表 |
一览表
列表表格按25行分页,每行同时显示处理趋势的迷你图和错误率的甜甜圈图表。翻页后图表状态保持不变。
| 列 | 说明 |
|---|---|
| 选择 | 用于批量部署·解除的复选框 |
| 状态 | 배포 / 해제 徽章 |
| 流程 ID | FLOW_NNNNN 格式的序列 ID |
| 名称 | 说明 | 运营人员指定的元信息 |
| 节点 | 包含的节点数量 |
| 触发器 | 触发器节点数量 |
| 执行次数 | 累计消息处理计数 |
| 处理时间 | 节点平均/最近处理时间 |
| 错误 | 累计错误计数 |
| 最终修改 | 图形最后保存时间 |
| 动作 | 편집 / 삭제 按钮 |
批量部署·解除
对选中的流程一次性执行部署(激活)或解除(停用)。只有已部署的流程才能接收触发事件。
全部重新部署
顶部全部重新部署按钮会重新加载所有已激活的流程。请在以下情况使用。
- 从外部批量导入图形后
- 想重新注册具有自身调度器的触发器(调度·外部MQTT订阅·外部DB轮询)时
- 怀疑运行中的缓存一致性出现问题时
编辑界面 (可视化画布)
从画布左侧调色板拖拽节点放置后,点击并拖拽节点的输出端口以连接到下一个节点。
顶部工具栏
| 按钮 | 动作 |
|---|---|
| 名称·说明 | 编辑流程元信息 |
| 部署/解除 | 立即激活/停用当前流程 |
| 导出 | 将整个图形下载为JSON文件 |
| 测试运行 | 注入任意JSON消息执行一次并确认结果 (参见下方测试运行使用方法) |
| 保存 | 将当前图形保存到服务器。保存时会自动加载快照,之后可回滚 |
若尝试保存没有触发器的图形,会显示"仅可通过手动测试运行启动"的警告。可以有意去掉触发器,将其作为仅手动执行的流程使用。
测试运行使用方法
按下테스트 실행按钮会打开JSON编辑对话框。运营人员可以直接编写消息发布一次,无需等待触发事件即可验证图形动作。
| 项目 | 说明 |
|---|---|
| 编辑器 | 带行号·语法高亮的JSON编辑器。可自由编写消息正文 |
| 验证显示 | 有type必填字段时显示✓ JSON OK (type=X),缺失时显示⚠ type 필수 |
| 发布 | 通过발행按钮将消息注入调度器。结果可在实时调试面板和执行历史中确认 |
基本模板示例
{
"type": "POST_TELEMETRY",
"originator": { "entity_type": "Asset", "id": "MOTOR-001" },
"data": { "speed": 1500, "temp": 75.3 },
"metadata": { "ts": 1746247200000, "site_id": "SITE-01" }
}
有未保存的更改时点击测试运行,会弹出"服务器将执行已保存的版本"的提示。要验证更改需先保存。
左侧 — 节点调色板
可按类别折叠展开,支持通过搜索输入即时过滤。
| 类别 | 颜色 | 节点数 |
|---|---|---|
| 触发器 (Trigger) | 灰色 | 22种 |
| 过滤器 (Filter) | 蓝色 | 6种 |
| 变换 (Transform) | 绿色 | 8种 |
| 动作 (Action) — 集成·存储·资产发布·命令调用 | 橙色 | 7种 |
| 动作 (Action) — 领域 CRUD | 橙色 | 34种 |
| 外部集成 (External) | 紫色 | 9种 |
| 流程控制 (Control) | 灰色 | 7种 |
| 边缘 (Edge) | 青色 | 22种 |
中央 — 画布
| 工具 | 快捷键 / 操作 | 动作 |
|---|---|---|
| 放大/缩小 | Ctrl/⌘ + / Ctrl/⌘ - · 鼠标滚轮 | 画布缩放 |
| 100% | Ctrl/⌘ 0 | 重置缩放 |
| 自适应画面 | Ctrl/⌘ 1 | 自动调整以显示所有节点 |
| 删除所选节点 | Del / Backspace | 删除所选节点/连线 |
| 画布移动 (panning) | 在空白区域左键点击后拖拽 | 在未点击节点的状态下抓住空白画布背景拖动,整个画布会随之移动 |
画布右下角显示小地图,点击小地图可立即跳转到该位置。
💡 panning使用提示: 在节点上拖拽会移动节点 — 要移动画布必须抓住没有节点/连线的空白背景区域。在大型流程中,panning比小地图更快。
右侧 — 节点设置
在画布上点击节点,右侧检查器会显示该节点的设置表单。输入字段会根据节点类型自动生成。
| 输入方式 | 说明 |
|---|---|
| 静态值 | 直接使用表单中输入的值 |
*_field 动态值 | 从消息负载的路径(如data.tag_id、metadata.site_id)提取值。无值时回退到静态值 |
脚本节点(过滤器·变换·switch)可以直接在检查器内用代码编辑器编辑,也可以展开到单独的对话框在更大屏幕上编写。
代码编辑器字体使用可读性更强的等宽字体堆栈(优先Cascadia Code · JetBrains Mono · Consolas · Menlo),中文注释也能稳定对齐。
计数重置
节点设置面板右上角以图标形式显示两个重置按钮。将鼠标悬停在每个按钮上会显示工具提示说明。
| 图标按钮 | 动作 |
|---|---|
| 🩹 (创可贴) — 错误计数重置 | 仅将该流程内所有节点累计的错误计数重置为0 |
| 🔄 (旋转箭头) — 全部计数重置 | 处理·错误·处理时间及流程单位统计全部归零。用于结束运维验证并重新开始统计时使用 |
两个按钮均无需确认对话框即可立即生效 — 仅重置统计,不影响节点动作或消息处理。
右侧 — 实时调试
实时调试面板显示在检查器下方。
| 项目 | 说明 |
|---|---|
| 刷新周期 | 2秒 |
| 级别颜色 | 左侧边框显示INFO(蓝)/WARN(黄)/ERROR(红) |
| 显示信息 | 节点显示名 · 处理时间(ms) · 消息预览 |
| 暂停 | 面板右上角开关,暂停刷新 |
| 清空 | 仅在屏幕上清空累积的调试项目 |
画布上节点右上角显示小指示灯。
| 颜色 | 含义 |
|---|---|
| 灰色 | 等待中 — 未接收消息状态 |
| 绿色 | 消息正在通过 |
| 红色 | 处理中发生错误 |
节点右下角以평균 X · 최근 Y的形式显示处理时间。
消息结构
在流程内部流动的消息由以下4个区域组成。
{
"type": "POST_TELEMETRY",
"originator": {
"entity_type": "Asset",
"id": "MOTOR-001"
},
"data": { "speed": 1500, "temp": 75.3 },
"metadata": {
"ts": 1746247200000,
"site_id": "SITE-01",
"shift": "DAY",
"tag_id": "MOTOR-001.SPEED"
}
}
| 区域 | 含义 |
|---|---|
type | 消息分类。过滤器节点的分支依据 |
originator | 消息主体实体 (针对哪个资产/标签/订单) |
data | 负载正文 |
metadata | 上下文 (时刻·站点·班次·标签ID等) |
消息类型
| 类型 | 触发入口 |
|---|---|
POST_TELEMETRY / TAG_POINT | 标签点接入 |
POST_ATTRIBUTES | 标签/资产元数据更新 |
TAG_ALARM | 标签单位报警 |
ENTITY_CREATED / UPDATED / DELETED | 实体生命周期事件 |
ASSET_DATA / ASSET_EVENT / ASSET_ALARM / ASSET_COMMAND / ASSET_AGGREGATION / ASSET_CONTEXT | 资产领域事件 |
ASSET_HEALTH_STATUS / ASSET_CONNECTION_STATUS | 资产周期评估 (健康/连接状态) |
OEE_EVENT / RAM_EVENT / EMS_EVENT | ISO 分析结果事件 |
OPC_STATUS / EDGE_STATUS | OPC/边缘设备状态 |
DIAGNOSTIC / DOMAIN_CHANGED | 诊断·领域变更 |
ALARM | 报警发生 |
WEBHOOK | HTTP webhook 接收 |
KAFKA_INBOUND / MQTT_INBOUND | 外部主题接收 |
TIMER | 调度触发 |
按触发器的负载示例
编写脚本节点时需要准确知道可以访问哪些字段。以下是各触发器生成的实际消息JSON示例。所有触发器消息通用地包含type、originator、data、metadata。
flow_on_tag_point — 标签点接入
标签每接收到一个值就触发。是最常见的触发器。
{
"type": "POST_TELEMETRY",
"originator": { "entity_type": "Tag", "id": "MOTOR-001.SPEED" },
"data": {
"value": 1500.7,
"quality": "GOOD",
"ts": 1746247200123
},
"metadata": {
"tag_id": "MOTOR-001.SPEED",
"site_id": "SITE-01",
"area_id": "AREA-A",
"line_id": "LINE-1",
"asset_id": "MOTOR-001",
"opc_id": "OPC-LINE-1",
"java_type": "Float",
"unit": "rpm",
"shift": "DAY"
}
}
| 字段 | 含义 | 脚本访问 |
|---|---|---|
data.value | 接收到的值 (数值/字符串/布尔) | msg.data.value |
data.quality | OPC 质量 (GOOD/BAD/UNCERTAIN) | msg.data.quality |
data.ts | 接收时刻 (epoch ms) | msg.data.ts |
metadata.tag_id | 标签 ID | msg.metadata.tag_id |
flow_on_tag_alarm — 标签报警发生
超过标签的报警波段(hi/lo/...)瞬间触发。
{
"type": "TAG_ALARM",
"originator": { "entity_type": "Tag", "id": "MOTOR-001.TEMP" },
"data": {
"alarm_band": "HI_HI",
"value": 95.3,
"threshold": 90.0,
"priority": "ERROR",
"band_message": "온도 위험"
},
"metadata": {
"tag_id": "MOTOR-001.TEMP",
"asset_id": "MOTOR-001",
"site_id": "SITE-01",
"ts": 1746247200123
}
}
data.alarm_band 值 | 含义 |
|---|---|
NORMAL / HI / LO / HI_HI / LO_LO / TRIP_HI / TRIP_LO | 数值型报警级别 |
BOOL_TRUE / BOOL_FALSE | 布尔型报警 |
flow_on_asset_data / flow_on_asset_event — 资产事件
当发生以资产为单位聚合的事件(CEP处理结果·插件评估结果等)时触发。
{
"type": "ASSET_EVENT",
"originator": { "entity_type": "Asset", "id": "MOTOR-001" },
"data": {
"event_type": "STARTUP",
"details": { "rpm_target": 1500 }
},
"metadata": {
"asset_id": "MOTOR-001",
"site_id": "SITE-01",
"ts": 1746247200123
}
}
flow_on_asset_data传递以资产为单位的时序数据(data.values为键值映射),flow_on_asset_aggregation传递分/时单位的聚合值。
flow_on_asset_health_status / flow_on_asset_connection_status — 周期评估
平台每1分钟周期评估资产单位的健康/连接状态。
{
"type": "ASSET_HEALTH_STATUS",
"originator": { "entity_type": "Asset", "id": "MOTOR-001" },
"data": {
"status": "WARN",
"info_count": 12,
"warn_count": 3,
"error_count": 0,
"prev_status": "OK"
},
"metadata": { "asset_id": "MOTOR-001", "ts": 1746247200123 }
}
data.status 值 | 含义 |
|---|---|
OK / WARN / ERROR / UNKNOWN | 健康级别 |
CONNECTED / LATENT / ERROR / DISCONNECTED / UNKNOWN | 连接级别 (connection_status) |
常用的模式是与prev_status比较,仅在状态发生转换的瞬间触发后续动作进行过滤。
flow_on_oee_event / flow_on_ram_event / flow_on_ems_event — 插件事件
工单单位OEE/RAM/EMS评估结果更新时触发。
{
"type": "OEE_EVENT",
"originator": { "entity_type": "WorkOrder", "id": "WO-20260512-001" },
"data": {
"oee": 0.78,
"availability": 0.95,
"performance": 0.85,
"quality": 0.97,
"good_count": 1560,
"bad_count": 42,
"target_count": 2000
},
"metadata": {
"order_id": "WO-20260512-001",
"asset_id": "LINE-1.PRESS",
"shift_id": "DAY-A",
"ts": 1746247200123
}
}
flow_on_opc_status / flow_on_edge_status — OPC/边缘状态
{
"type": "OPC_STATUS",
"originator": { "entity_type": "OPC", "id": "OPC-LINE-1" },
"data": {
"connection_status": "CONNECTED",
"scan_status": "START",
"prev_status": "DISCONNECTED"
},
"metadata": { "opc_id": "OPC-LINE-1", "edge_id": "EDGE-A", "ts": 1746247200123 }
}
flow_on_diagnostic — 系统诊断消息
服务器模块发出的诊断消息进入时触发。
{
"type": "DIAGNOSTIC",
"originator": { "entity_type": "Module", "id": "cep-engine" },
"data": {
"level": "WARN",
"code": "PATTERN_LAG",
"summary": "EQL 패턴 평가 지연 1.2s",
"module": "cep-engine"
},
"metadata": { "ts": 1746247200123 }
}
data.level | 含义 |
|---|---|
INFO / WARN / ERROR | 诊断严重级别 |
flow_on_domain_changed — 领域变更事件
当资产·标签·站点·工单等领域实体发生CRUD操作时触发。
{
"type": "DOMAIN_CHANGED",
"originator": { "entity_type": "Asset", "id": "MOTOR-001" },
"data": {
"action": "UPDATED",
"before": { "asset_name": "Motor1" },
"after": { "asset_name": "Motor 01 - Renamed" },
"changed_by": "admin"
},
"metadata": { "ts": 1746247200123 }
}
flow_on_entity_event — 实体生命周期
以ENTITY_CREATED / ENTITY_UPDATED / ENTITY_DELETED三种类型统一触发。用单个节点即可接收全部三种生命周期。
{
"type": "ENTITY_CREATED",
"originator": { "entity_type": "Customer", "id": "CUST-9001" },
"data": { "customer_name": "신규 고객", "external_id": "ERP-CUST-9001" },
"metadata": { "ts": 1746247200123 }
}
flow_on_webhook — 外部HTTP推送
外部系统发送到POST /api/v4/flow/webhook/{flow_id}的负载会原样转换为消息。通过URL路径的{flow_id}和请求头X-API-Key进行认证。
{
"type": "WEBHOOK",
"originator": { "entity_type": "External", "id": "ERP" },
"data": {
"order_no": "PO-20260512-001",
"customer": "ACME",
"quantity": 1000
},
"metadata": {
"http_method": "POST",
"remote_addr": "10.20.0.55",
"request_id": "req-7c0a...",
"ts": 1746247200123
}
}
外部系统发送的整个JSON正文会原样进入
data。HTTP请求头仅部分(remote_addr/method/request_id)显示在metadata中。
flow_on_mqtt_subscribe — MQTT主题订阅
{
"type": "MQTT_INBOUND",
"originator": { "entity_type": "Topic", "id": "factory/line1/events" },
"data": { "event": "STARTUP", "rpm": 1500 },
"metadata": {
"topic": "factory/line1/events",
"qos": 1,
"broker": "tcp://mqtt.example.com:1883",
"ts": 1746247200123
}
}
flow_jdbc_poll — 外部DB周期轮询
定期执行设定的SELECT查询,每一行都触发一条消息。
{
"type": "KAFKA_INBOUND",
"originator": { "entity_type": "DB", "id": "mes_db" },
"data": {
"PO_NO": "PO-20260512-001",
"CUSTOMER": "ACME",
"QTY": 1000,
"DUE_DATE": "2026-05-20"
},
"metadata": {
"datasource": "mes_db",
"query": "SELECT * FROM po WHERE status='NEW'",
"row_index": 0,
"ts": 1746247200123
}
}
若有100行,则100条消息将依次触发。为避免重复处理同一行,请在SELECT查询中同时更新处理标志,或加入
processed_at列的比较条件。
flow_schedule — 基于时间的触发
以Cron表达式或固定周期触发。负载为空,仅填充metadata.ts。
{
"type": "TIMER",
"originator": { "entity_type": "Schedule", "id": "daily-report" },
"data": {},
"metadata": {
"cron": "0 0 8 * * ?",
"fired_at": 1746247200000,
"ts": 1746247200000
}
}
负载转换时的注意事项
originator.id是领域ID — 若用flow_change_originator节点更改,之后的flow_save_attributes·flow_publish_asset_*动作节点将基于新的originator运行。- 用脚本替换
data/metadata时,建议采用直接赋值而非浅拷贝(Object.assign) — 若更改原始对象,可能影响订阅同一触发器的其他流程。 - 所有epoch时刻字段均为毫秒(ms)。需要秒单位时请
Math.floor(msg.metadata.ts / 1000)。
节点目录
详细的节点ID和选项可在编辑界面的检查器中确认。
触发器 (22种)
所有流程的起点都是触发器节点。没有单独的入口点/终点节点。
领域自动接收 — 系统内部事件会自动分派。
| 类别 | 触发器节点 |
|---|---|
| 标签 | flow_on_tag_point (遥测), flow_on_tag_alarm |
| 资产 | flow_on_asset_data, flow_on_asset_event, flow_on_asset_alarm, flow_on_asset_command, flow_on_asset_aggregation, flow_on_asset_context, flow_on_asset_health_status, flow_on_asset_connection_status |
| 插件 | flow_on_oee_event, flow_on_ram_event, flow_on_ems_event |
| OPC/边缘 | flow_on_opc_status, flow_on_edge_status |
| 诊断·领域 | flow_on_diagnostic, flow_on_domain_changed |
| 实体 | flow_on_entity_event |
flow_on_alarm的节点报警触发器根据对象分为两种 — 标签报警是**flow_on_tag_alarm,
资产报警是flow_on_asset_alarm**。旧文档的示例使用的是flow_on_alarm,
但由于调色板中不存在该名称,照原样使用会找不到该节点。
所有触发器节点都可以用
*_pattern选项(通配符:*、?)进行消息单位的预过滤。未匹配模式的消息不会传递给后续节点,执行计数也不会增加(处理为SKIPPED)。建议在触发器阶段先行过滤以最小化运维负载。
外部进入
| 节点 | 动作 |
|---|---|
flow_on_webhook | 接收外部系统以HTTP推送的负载 |
flow_on_mqtt_subscribe | 订阅外部MQTT代理主题 |
flow_jdbc_poll | 定期读取外部数据库SELECT结果,每行触发一次 |
基于时间
| 节点 | 动作 |
|---|---|
flow_schedule | Cron/周期调度 — 触发TIMER消息 |
过滤器 (6种)
| 节点 | 说明 |
|---|---|
flow_msg_type_filter | 若type包含在指定列表中则TRUE |
flow_originator_type_filter | 若originator.entity_type包含在指定列表中则TRUE |
flow_script_filter | 用脚本进行boolean评估 |
flow_check_existence_field | data/metadata特定字段是否存在 |
flow_switch | 多重case分支 (每个case不同的relation) |
flow_check_relation | 根据上一步骤的relation进行分支 |
变换 (8种)
| 节点 | 界面名称 | 说明 |
|---|---|---|
flow_script_transform | 脚本变换 | 用脚本变换data/metadata |
flow_change_originator | 发送者变更 | 将originator变更为其他实体 |
flow_rename_keys | 键名变更 | 批量变更data字段名 |
flow_template | 模板 | 生成${path}替换文本 |
flow_split | 分割 | 若data是数组,按每个元素分割为消息 |
flow_merge | 合并 | 将多条消息按时间窗口合并 — flow_split的反义词 |
flow_flatten | 平坦化 | 将嵌套对象的子键提升到最上层 (data.data_map.x → data.x) |
flow_to_email | 邮件转换 | 将消息转换为邮件格式 |
flow_flatten用于像标签点这样,希望在data_map内多一层嵌套的值, 让后续节点直接以${data.x}引用时使用。
动作 — 集成·存储 (2种)
| 节点 | 说明 |
|---|---|
flow_save_tag_point | 加载标签点(与正常摄取路径相同) |
flow_save_attributes | 标签/资产元数据部分更新 |
若您在寻找
flow_dds_publish(在内部消息通道自由发布),它不在此处而在 外部集成 调色板组中。
动作 — 发布资产事件 (4种)
以与CEP(EQL)相同的处理路径发布资产事件(永久保存 + 缓存 + 时间线 + 插件 + 消息通道一致处理)。
| 节点 | 通道 |
|---|---|
flow_publish_asset_event | 资产事件 |
flow_publish_asset_context | 资产上下文 |
flow_publish_asset_aggregation | 资产聚合 |
flow_publish_asset_command | 资产命令 |
flow_asset_command_invoke — 不是发布而是执行
| 节点 | 界面名称 | 所做的事 |
|---|---|---|
flow_asset_command_invoke | 调用资产命令 | 同步执行资产中定义的命令并等待结果 |
flow_publish_asset_command混淆flow_publish_asset_command— 将资产命令事件发布到通道。发布后即结束。flow_asset_command_invoke— 实际执行资产中定义的命令并等待接收结果。
如果本想让命令被执行,却使用了发布节点,会看起来什么都没有发生。
动作 — 领域 CRUD (34种)
所有领域操作都委托给同一个领域服务处理,以保持审计·完整性。可用*_field动态选项从消息负载中提取值。
| 领域 | Create | Update | Delete |
|---|---|---|---|
| 资产 (Asset) | flow_create_asset | flow_update_asset | flow_delete_asset |
| 标签 (Tag) | flow_create_tag | flow_update_tag | flow_delete_tag |
| 站点/区域/产线 | flow_create_site | flow_update_site | flow_delete_site |
| 工单 (WorkOrder) | flow_create_work_order | flow_update_work_order | flow_delete_work_order |
| 报警设置 (EQL) | flow_create_alarm_config | flow_update_alarm_config | flow_delete_alarm_config |
| 客户 (Customer) | flow_create_customer | flow_update_customer | flow_delete_customer |
| 产品 (Product) | flow_create_product | flow_update_product | flow_delete_product |
| 作业人员 (Employee) | flow_create_employee | flow_update_employee | flow_delete_employee |
| 日历 (班次) | flow_create_calendar | flow_update_calendar | flow_delete_calendar |
标签报警波段部分更新 (2种)
| 节点 | 说明 |
|---|---|
flow_update_tag_alarm_band_numeric | 仅部分更新已输入的数值型报警波段(hi/lo/hi_hi/lo_lo/trip_hi/trip_lo/band_message/use_alarm)字段 |
flow_update_tag_alarm_band_boolean | 仅部分更新已输入的布尔型报警波段(bool_true/bool_false/优先级/消息/use_alarm)字段 |
工单(WorkOrder)状态转换 (5种)
不直接更新数值列,而是调用领域服务的状态转换方法,因此可保持OEE/RAM/EMS可见性。
| 节点 | 转换 |
|---|---|
flow_start_work_order | WAIT → START |
flow_pause_work_order | START → PAUSED |
flow_resume_work_order | PAUSED → START |
flow_end_work_order | START 或 PAUSED → END |
flow_abort_work_order | START 或 PAUSED → ABORTED (选项 abort_code、notes) |
NOT NULL自动补全 — Create节点会自动为NOT NULL列填充default值。例如:工单
status="WAIT"/master_id是MES主ID(无则自动回填)、客户经理信息"admin"/"admin@example.com"、作业人员org_id回退为site_id、所有行的insert_user_id="flow"。FK列(如customer_id/product_id)中的空字符串会被转换为NULL。
Update节点 — 部分更新 — Customer/Product/Employee/Calendar以及报警波段Update节点会先查询现有记录,再仅合并输入的字段进行保存。空字符串·null值会被忽略,保留原有值。若需要完全覆盖,请使用Delete + Create组合。
有意排除了直接触发报警的节点。报警必须仅通过
flow_create_alarm_config路径产生,才能保持报警历史的一致性。
边缘 (Edge, 22种)
通过调用边缘设备(OPC Agent)的REST API,实现OPC服务器注册·标签CRUD·标签值读写· 查询·Docker应用控制的自动化。所有节点在语义上共享相同的动作,仅默认HTTP方法(GET/POST/PUT/DELETE)不同。
| 组 | 节点 |
|---|---|
| OPC服务器管理 | flow_edge_opc_create, flow_edge_opc_update, flow_edge_opc_delete, flow_edge_opc_start, flow_edge_opc_stop, flow_edge_opc_list |
| 标签管理 | flow_edge_tag_create, flow_edge_tag_update, flow_edge_tag_delete, flow_edge_tag_read, flow_edge_tag_write, flow_edge_tag_list |
| 查询 | flow_edge_monitoring, flow_edge_info, flow_edge_transfer |
| Docker应用控制 | flow_edge_app_list, flow_edge_app_inspect, flow_edge_app_start, flow_edge_app_stop, flow_edge_app_restart, flow_edge_app_logs, flow_edge_app_stats |
查询3种
| 节点 | 界面名称 | 所做的事 |
|---|---|---|
flow_edge_monitoring | 监控 | 查询边缘监控指标 |
flow_edge_info | 边缘信息 | 查询边缘ID · 版本 · uptime |
flow_edge_transfer | 传输健康 | MQTT/Sparkplug传输状态 |
| 节点 | 界面名称 | 所做的事 |
|---|---|---|
flow_edge_opc_list | OPC列表 | 查询边缘的OPC服务器列表 |
flow_edge_tag_list | 标签列表 | 查询OPC的标签列表 |
Docker应用控制7种
将边缘详情界面Docker面板中手动做的事 用流程自动化。
| 节点 | 界面名称 | 所做的事 |
|---|---|---|
flow_edge_app_list | 应用列表 | 边缘Docker容器列表 |
flow_edge_app_inspect | 应用详情 | 容器详情(inspect) |
flow_edge_app_start | 应用启动 | 启动容器 |
flow_edge_app_stop | 应用停止 | 停止容器 |
flow_edge_app_restart | 应用重启 | 重启容器 |
flow_edge_app_logs | 应用日志 | 容器日志 (line=N) |
flow_edge_app_stats | 应用统计 | 容器CPU / 内存 / I/O |
app_stop · app_restart真的会停止并重新启动运行在现场边缘上的容器。
在触发器上挂接以自动执行时,请严格设定条件 — 若在会震荡的报警上挂接
app_restart,容器会反复重启。
通用设置
| 选项 | 说明 |
|---|---|
url | 边缘REST端点。支持${data.x}/${metadata.y}模板替换 |
method | HTTP方法(未设置时为按节点默认值 — 例如create=POST, update=PUT, delete=DELETE, read/monitoring=GET) |
headers | JSON请求头 (例如{"Authorization":"Bearer ${TOKEN}"}) |
body_template | 请求正文(未设置时原样发送data, GET/DELETE不发送正文) |
timeout_ms | 超时(默认5000) |
响应·分支
data.response_status— HTTP状态码data.response— 响应正文(字符串)data.error— 错误消息(失败时)SUCCESS(200~399) /FAILURE(其他或异常)
外部集成 (9种)
所有外部节点都支持*_field动态选项。
| 节点 | 动态选项 |
|---|---|
flow_dds_publish | 在内部消息通道自由发布(虽不是外部调用,但在此组内) |
flow_http_request | url_field / method_field / body_field |
flow_kafka_publish | topic_field / key_field |
flow_mqtt_publish | topic_field |
flow_webhook_callback | url_field |
flow_send_email | to_field / cc_field / subject_field / body_field |
flow_send_sms | to_field / text_field |
flow_send_push | title_field / body_field |
flow_jdbc_query | SQL静态(SELECT/INSERT/UPDATE/DELETE) |
流程控制 (7种)
| 节点 | 说明 |
|---|---|
flow_log | 调试日志 (level / prefix) |
flow_noop | 通过 |
flow_delay | delay_ms后传递给下一个节点 |
flow_throttle | max_msgs / window_ms限制(超过时走THROTTLED relation) |
flow_debounce | window_ms稳定后仅触发最后一条消息 |
flow_merge | 在window_ms期间累积输入以data.merged数组形式一次性emit |
flow_subflow | target_flow_id — 调用其他流程 |
flow_retry | max_attempts(默认3) / backoff_ms(默认1000) / backoff_multiplier(默认2.0)。等待退避后经SUCCESS分支传递消息,达最大尝试次数时走EXHAUSTED分支。自动更新metadata.retry_count / metadata.retry_exhausted |
Retry推荐接线模式
[risky_node] ─[FAILURE]─▶ [flow_retry] ─[SUCCESS]─▶ (다시 risky_node 로 루프 연결)
└[EXHAUSTED]─▶ [에러 핸들러 / 알림]
节点选项详细参考
不用一行表格,而是详细说明复杂节点选项的按选项默认值·示例·失败处理。
flow_script_filter — 脚本过滤器
| 选项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
script | text | (必填) | 求值式 — 返回boolean。true → TRUE分支 / false → FALSE分支 |
script_type | enum | javascript | javascript / eql |
on_error | enum | FALSE | 脚本异常时 — 是否发送到TRUE / FALSE / FAILURE分支 |
脚本上下文
| 变量 | 含义 |
|---|---|
msg.type | 消息类型 (POST_TELEMETRY等) |
msg.data | 负载 (可修改,但在过滤器中无意义) |
msg.metadata | 上下文 |
msg.originator | originator对象 |
示例
// 온도가 임계 초과 + 야간 시프트만
msg.data.value > 80 && msg.metadata.shift === 'NIGHT'
// 사이트별 임계 분기
var th = {'SITE-A': 80, 'SITE-B': 90, 'SITE-C': 75};
msg.data.value > (th[msg.metadata.site_id] || 100);
flow_script_transform — 脚本变换
| 选项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
script | text | (必填) | 变换式 — 修改msg对象或返回新对象 |
script_type | enum | javascript | javascript / eql |
mode | enum | mutate | mutate(就地修改) / return(使用返回值) |
示例
// data 에 계산 필드 추가
msg.data.fahrenheit = msg.data.value * 9/5 + 32;
msg.metadata.processed_at = Date.now();
// 페이로드 통째로 교체 (mode=return)
return {
type: 'WEBHOOK',
originator: msg.originator,
data: { temp: msg.data.value, level: msg.data.value > 80 ? 'HIGH' : 'OK' },
metadata: msg.metadata
};
在
mutate模式下即使有return语句也会被忽略。要替换为新对象需设置为mode=return。
flow_switch — 多重分支
| 选项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
cases | array | (必填) | [{expression, relation}]数组 — 从上到下依次评估,采用第一个匹配 |
default_relation | string | DEFAULT | 所有case均未匹配时 |
script_type | enum | javascript | — |
示例
// cases 설정
[
{ "expression": "msg.data.value > 90", "relation": "CRITICAL" },
{ "expression": "msg.data.value > 80", "relation": "WARN" },
{ "expression": "msg.data.value > 70", "relation": "INFO" }
]
// default_relation: "NORMAL"
在后续节点中可通过4个分支标签(CRITICAL/WARN/INFO/NORMAL)分别进行不同处理。
flow_retry — 自动重试
自动恢复外部IO节点的临时性故障。
| 选项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
max_attempts | int | 3 | 最大尝试次数(超过此值走EXHAUSTED) |
backoff_ms | long | 1000 | 首次等待时间(ms) |
backoff_multiplier | double | 2.0 | 指数退避倍率 — 第1次1秒 → 第2次2秒 → 第3次4秒 |
max_backoff_ms | long | 30000 | 单次等待上限 |
jitter_pct | int | 0 | 在退避上加入±N%随机抖动(避免thundering herd) |
metadata自动补全
| 字段 | 含义 |
|---|---|
metadata.retry_count | 目前已尝试次数 |
metadata.retry_exhausted | 为true时进入EXHAUSTED分支 |
metadata.retry_last_error | 最后一次失败原因 |
退避总和超过60秒可能导致触发处理队列被阻塞。若外部系统响应持续缓慢,请先用
flow_throttle限制进入速度。
flow_on_webhook — HTTP触发器
| 选项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
auth_required | boolean | true | 是否必须携带X-API-Key请求头 — 关闭后任何人都可以调用 |
allowed_origins | csv | * | CORS Origin白名单 |
max_body_kb | int | 256 | 正文大小上限 (超过时返回413) |
payload_pattern | glob | * | 消息预过滤通配符 |
调用方法
curl -X POST \
https://platform.example.com/api/v4/flow/webhook/{flow_id} \
-H "X-API-Key: {edge_or_token_key}" \
-H "Content-Type: application/json" \
-d '{"order_no":"PO-001","customer":"ACME","quantity":1000}'
- 路径中的
{flow_id}可在列表界面中复制 X-API-Key是在边缘API Key或API认证令牌界面发放的令牌- 响应:
200 OK(已enqueue到消息队列) /401(认证失败) /404(无流程或未部署) /413(超出正文大小)
flow_http_request — 外部HTTP调用
| 选项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
url | string | (必填) | 要调用的URL。支持${data.x}模板替换 |
url_field | string | — | 从负载中动态获取URL时 — 如data.endpoint |
method | enum | GET | GET / POST / PUT / DELETE / PATCH |
method_field | string | — | 从负载中动态获取方法 |
headers | json | {} | {"Authorization": "Bearer ${TOKEN}"}格式 |
body | text | — | 静态正文 (支持模板替换) |
body_field | string | — | 从负载中获取正文时 — 通常为data |
timeout_ms | int | 5000 | 响应等待上限 |
follow_redirect | boolean | true | 自动跟随3xx重定向 |
verify_ssl | boolean | true | TLS证书验证 (仅测试用请关闭) |
响应负载补全
| 字段 | 含义 |
|---|---|
data.response_status | HTTP状态码 (200 / 404 / 500 ...) |
data.response_body | 响应正文 (若为JSON则自动解析) |
data.response_headers | 响应请求头对象 |
分支
SUCCESS— 2xx/3xxFAILURE— 4xx/5xx或异常/超时
flow_send_email — 邮件发送
| 选项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
to | string | — | 静态收件人(逗号分隔) |
to_field | string | — | 从负载中提取收件人 — 如data.recipient |
cc / cc_field | string | — | 抄送 |
bcc / bcc_field | string | — | 密抄 |
subject / subject_field | string | (必填其一) | 主题 — 支持模板替换 |
body / body_field | text | (必填其一) | 正文 (允许HTML) |
is_html | boolean | true | 纯文本邮件时请关闭 |
attachments | json | [] | [{"url":"...","filename":"..."}] |
SMTP设置由运营人员预先在系统 → 设置的邮件设置中注册。未注册前,所有邮件节点都会走
FAILURE分支。
flow_jdbc_poll — 外部DB轮询
| 选项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
datasource_id | string | (必填) | 在系统 → 设置中注册的外部DB标识符 |
query | sql | (必填) | SELECT查询 — 一次最多返回1,000行 |
poll_interval_ms | int | 60000 | 轮询周期 (默认1分钟) |
marker_column | string | — | "最后处理时间"列 — 仅SELECT标记之后的行 |
marker_initial | string | 1970-01-01 00:00:00 | 首次轮询时标记的起始值 |
row_limit | int | 1000 | 每次轮询的最大行数 (超过此值也安全) |
on_error_continue | boolean | true | DB错误时仅记录诊断,继续下一次轮询 |
marker_column使用示例
SELECT po_no, customer, qty, created_at
FROM po
WHERE created_at > :marker
ORDER BY created_at
→ :marker处会自动填入上次轮询中获得的最大created_at值。
flow_kafka_publish / flow_mqtt_publish — 外部发布
| 选项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
broker | string | (必填) | kafka:9092 或 tcp://mqtt:1883 |
topic | string | — | 静态主题 — 可进行${data.x}替换 |
topic_field | string | — | 从负载提取主题(如data.target_topic) |
key / key_field | string | — | (仅Kafka) 消息键 |
body / body_field | json/text | (必填其一) | 发布正文 — 未设置时原样使用data |
qos | int | 1 | (仅MQTT) 0/1/2 |
retain | boolean | false | (仅MQTT) Retained标志 |
flow_publish_asset_* — 发布资产事件
资产事件(事件/上下文/聚合/命令)通过CEP/插件/时间线同时识别的统一通道发布。4个节点的通用选项:
| 选项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
asset_id | string | — | 静态资产ID |
asset_id_field | string | — | 从负载提取资产ID (通常为metadata.asset_id) |
event_type | string | — | 资产事件分类 (例如STARTUP、SHUTDOWN、MAINTENANCE) |
event_type_field | string | — | 从负载提取分类 |
payload | json | ${data} | 要发布的正文 — 未设置时原样使用data |
对于
flow_publish_asset_command,会立即传递到资产的命令接收主题(asset_id/cmd/{event_type}),到达边缘。
flow_create_* / flow_update_* — 领域 CRUD
所有Create/Update节点共享以下选项模式。
| 选项 | 类型 | 说明 |
|---|---|---|
{컬럼명} | string | 静态值 (未输入则为NULL/default) |
{컬럼명}_field | string | 从负载提取值 (data.foo / metadata.bar) |
id_strategy | enum | auto(系统颁发) / field(从{도메인}_id_field中提取) |
on_duplicate | enum | error(默认) / skip / update — 仅限Create节点 |
NOT NULL自动补全
Create节点会自动为NOT NULL列填充default值。
| 领域 | 自动填充的列 |
|---|---|
| 工单 | status="WAIT" / master_id(未设置则自动回填) / insert_user_id="flow" |
| 客户 | 经理信息"admin" / "admin@example.com" (若未注册) |
| 作业人员 | org_id=site_id(回退) |
| 通用 | insert_date=now() / insert_user_id="flow" |
FK空字符串处理
FK列(customer_id/product_id等)中输入空字符串""会自动转换为NULL。在JS中,比起delete msg.data.customer_id,msg.data.customer_id = ''更安全。
flow_edge_* — 边缘REST调用
调用边缘设备REST端点的22种节点共享通用选项。
| 选项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
edge_id | string | — | 静态边缘ID |
edge_id_field | string | — | 从负载提取边缘ID |
path | string | (按节点默认) | 边缘REST路径 (例如/api/v1/opc、/api/v1/app/grafana/start) |
body_template | text | — | 请求正文 — 未设置时原样使用data |
timeout_ms | int | 5000 | — |
自动认证
边缘节点在调用时会自动从mm_edge主表中查找该边缘的api_key,并将其附加为X-API-Key请求头。运营人员无需另行设置。
响应负载
| 字段 | 含义 |
|---|---|
data.edge_response_status | 边缘响应码 |
data.edge_response | 响应正文 |
data.edge_id | 调用对象的边缘ID (用于确认) |
分支与flow_http_request相同(SUCCESS / FAILURE)。
脚本节点编写
flow_script_filter / flow_script_transform / flow_switch节点支持两种表达式。
JavaScript (默认推荐)
标准ECMAScript语法。全部支持多行·var/let/const·函数·对象字面量。
绑定
| 变量 | 说明 |
|---|---|
msg | 整条消息。msg.data.x、msg.metadata.topic、msg.type、msg.originator.id均可直接访问 |
data | msg.data的简写别名 |
metadata | msg.metadata的简写别名 |
过滤器示例
data.temp > 80
变换示例
data.temp_f = data.temp * 1.8 + 32;
data.alert = data.temp > 80 ? 'HIGH' : 'OK';
msg
Switch示例 (按case的boolean)
data.t > 100 // case "Critical"
常用变换模式
// 1) 단위 변환 (섭씨 → 화씨) + 라벨링
data.temp_f = data.temp * 1.8 + 32;
data.alert = data.temp > 80 ? 'HIGH' : 'OK';
msg
// 2) 메타데이터 보강 — 시간대·시프트 자동 부여
const h = new Date(metadata.ts).getHours();
metadata.shift = (h >= 6 && h < 18) ? 'DAY' : 'NIGHT';
msg
// 3) 외부 페이로드를 도메인 모델로 매핑 (MES PO → WorkOrder)
const po = data;
data = {
master_id: 'WO-MES-' + po.po_no,
asset_id: po.line_id || 'UNASSIGNED',
title: po.product_name + ' (' + po.qty + ')',
due_date: po.delivery_date,
qty: po.qty
};
msg
// 4) 실패 분기로 명시적 라우팅 (필수 필드 누락 시)
if (!data.tag_id || data.value == null) throw new Error('필수 필드 누락');
msg
// 5) 배열 분할 후 데이터 정제 — split 노드 후에 사용
data.value = parseFloat(data.raw);
data.threshold = data.value > 100;
msg
常用过滤器模式
// 우선순위 화이트리스트
['ERROR', 'CRITICAL'].includes(data.priority)
// 시간대 기반 필터 (주간만 허용)
new Date(metadata.ts).getHours() >= 8 && new Date(metadata.ts).getHours() < 20
// 자산 ID 패턴 매칭
/^MOTOR-.*$/.test(originator.id)
// 임계값 + 안정성 (값이 5번 이상 누적된 경우)
data.value > data.threshold && data.consecutive_count >= 5
用户脚本在安全沙箱中执行,文件·网络·线程·任意类访问均被阻断。若需要调用外部系统,请另行接线
flow_http_request等外部集成节点。
EQL表达式
用于兼容现有EQL用户。仅支持单个表达式(不支持多行·分号分隔),仅可使用带副作用的模式。
#msg.getData().getInt('temp') > 80
如需多个动作,请使用JavaScript。
JavaScript执行环境规范
脚本在隔离的沙箱中执行。准确了解哪些功能可用/不可用,才能编写稳定的脚本。
可用 (✅)
| 功能 | 备注 |
|---|---|
| 标准ECMAScript | var/let/const·函数·类·解构·spread·?.·??等 |
| 对象字面量 | { key: value, ... } |
| 数组方法 | map / filter / reduce / forEach / find / some / every / flat / slice |
| 字符串方法 | split / replace / includes / match / padStart / repeat |
| 数学函数 | 整个Math.* |
| JSON | JSON.parse / JSON.stringify (但若data已是对象,无需再次stringify) |
| Date | new Date() / Date.now() / getHours() / toISOString() 等 |
| 正则表达式 | /pattern/字面量 + RegExp构造函数 |
| 抛出错误 | throw new Error('...') — 自动分支到FAILURE |
| try/catch/finally | 异常处理 |
不可用 (❌)
| 功能 | 原因 / 替代方案 |
|---|---|
网络调用 (fetch / XMLHttpRequest) | 沙箱阻断 — 另行接线flow_http_request节点 |
文件系统 (require('fs')) | 沙箱阻断 |
线程 (setTimeout / setInterval / Worker) | 仅允许同步执行 — 如需延迟,请使用flow_delay节点 |
require / import | 无法加载外部模块 — 所需函数应定义在同一脚本内 |
eval / new Function(string) | 出于安全考虑被阻断 |
| 任意Java类 | 即使在EQL兼容模式下也不会暴露给用户代码 |
process / global / window | 未定义 |
| WebSocket / EventSource | 阻断 — 消息接收由触发器节点负责 |
执行限制 — 没有
旧文档曾刊载执行时间500ms · 内存16MB · 栈1024 · 输出2MB的表格, 但四者均未实际生效。
- 看起来节点设置中可以使用
timeout_ms,但脚本节点不会读取 该值。(ScriptTransformNode·ScriptFilterNode仅读取script和language。) - JS执行上下文中也没有设置时间·语句数·内存限制。
因此像while (true) {}这样的脚本不会被中断,会一直占用工作线程。
流程整体限制的30秒(flow_timeout_ms)是在执行循环取出下一个任务时确认的,
所以脚本内部运行的无限循环无法到达该检查点。
**只写会自行结束的脚本。**避免无限循环·巨量重复·累积巨大字符串, 循环中务必设置上限。
实际被阻断的 — 安全沙箱
与时间·内存不同,触及外部的行为确实被阻断了。
| 项目 | 状态 |
|---|---|
Java类访问 (Java.type等) | 阻断 |
| 文件·网络IO | 阻断 |
| 线程创建 | 阻断 |
| 原生访问 | 阻断 |
| 实验性选项 | 阻断 |
| 宿主对象访问 | 仅限白名单中的内容 |
如需外部调用,请不要在脚本中进行,而是用
flow_http_request等专用节点 分离出来 — 脚本中无论如何都行不通,而专用节点中timeout_ms实际生效。
请不要将繁重处理集中在1个脚本上,而是拆分到多个节点。
msg返回规约
// flow_script_transform 의 두 가지 모드
// 1) mutate 모드 (기본) — msg 객체를 직접 수정, 마지막에 msg 또는 아무 값 반환
data.foo = 'bar';
msg
// 2) return 모드 — 완전히 새 객체로 교체
return {
type: msg.type,
originator: msg.originator,
data: { ...data, foo: 'bar' },
metadata: msg.metadata
};
标准时刻·区域
| 项目 | 动作 |
|---|---|
| 服务器系统时区 | UTC(推荐直接使用epoch ms) |
| 显示韩国时间 | toLocaleString('ko-KR', { timeZone: 'Asia/Seoul' }) |
| 计算韩国时间时差 | new Date(ts + 9*3600000).getUTCHours()或上述区域设置 |
| 0秒单位timestamp | Math.floor(Date.now() / 1000) * 1000 |
内存安全编写模式
脚本节点有16MB内存上限,但在频繁调用的节点中,若小内存泄漏积累,引擎GC负担会加大。请避免以下模式。
反模式1 — 闭包持有大数据
// ❌ 나쁜 예 — 큰 배열을 변환 함수에 캡쳐
const heavy = data.records || []; // 1만 건
const summarize = (item) => heavy.find(r => r.id === item.id);
data.matched = data.targets.map(summarize);
msg
// ✅ 좋은 예 — 인덱스 미리 만들고 함수 안에서만 사용
const index = {};
(data.records || []).forEach(r => { index[r.id] = r; });
data.matched = data.targets.map(t => index[t.id]);
data.records = undefined; // 변환 후 큰 원본 제거
msg
反模式2 — 原样传递大容量数组
// ❌ data.records 가 10,000 건이면 모든 후속 노드에서 메모리 차지
msg
// ✅ 필요한 통계만 남기고 원본 제거
data.summary = {
count: (data.records || []).length,
total: (data.records || []).reduce((s, r) => s + r.value, 0)
};
delete data.records;
msg
反模式3 — 深度对象复制
// ❌ JSON.parse(JSON.stringify(obj)) 는 큰 객체에서 매우 느림
data.copy = JSON.parse(JSON.stringify(data.original));
// ✅ 얕은 복사 또는 필요한 필드만 직접 선택
data.summary = { id: data.original.id, name: data.original.name };
反模式4 — 正则表达式爆炸 (Catastrophic Backtracking)
// ❌ (a+)+b 형태의 중첩 그룹은 입력에 따라 지수 시간
const re = /^(a+)+b$/;
if (re.test(data.text)) ...
// ✅ 비포획 그룹 + atomic 그룹 패턴 또는 단순 매칭
const re = /^a+b$/;
if (re.test(data.text)) ...
反模式5 — 在函数外累积let变量
// ❌ 함수 밖 let 은 매 노드 호출마다 0으로 리셋되지만, 의도와 다르게 동작 가능
let counter = 0;
data.items.forEach(() => counter++);
data.counter = counter;
// ✅ 명시적으로 함수 안에서 선언
data.counter = data.items.length; // 같은 결과, 더 명확
反模式6 — 嵌套try/catch无限尝试
// ❌ 실패해도 다시 던지지 않으면 그래프가 잘못된 분기로 이동
try {
riskyCall();
} catch (e) { /* 무시 */ }
msg
// ✅ 실패 의도면 throw, 정상 의도면 명시적으로 기록
try {
riskyCall();
data.status = 'OK';
} catch (e) {
data.status = 'ERROR';
data.error_msg = e.message;
}
msg
脚本中
throw的异常会走FAILURE分支且消息会被保留。外部IO节点自带独立分支,不要用try/catch包裹。
内存压力诊断
NODE_SCRIPT_OOM这样的代码旧文档曾将NODE_SCRIPT_OOM · FLOW_MSG_TRIMMED · ENGINE_GC_PRESSURE介绍为内存压力
信号,但**产品不会产生这种代码。**脚本内存上限也不存在,
也没有负载自动截断。
内存压力应通过指标和症状判断,而非代码。
| 要看的 | 压力信号 |
|---|---|
FLOW_EXECUTOR_DROPPED_COUNT | 从0开始上升 — 处理池饱和正在丢弃消息 |
FLOW_EXECUTOR_QUEUE_SIZE | 不减少持续增加 |
FLOW_EXECUTION_TIME | 最大值飙升到平时的数倍 |
| 服务器日志 | work_queue full ... task dropped警告、FlowExecutor timeout警告 |
| JVM | GC时间增加 · 堆使用率上升 |
若出现上述信号,请先检查处理大负载的流程, 确认脚本中是否累积了大字符串·数组。
移动推送通知渠道
用flow_send_push节点向运营人员的移动应用发送推送通知。
节点选项
| 选项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
to / to_field | string | (必填其一) | 收件人用户ID (逗号分隔) 或负载路径 |
title / title_field | string | (必填其一) | 通知标题 (建议50字以内) |
body / body_field | text | (必填其一) | 正文 (建议120字以内) |
priority | enum | NORMAL | NORMAL / HIGH — HIGH会在锁屏上也显示 |
sound | enum | default | default / silent / 自定义声音 |
data | json | {} | 应用接收后处理的附加数据 (负载4KB上限) |
deep_link | string | — | 点击通知后打开的界面 (例如pp://asset/MOTOR-001) |
ttl_sec | int | 86400 | 未接收时的保管时间(秒) — 过期后自动删除 |
按用户令牌自动路由
运营人员首次登录移动应用后,设备令牌会自动注册到安全 → API认证令牌。在流程中无需直接处理令牌,只需指定用户ID即可。
- 若iOS/Android两端令牌均已注册,两个设备都会发送
- 若令牌无效(应用已删除等)会自动注销
简单发送示例
[flow_on_asset_alarm]
↓ where priority='ERROR'
[flow_send_push]
to_field: "metadata.responsible_user"
title: "🚨 ${metadata.asset_id} 알람"
body_field: "data.band_message"
priority: HIGH
deep_link: "pp://alarm/view/${data.alarm_id}"
多语言推送
// flow_script_transform — 사용자 로케일에 따라 본문 분기
const locale = metadata.user_locale || 'ko-KR';
const templates = {
'ko-KR': { title: '🚨 ${asset} 위험 알람', body: '값 ${value}, 즉시 점검 바랍니다.' },
'en-US': { title: '🚨 ${asset} Critical Alarm', body: 'Value ${value}, please inspect immediately.' },
'ja-JP': { title: '🚨 ${asset} 危険警報', body: '値 ${value}, 即時点検が必要です。' }
};
const tpl = templates[locale] || templates['ko-KR'];
data.push_title = tpl.title.replace('${asset}', metadata.asset_id).replace('${value}', data.value);
data.push_body = tpl.body.replace('${value}', data.value);
msg
之后在flow_send_push的title_field=data.push_title / body_field=data.push_body中使用。
分组通知合并 (Inbox-style)
将多个报警一次性合并为一个通知(5分钟单位):
[flow_on_asset_alarm]
↓
[flow_merge window=300s]
↓ (data.merged 배열)
[flow_script_transform — 요약 만들기]
↓ data.push_body="알람 N건: A자산, B자산, ..."
[flow_send_push]
用户不在时自动升级
推送未确认30分钟后自动升级为短信或邮件。
[flow_on_asset_alarm priority=ERROR]
↓
[flow_send_push]
↓ SUCCESS
↓ data.alert_id = response_id
[flow_delay 1800s (30분)]
↓
[flow_http_request GET /api/alert/${alert_id}/status]
↓ (response_body.read=false 면)
[flow_script_filter (data.response_body.read === false)]
↓ TRUE
[flow_send_sms]
to_field: "metadata.responsible_phone"
text: "푸시 미확인 30분 경과: ${data.title}"
建议限制发送频率
| 优先级 | 建议频率 |
|---|---|
| HIGH (锁屏显示) | 每用户每小时5件以内 — 防止报警疲劳 |
| NORMAL | 每用户每小时20件以内 |
请在发送前接线flow_throttle节点,或对同一资产的报警使用flow_debounce稳定后再发送。
外部认证令牌自动刷新
自动刷新调用外部系统时会过期的OAuth2等令牌的模式。
简单刷新 — 基于时间(cron)
[flow_schedule cron='0 */50 * * * ?'] ← 50분마다 (만료 1시간 전)
↓
[flow_http_request]
url: "https://auth.example.com/oauth2/token"
method: POST
body: "grant_type=client_credentials&client_id=${creds.client_id}&client_secret=${creds.client_secret}"
↓
[flow_script_transform]
↓ data.access_token = data.response_body.access_token
[flow_save_attributes] ← 자격 증명 저장소에 저장
target_id: "creds:erp_api"
attributes: {"access_token": "${data.access_token}", "expires_at": ${data.response_body.expires_in * 1000 + ts}}
之后在其他流程的flow_http_request中:
Headers: Authorization: Bearer ${creds:erp_api.access_token}
主动刷新 — 收到401时
[main flow]
↓
[flow_http_request]
↓ SUCCESS → 정상 처리
↓ FAILURE
[flow_script_filter (response_status === 401)]
↓ TRUE
[flow_subflow target_flow_id="refresh-token"] ← 토큰 갱신
↓
[원래 노드로 루프] ← 갱신된 토큰으로 재시도
使用Refresh Token
[flow_http_request]
url: "https://auth.example.com/oauth2/token"
method: POST
body: "grant_type=refresh_token&refresh_token=${creds.refresh_token}"
↓
[flow_script_transform]
// 새 access_token + 새 refresh_token (rotation)
data.access_token = data.response_body.access_token;
data.refresh_token = data.response_body.refresh_token;
data.expires_at = Date.now() + data.response_body.expires_in * 1000;
msg
↓
[flow_save_attributes]
令牌过期阈值自动报警
[flow_schedule cron='0 0 * * * ?'] ← 매시
↓
[flow_jdbc_query]
sql: "SELECT id, expires_at FROM credentials WHERE expires_at < NOW() + INTERVAL '1 day'"
↓
[flow_split]
↓
[flow_send_email]
subject: "API 토큰 만료 임박: ${data.id}"
body: "${data.id} 토큰이 ${data.expires_at} 만료 예정입니다."
凭证保管推荐位置
| 种类 | 推荐位置 |
|---|---|
| 静态(几乎不变化) | 系统 → 设置的凭证存储库 |
| 动态(自动刷新) | 使用上述模式,flow_save_attributes存储在资产元数据中 |
| 按用户OAuth | 安全 → API认证令牌界面 |
请勿将凭证以明文形式留在图形中,务必以外部存储引用(
${creds.x})的形式编写。导出/导入时明文令牌也不会包含在JSON中。
用language选项选择使用哪种语言 (JS或EQL)。
执行历史界面
可按时间顺序查询节点单位的执行事件。
检索条件
| 项目 | 说明 |
|---|---|
| 时间 | 指定查询期间 |
| 级别 | 전체 / INFO / WARN / ERROR |
| 流程 ID | 仅过滤特定流程 |
| 消息 ID | 用于追踪单一消息 |
| 数量 | 最近200/500/1,000条 |
快速时间范围
用时间线面板顶部的按钮立即跳转: 10분전 / 30분전 / 1시간전 / 6시간전 / 12시간전 / 전체기간。
结果列
| 列 | 说明 |
|---|---|
| 级别 | INFO / WARN / ERROR |
| 时间 | 事件发生时刻 |
| 事件 | FLOW_START / NODE_IN / NODE_OUT / NODE_ERROR / FLOW_END |
| 流程 ID | 是哪个流程 |
| 节点 | 节点显示名 |
| 节点类型 | 例如flow_script_transform |
| 关系 | SUCCESS / FAILURE / TRUE / FALSE / MATCH / NO_MATCH / DEFAULT / THROTTLED / EXHAUSTED (均为大写) |
| 消息 | 消息类型 · 主体 · 关系 · 处理时间 · 数据预览摘要 |
可用CSV下载按钮导出当前查询结果。
执行历史的保留周期为7天。如需长期保存请导出到外部日志系统。
通用工作流
- 创建新流程 — 列表界面
새 플로우按钮 → 输入名称·说明 - 进入编辑界面 — 自动打开空画布
- 放置触发器节点 — 从左侧调色板拖拽触发器节点
- 添加处理节点 — 按过滤器 → 变换 → 动作的顺序放置并用连线连接
- 节点设置 — 点击各节点在右侧检查器中输入选项
- 保存 — 右上角
저장按钮 (自动加载快照) - 测试运行 — 注入任意消息确认结果
- 部署 — 用
배포开关激活 → 有触发事件进入时自动执行 - 监控 — 通过实时调试面板和执行历史界面确认
示例流程
各示例由节点接线图 + 关键节点设置 + 动作说明组成。JSON格式的节点设置与您在右侧检查器中输入的值1:1对应。
示例1: MES工单自动创建
每分钟从MES轮询新PO并转换为工单。
[flow_schedule: 매분]
│
▼ TIMER
[flow_http_request: MES /api/po/list?status=NEW]
│
▼ SUCCESS
[flow_split: data → 각 PO별 메시지]
│
▼
[flow_script_transform: PO → WorkOrder 매핑]
│
▼
[flow_create_work_order]
│
├── SUCCESS → [flow_log: 워크오더 생성됨]
└── FAILURE → [flow_send_email: 실패 알림]
节点设置
| 节点 | 关键设置 |
|---|---|
flow_schedule | cron: 0 * * * * ? (每分钟0秒) |
flow_http_request | method: GET, url: https://mes.example.com/api/po/list?status=NEW, headers: {"Authorization":"Bearer ${MES_TOKEN}"} |
flow_split | path: data (分割为数组) |
flow_script_transform | language: JS, 参见下方脚本 |
flow_create_work_order | master_id_field: data.master_id, asset_id_field: data.asset_id, title_field: data.title |
flow_send_email | to_field: metadata.alert_to, subject: [MES 동기화 실패] ${data.po} |
Script Transform示例
data.master_id = 'WO-MES-' + data.po;
data.title = data.product_name + ' (' + data.qty + ')';
data.asset_id = data.line_id || 'UNASSIGNED';
data.due_date = data.delivery_date;
metadata.alert_to = 'ops@example.com';
msg
示例2: 通过外部webhook发布资产命令
验证外部系统发送的HTTP负载后转换为内部资产命令。
[flow_on_webhook]
│
▼ WEBHOOK
[flow_script_filter: payload 검증]
│
├── TRUE
▼
[flow_publish_asset_command]
│
├── SUCCESS → [flow_log]
└── FAILURE → [flow_webhook_callback: 외부에 실패 통보]
调用方法: 以POST /flow/webhook/{flow_id}发送JSON正文,会原样载入data区域并触发。
Script Filter示例
// 인증 토큰 일치 + 필수 필드 존재 검사
if (data.token !== 'EXPECTED_TOKEN') return false;
if (!data.asset_id || !data.cmd_key) return false;
true
示例3: 报警 → 自动发布紧急工单
收到ERROR级别以上的报警时,针对上级资产发布紧急维护工单。
[flow_on_tag_alarm]
│
▼ ALARM
[flow_script_filter: priority = ERROR/CRITICAL]
│
├── TRUE
▼
[flow_change_originator: Tag → 상위 Asset]
│
▼
[flow_create_work_order: 긴급 정비]
│
└── SUCCESS → [flow_send_email: 정비 담당자]
Script Filter示例
['ERROR', 'CRITICAL'].includes(data.priority)
示例4: 外部DB同步 — 批量注册资产
轮询遗留DB中的新设备行,自动注册为资产。
[flow_jdbc_poll: SELECT * FROM legacy_assets WHERE sync_status='NEW']
│
▼
[flow_script_transform: 컬럼 매핑]
│
▼
[flow_create_asset]
│
├── SUCCESS → [flow_jdbc_query: UPDATE legacy_assets SET sync_status='OK' WHERE id=?]
└── FAILURE → [flow_log: ERROR + flow_send_push]
节点设置
| 节点 | 关键设置 |
|---|---|
flow_jdbc_poll | dsn: 外部DB连接, sql: SELECT * FROM legacy_assets WHERE sync_status='NEW' LIMIT 100, interval_ms: 60000 |
flow_create_asset | asset_id_field: data.legacy_id, asset_name_field: data.name, site_id_field: data.plant_code |
flow_jdbc_query | sql: UPDATE legacy_assets SET sync_status='OK' WHERE id=?, params_field: data.legacy_id |
示例5: 工单自动状态转换
用资产产生的运转/停止事件自动转换工单状态。
[flow_on_asset_event]
│
▼ ASSET_EVENT
[flow_switch: data.event_type 기준]
│
├── case "RUN" ─▶ [flow_start_work_order] ─▶ [flow_log]
├── case "STOP" ─▶ [flow_pause_work_order] ─▶ [flow_log]
├── case "DONE" ─▶ [flow_end_work_order] ─▶ [flow_log]
└── DEFAULT ─▶ [flow_noop]
Switch节点case示例
RUN:data.event_type === 'RUN'STOP:data.event_type === 'STOP'DONE:data.event_type === 'DONE' && data.qty_done >= data.qty_planned
状态转换失败(错误的当前状态)的情况会自动路由到
FAILURE,因此无需单独的Filter即可安全接线。
示例6: 自动调整报警波段
基于资产聚合(如6小时平均值)动态调整标签的上限/下限报警波段。
[flow_on_asset_aggregation]
│
▼ ASSET_AGGREGATION
[flow_script_transform: 통계 → 임계값 산출]
│
▼
[flow_update_tag_alarm_band_numeric]
│
├── SUCCESS → [flow_log]
└── FAILURE → [flow_send_email: 운영자 알림]
Script Transform示例
// 6h 평균 ± 3σ 를 임계값으로 사용
const mean = data.mean;
const stddev = data.stddev || 1;
data.tag_id = originator.id + '.TEMP';
data.hi = mean + 3 * stddev;
data.lo = mean - 3 * stddev;
data.hi_hi = mean + 4 * stddev;
data.lo_lo = mean - 4 * stddev;
data.use_alarm = true;
msg
示例7: 自动注册边缘设备
收到新OPC服务器信息时批量注册到边缘设备。
[flow_on_webhook] (POST 본문에 OPC + 태그 목록)
│
▼
[flow_edge_opc_create]
│
▼ SUCCESS
[flow_split: data.tags 배열]
│
▼
[flow_edge_tag_create]
│
▼ SUCCESS (모든 태그 등록 완료 후)
[flow_edge_opc_start]
│
└── SUCCESS → [flow_log: 엣지 디바이스 가동]
调用正文示例
{
"type": "WEBHOOK",
"data": {
"opc_id": "OPC-LINE-A",
"endpoint": "opc.tcp://line-a.local:4840",
"tags": [
{ "tag_id": "MOTOR-001.SPEED", "address": "ns=2;s=Motor1.Speed" },
{ "tag_id": "MOTOR-001.TEMP", "address": "ns=2;s=Motor1.Temp" }
]
}
}
示例8: 增强外部API可靠性 (Retry)
对间歇性失败的外部API调用应用退避重试,超过最大尝试次数时通知运营人员。
[flow_on_asset_event]
│
▼
[flow_http_request: 외부 ERP API]
│
├── SUCCESS → [flow_log]
└── FAILURE ─▶ [flow_retry]
│
├── SUCCESS ─▶ (다시 flow_http_request 로 루프 결선)
└── EXHAUSTED ─▶ [flow_send_email: 'ERP 동기화 N회 실패']
Retry节点设置
max_attempts:5backoff_ms:2000backoff_multiplier:2.0→ 间隔2秒、4秒、8秒、16秒、32秒
示例9: 插件事件路由 (OEE低下通知)
OEE值降到阈值以下时向产线管理员发送推送通知。
[flow_on_oee_event]
│
▼
[flow_script_filter: data.availability * data.performance * data.quality < 0.6]
│
├── TRUE
▼
[flow_template: '라인 ${originator.id} OEE ${data.oee_pct}%']
│
▼
[flow_send_push: 라인 매니저]
Template示例
body:라인 ${originator.id} OEE ${data.oee_pct}% (목표 60% 미달, 주요 손실: ${data.top_loss})
示例10: 消息净化管道 (Throttle + Debounce)
将高频标签变化限制为每分钟1次,另外仅在5秒内无变化时才触发后续节点。
[flow_on_tag_point]
│
▼
[flow_throttle: max_msgs=1, window_ms=60000]
│
├── SUCCESS → [flow_debounce: window_ms=5000]
│ │
│ ▼
│ [flow_save_attributes]
└── THROTTLED → [flow_log: 차단됨]
示例11: 多重触发器合计 (Merge)
将OEE/RAM/EMS三种事件以30秒为单位合并为一次报表消息发送。
[flow_on_oee_event] ─┐
[flow_on_ram_event] ─┼─▶ [flow_merge: window_ms=30000]
[flow_on_ems_event] ─┘ │
▼
[flow_script_transform: 요약 메시지 생성]
│
▼
[flow_send_email: 일일 요약]
在同一个流程内放置多个触发器节点时,都会成为入口点。
flow_merge的data.merged数组中会累积窗口期间进入的所有消息。
示例12: 用子流程整合通用处理
将通用消息净化逻辑(重复检查 + 单位转换 + 加载)分离为独立流程,从多个触发器调用。
主流程 (各自独立)
[flow_on_tag_point] ─▶ [flow_subflow: target_flow_id=FLOW_00099]
[flow_on_asset_data] ─▶ [flow_subflow: target_flow_id=FLOW_00099]
子流程 FLOW_00099
(트리거 없음 — 호출 전용)
[flow_log: 'subflow in']
│
▼
[flow_check_existence_field: data.value 존재]
│
├── TRUE
▼
[flow_script_transform: 단위 변환]
│
▼
[flow_save_tag_point]
子流程即使没有触发器也可保存,可用作纯外部调用(保存时会显示警告)。
示例13: 交班时自动发送日报
在每日夜班结束时刻(例如06:00),将昨天~今天的生产·质量·停机摘要发送邮件。
[flow_schedule: 0 0 6 * * ?] (매일 06:00)
│
▼ TIMER
[flow_http_request: 내부 통계 API /api/report/daily]
│
▼ SUCCESS
[flow_script_transform: 본문 마크다운 생성]
│
▼
[flow_send_email]
Script Transform示例
const r = data.report;
data.subject = `[${r.site_id}] ${r.date} 일일 운영 요약`;
data.body =
`■ 생산: ${r.qty_done}/${r.qty_planned} (${(100*r.qty_done/r.qty_planned).toFixed(1)}%)\n` +
`■ 가동률: ${r.availability}%\n` +
`■ 품질률: ${r.quality}%\n` +
`■ 다운타임 Top3:\n` +
r.downtimes.slice(0,3).map(d => ` - ${d.code} ${d.minutes}분`).join('\n');
msg
示例14: 自动汇总停机时间
资产STOP事件进入时更新按班次累计的停机时间,超过阈值时发送通知。
[flow_on_asset_event]
│
▼
[flow_msg_type_filter: type in [ASSET_EVENT]]
│
▼ TRUE
[flow_script_filter: data.event_type === 'STOP']
│
▼ TRUE
[flow_script_transform: 다운타임 분 단위 계산]
│
▼
[flow_save_attributes: 자산 누적 다운타임 갱신]
│
▼ SUCCESS
[flow_script_filter: data.shift_downtime_min > 30]
│
▼ TRUE
[flow_send_push: '시프트 다운타임 30분 초과']
示例15: 质量不良产线自动隔离
质量检查中连续5件以上报告不良时,向产线资产发布停止命令并中止工单。
[flow_on_asset_event] (event_type=QUALITY_FAIL)
│
▼
[flow_script_transform: data.consecutive_fail = (... + 1)]
│
▼
[flow_save_attributes]
│
▼ SUCCESS
[flow_script_filter: data.consecutive_fail >= 5]
│
▼ TRUE
[flow_publish_asset_command: cmd_key='STOP']
│
▼
[flow_abort_work_order: abort_code='QUALITY']
│
└── SUCCESS → [flow_send_email: 품질 매니저 + 라인 매니저]
示例16: 能源超阈值 — 建议产线暂停
产线的小时能耗超出预算时,向运营人员发送带建议信息的短信。
[flow_on_ems_event]
│
▼
[flow_script_filter: data.power_kwh > data.budget_kwh * 1.2]
│
▼ TRUE
[flow_template: '${originator.id} 시간당 ${data.power_kwh}kWh (예산 ${data.budget_kwh}kWh 초과)']
│
▼
[flow_send_sms]
│
└── SUCCESS → [flow_save_attributes: 자산에 마지막 경고 시각 기록]
示例17: 外部系统双向同步 — 工单状态镜像
将来自外部ERP的状态与内部工单状态进行双向同步。
下行 (ERP → 内部)
[flow_on_mqtt_subscribe: erp/work-order/status]
│
▼
[flow_script_transform: 메시지 → 도메인 매핑]
│
▼
[flow_switch: data.status]
│
├── case "STARTED" ─▶ [flow_start_work_order]
├── case "PAUSED" ─▶ [flow_pause_work_order]
├── case "DONE" ─▶ [flow_end_work_order]
└── DEFAULT ─▶ [flow_log: WARN]
上行 (内部 → ERP)
[flow_on_entity_event] (originator.entity_type=Order)
│
▼
[flow_msg_type_filter: type in [ENTITY_UPDATED]]
│
▼ TRUE
[flow_template: ERP 형식으로 변환]
│
▼
[flow_mqtt_publish: erp/work-order/status]
为避免双向同步时出现无限循环,请用如
metadata.source这样的键标记消息来源,并在触发器阶段过滤自身发布的消息。
导入/导出
可将流程定义序列化为JSON,迁移到其他环境或进行备份。
| 动作 | 位置 | 说明 |
|---|---|---|
| 导出 | 编辑界面顶部내보내기按钮 | 将图形 + 节点 + 连线 + 节点设置整体下载为JSON |
| 导入 | 列表界面顶部가져오기按钮 | 粘贴或上传JSON文本 |
| 自动快照 | 保存时自动 | 保存图形时按版本单位加载(用于回滚) |
导入动作
- 总是会颁发新的流程ID(防止覆盖现有ID)
- 节点ID也会重新赋予,连线链接会自动重新映射
- 导入后的流程会以解除状态加载,运营人员需审核后自行部署
错误处理
节点单位
- 节点处理过程中发生异常时会自动路由到
FAILURErelation。 - 若
FAILURE输出没有连接的节点,消息将被drop,仅留下错误日志。 - 触发器节点未匹配
*_pattern的消息会被处理为SKIPPED,不传递给后续节点,也不计入处理/错误计数器。 - 所有错误都会在实时调试面板和执行历史界面中显示。
流程单位
| 项目 | 默认值 | 说明 |
|---|---|---|
| max_depth | 100 | 限制一个消息处理过程中累积的节点访问次数 |
| max_revisit | 3 | 限制对同一节点的再访问次数 (防止循环无限循环) |
| flow_timeout_ms | 30,000 | 处理时间超时时强制终止 |
外部IO节点重试
HTTP·Kafka·外部DB等外部IO节点可通过retry_count / retry_delay_ms选项在节点内部立即重试,所有重试均失败时路由到FAILURE relation。若需要更精细的退避或EXHAUSTED分支处理,请单独使用flow_retry节点。
运维诊断
管理员可在系统菜单中确认流程引擎调度队列状态(worker是否运行、待处理队列大小、累计处理/失败/丢弃数量、最后错误)。
服务器会在后台定期检查流程worker池的负载,若负载持续一定时间以上,会在运维日志中以一行记录状态变化(HEALTHY → DEGRADED → CRITICAL)。恢复正常后还会再留下一行恢复日志。
权限
流程不引入单独的权限模型,直接使用现有系统认证/权限。
| 功能 | 所需权限 |
|---|---|
| 列表查询 / 执行历史查询 | 所有已认证用户 |
| 流程创建·编辑·部署·删除 | ADMIN |
| 导入/导出 / 全部重新部署 | ADMIN |
运维模式 (Recipes)
汇集常用接线形式的库。各模式均可直接复制作为新流程的出发点使用。
模式1 — 处理 + 通知分支
根据处理结果,SUCCESS加载,FAILURE以通知的方式二重分支。
[Action Node]
├── SUCCESS → [flow_log] / [flow_save_attributes] / ...
└── FAILURE → [flow_send_email] / [flow_send_push]
模式2 — 安全重试
配置退避 + EXHAUSTED处理器,使外部IO能应对临时性故障。
[risky_node] ─[FAILURE]→ [flow_retry] ─[SUCCESS]→ (risky_node 로 루프)
└[EXHAUSTED]→ [에러 핸들러]
模式3 — 用预过滤阻断负载
在触发器阶段用*_pattern预先过滤消息,减少后续处理量。未匹配模式的会处理为SKIPPED,不计入计数器。
[flow_on_tag_point] (옵션 tag_id_pattern: "MOTOR-*.SPEED")
│
▼
[필터/변환/액션 ...]
模式4 — 多重case分支
根据状态/类型分为多个分支。
[flow_switch] (case별 boolean 표현식)
├── case A → [...]
├── case B → [...]
└── DEFAULT → [...]
模式5 — 窗口累积 + 一次性emit
将高频输入在一定时间内收集,转换为一条消息。
[High-rate Trigger]
│
▼
[flow_merge: window_ms=10000] → data.merged 배열에 누적
│
▼
[flow_script_transform: 요약]
│
▼
[Action / 외부 발송]
模式6 — 重试 + 超过重试上限后绕行
重试结束后绕行到备份路径(其他API、通知、DB加载)。
[Primary HTTP] ─[FAILURE]→ [flow_retry]
├─[SUCCESS] → (Primary HTTP)
└─[EXHAUSTED] → [Backup HTTP] ─[FAILURE]→ [flow_log/Email]
模式7 — 领域变更后的后续动作
将标签单位报警转换为上级资产单位,委托给按资产的处理。
[flow_on_tag_alarm]
│
▼
[flow_change_originator: Tag → Asset]
│
▼
[자산 단위 액션 (Create Work Order / Publish Asset Event 등)]
模式8 — 结合Throttle与Debounce
限制每分钟不超过1次 + 仅在5秒内无变化时才处理。
[High-rate Trigger]
│
▼
[flow_throttle: max_msgs=1, window_ms=60000]
│
▼ SUCCESS
[flow_debounce: window_ms=5000]
│
▼
[Action]
模式9 — 用子流程模块化通用逻辑
当多个入口点需要共享相同的后续处理(验证·净化·加载)时分离为子流程。
메인1: [Trigger A] → [flow_subflow: target=FLOW_99]
메인2: [Trigger B] → [flow_subflow: target=FLOW_99]
서브 (FLOW_99): (트리거 없음)
[검증] → [정제] → [적재]
模式10 — 直接触发绕行
若要将其他流程的结果流入下一个流程的输入,用flow_dds_publish发布到领域通道,另一个流程以同一消息类型作为触发器接收。
플로우 A: [...] → [flow_dds_publish: type=ASSET_EVENT, originator=...]
플로우 B: [flow_on_asset_event] → [...]
应用示例 (一览)
| 场景 | 构成 |
|---|---|
| MES集成 | flow_schedule → flow_http_request → flow_script_transform → flow_create_work_order |
| 事件外部中继 | flow_on_asset_event → flow_msg_type_filter → flow_mqtt_publish |
| 数据净化·加载 | flow_on_tag_point → flow_script_transform → flow_save_tag_point |
| 报警自动化 | flow_on_tag_alarm → flow_script_filter → flow_send_email + flow_create_work_order |
| 外部DB同步 | flow_jdbc_poll → flow_script_transform → flow_create_asset |
| Webhook接收 | flow_on_webhook → flow_script_filter → flow_publish_asset_command |
| 报警波段自动调整 | flow_on_asset_aggregation → flow_script_transform → flow_update_tag_alarm_band_numeric |
| 工单状态自动化 | flow_on_asset_event → flow_switch → flow_start/end/pause/resume_work_order |
| API可靠性增强 | flow_http_request ─FAILURE→ flow_retry ─SUCCESS→ 循环 / EXHAUSTED→ 通知 |
| OEE低下通知 | flow_on_oee_event → flow_script_filter → flow_template → flow_send_push |
| 边缘批量注册 | flow_on_webhook → flow_edge_opc_create → flow_split → flow_edge_tag_create → flow_edge_opc_start |
| 高频净化 | flow_throttle → flow_debounce → 后续 |
| 多重事件合计 | 多重触发器 → flow_merge → 摘要变换 → 发送 |
分步调试指南
当流程未按预期动作时应遵循的标准流程。
步骤1 — 确认是否触发
在列表界面查看该流程的执行次数是否增加。
| 观察 | 含义·下一步动作 |
|---|---|
| 计数为0 | 未触发。检查流程部署状态 + 触发器节点接线 + *_pattern选项 |
| 计数在增加但错误也一同增加 | 动作节点中失败。转到步骤2 |
| 计数在增加但后续处理未执行 | 分支接线遗漏。检查FAILURE/THROTTLED/EXHAUSTED等所有输出是否已处理 |
步骤2 — 用实时调试确认按节点的流
打开编辑界面,开启实时调试面板(2秒刷新)。
| 观察 | 含义 |
|---|---|
| 特定节点的灯保持灰色 | 消息未到达 — 前一节点走了FAILURE分支或被过滤器阻断 |
| 灯为红色 + ERROR行 | 节点处理中发生异常。检查消息data.error字段和/flow/log的NODE_ERROR事件 |
| 灯为绿色但下一节点为灰色 | 输出relation标签不匹配。确认连线标签是否与节点的输出relation(SUCCESS/TRUE等)完全一致 |
步骤3 — 用执行历史界面追踪消息
在/flow/log界面按消息ID追踪单一消息的整个流程。
- 在消息ID过滤器中输入在实时调试中看到的
msg_id - 按
FLOW_START→NODE_IN→NODE_OUT→FLOW_END顺序按时间排序 - 若出现
NODE_ERROR,则该节点的error_message即为原因
步骤4 — 检查节点设置
常见错误:
| 错误 | 检查 |
|---|---|
*_field路径拼写错误 | 是否正确?是data.tag_id吗?还是metadata.tag_id?用实时调试的消息预览确认实际的键 |
${...}模板变量未替换 | 变量路径是否存在于消息中,是否有拼写错误 |
| 外部IO超时 | 增加timeout_ms。外部系统自身的响应时间 |
| 权限·认证请求头 | headers JSON格式是否正确,令牌是否有效 |
步骤5 — 隔离测试
将问题节点单独移到新的临时流程中,用测试运行触发一条消息,隔离验证结果。若正常动作,则怀疑原流程的接线·上一阶段消息的形式。
步骤6 — 重置计数器后重现
全部计数重置后触发一条消息,用干净的统计重现问题会更清晰。
常见问题
| 症状 | 原因·处理 |
|---|---|
| 触发器未触发 | 确认流程是否处于部署状态,触发器节点的*_pattern是否过窄。调度/外部订阅节点请在全部重新部署后重试 |
| 实时调试为空 | 确认消息计数器是否未增加 — 触发器模式未匹配(SKIPPED)不计入计数器。同时确认面板顶部的暂停开关是否开启 |
| 计数器数值异常累积 | 用节点设置右上角全部计数重置按钮重置统计窗口 |
| 同一消息被重复处理 | 怀疑存在循环。max_revisit默认3次内才允许同一节点再访问。用flow_subflow分离后用flow_throttle/flow_debounce调整输入量 |
| 外部系统调用间歇性失败 | 用retry_count/retry_delay_ms节点选项或flow_retry节点接线退避重试 |
| 导入后调度不动作 | 执行列表顶部的全部重新部署 |
| Update节点执行后部分字段未变 | 这是预期动作。Update节点仅更新输入的字段,其余保持不变。全量更新请使用Delete + Create组合 |
| Create节点因NOT NULL错误失败 | 确认节点设置中*_field动态选项路径是否实际存在于消息中。部分NOT NULL列会自动补全,但核心ID(asset_id、tag_id等)需要直接填写 |
flow_kafka_publish / flow_mqtt_publish未发布 | 确认外部代理连接信息(bootstrap_servers/broker_url)和主题权限。在旁边接线flow_log确认发布前的消息是否到达 |
flow_http_request无响应 | 增加timeout_ms(默认5000),正文直接以body_template中的${msg.data.x}形式明确指定。响应加载到data.response_status / data.response |
flow_retry未被重试 | 是接线错误的情况。需将flow_retry的SUCCESS输出重新循环连接到原失败节点,才会发生重试(参见下方模式2) |
flow_jdbc_poll这样的行每次重新进入 | 轮询SQL的WHERE条件必须包含处理后的状态更新(例如WHERE sync_status='NEW' + 在同一流程末尾用flow_jdbc_query更新UPDATE ... sync_status='OK') |
| 找不到直接触发报警的节点 | 是预期的排除。报警必须仅通过flow_create_alarm_config路径产生,才能保持报警历史一致性 |
脚本节点因Java class访问错误而失败 | 脚本沙箱中被阻断。外部调用请另行接线外部集成节点 |
工单状态转换节点仅出现FAILURE | 当前状态不是可转换的起始状态。例如flow_pause_work_order仅在START状态下动作。请事先用flow_check_existence_field / flow_script_filter确认状态 |
| 测试运行有动作但图形仍是已更改状态 | 测试运行总是触发已保存的版本。要验证更改请先保存 → 测试运行 |
| 导入的流程处于非激活状态 | 是预期动作。审核后请直接打开部署开关 |
常见问题 (FAQ)
汇集了运营人员初次使用流程时常问的问题。
Q. 一个流程可以放置多个触发器吗?
A. 可以。在同一流程内放置多个触发器节点,都会成为入口点,各自独立触发。与flow_merge结合可将多种事件合并为一个后续处理。
Q. 也可以创建没有触发器的流程吗?
A. 可以。保存时会显示警告,但可以保存。可作为仅通过测试运行或其他流程的flow_subflow调用触发的"库"形式使用。
Q. 可以让一条消息同时经过多个节点吗? A. 可以。在一个节点的输出端口连接多条连线,同一消息会同时分支,传递给所有后续节点。
Q. 变换节点可以创建消息本身并转换为其他消息吗?
A. 可以。在flow_script_transform中,msg.type、msg.originator、msg.data、msg.metadata都可以自由更改。但即使让消息本身看起来像其他触发器类型,也不会自动调用其他流程(保留原有的触发器映射)。要调用其他流程,请使用flow_subflow或flow_dds_publish。
Q. 自动快照最多能保存多少个?可以直接恢复吗?
A. 每次保存图形都会自动加载,保留策略遵循运营环境设置。界面上未提供直接恢复的UI,需要时可委托管理员恢复到特定版本。最安全的方法是在变更前用내보내기将JSON文件保管在外部。
Q. 执行历史保存多久?
A. 默认7天。如需长期保存,请用flow_kafka_publish或flow_jdbc_query加载到外部日志系统。
Q. 想单独重置某个流程的统计。 A. 请使用编辑界面节点设置右上角的错误计数重置(仅错误)或全部计数重置(全部)。不会影响其他流程。
Q. 想将流程迁移到其他环境(开发/预发布/生产)。
A. 用编辑界面的내보내기下载JSON,在目标环境的列表界面用가져오기上传即可。导入的流程始终以新ID + 解除状态加载,运营人员审核后可自行部署。
Q. 同一负载被处理两次的情况发生了。如何避免?
A. 在触发器阶段用*_pattern收窄,或用flow_throttle/flow_debounce限制频率。如果是外部输入(webhook·MQTT),可在发送方设置幂等键,并用flow_check_existence_field或基于消息ID的过滤器过滤重复。
Q. 循环接线安全吗?
A. 像flow_retry这样的有意循环是安全的。非此类循环会被max_revisit(默认3次)自动阻断。但运营中建议用flow_log将流程可视化,及早发现无意的循环。
Q. 流程更改会实时生效吗? A. 保存图形时会立即生效。但调度·外部订阅·外部DB轮询等自带调度器的触发器请通过全部重新部署或流程解除→部署开关重新注册。
Q. 可以更改已注册节点的ID吗? A. 节点ID由系统自动分配。显示名可在节点设置表单中更改,实时调试·执行历史中使用显示名。
Q. 没有权限的用户可以误操作更改流程吗? A. 流程创建·编辑·部署·删除需要ADMIN权限。普通用户只能查看列表和执行历史。
运维最佳实践
在生产环境中安全运维流程时应遵循的原则。
负载管理
- 在触发器阶段用
*_pattern预先过滤,减少后续处理量。每分钟接入数万件的标签点是成本最高的入口点。 - 外部系统调用(
flow_http_request/flow_jdbc_query等)与flow_throttle/flow_debounce结合以控制外部负载。 - 若同一输入需要挂接多个后续节点,用
flow_subflow分离以减少整体图形节点数。节点越多实时调试轮询成本也越高。 - 调试模式仅在运维验证阶段开启。加载量越大,执行历史7天窗口越快填满。
安全编写
- 所有动作节点必须接线
FAILURE分支(最起码是flow_log)。若FAILURE为空,消息会被无声drop,导致事故追踪困难。 - 外部IO节点尽可能与
flow_retry结合以吸收临时性故障。 - Update节点仅更新输入字段(部分更新)。全量覆盖请使用Delete + Create组合。
- 无触发器节点的图形仅用于手动测试运行。请确认保存时的警告,避免不慎遗漏触发器。
- 循环(回到同一节点的接线)必须以
flow_retry等有意的形式创建,其他情况虽有max_revisit保护,但对可疑的图形请用flow_log可视化流程。
变更流程 (Change Management)
变更运行中流程的推荐顺序。
- 备份 — 用编辑界面的
내보내기将当前图形下载为JSON文件保管(虽有自动快照,但外部保管更安全)。 - 复制或用新流程作业 — 与其立即修改正在运行的流程,不如在新流程(或导入创建的副本)中修改后验证。
- 测试运行 — 直接注入JSON消息,逐一触发所有分支(SUCCESS/FAILURE/EXHAUSTED等)在执行历史中确认结果。
- 反映到运营 — 将验证过的流程图形JSON
내보내기→ 在运营环境가져오기→ 运营人员审核后部署开关。 - 发生问题时回滚 — 立即用解除开关停用。因已加载自动快照,可委托开发人员恢复到之前版本。
计数器·统计活用
- 若列表界面处理趋势迷你图突然变平,请检查是否为未触发或SKIPPED处理的可能性。
- 若错误率甜甜圈图表异常,请怀疑节点设置的
*_field路径。消息负载中若缺少键则经常失败。 - 运维验证后请用全部计数重置重置统计窗口,重新测量平常的基线。
紧急应对流程
运行中发生问题时应快速使用的分步流程。
场景1 — 特定流程失控(消息暴增)
症状: 某流程的执行次数相比正常值急剧增加数十~数百倍,错误也随之增加
处理
- 立即用该流程的解除开关停用 (在列表界面1秒即可)
- 在实时调试面板·执行历史中确认是哪个触发器引起的失控
- 收窄触发器节点的
*_pattern选项,或紧接着接线flow_throttle/flow_debounce - 必要时用
flow_check_existence_field或flow_msg_type_filter限定消息类型 - 修改后用测试运行重现 → 确认正常后重新部署
场景2 — 外部系统故障导致批量失败
症状: HTTP/Kafka/外部DB等外部IO节点的错误同时累积
处理
- 若影响范围广,批量解除相关流程 (在列表界面批量选择后批量解除)
- 确认外部系统恢复
- 若外部IO节点没有
flow_retry接线,添加接线 - 若外部系统响应变慢,调整
timeout_ms - 根据运营环境全部重新部署后重新部署
场景3 — 循环导致的无限循环
症状: 单个消息持续经过同一节点,处理时间累积
处理
- 虽有
max_revisit保护会自动在最多3次时被阻断,但为了运维安全应先解除后检查 - 在图形中可视化追踪循环 (在编辑界面沿连线追踪)
- 若是有意的循环(
flow_retry),确认max_attempts是否适当 - 若是无意的循环,删除接线或用
flow_check_relation添加分支 - 修改后全部计数重置 → 重新部署
场景4 — 疑似数据损坏 (错误的自动更新)
症状: 自动化导致领域数据与意图不符地被更新
处理
- 立即解除该流程
- 从外部保管的前一图形JSON备份或自动快照中确认之前的版本 (需管理员协助)
- 图形分析: 确认非预期的Update节点接线·错误的
*_field路径·脚本转换错误 - 受影响的领域数据通过单独流程在后台办公界面纠正
- 修改后的图形通过测试运行验证后重新部署
场景5 — 系统检修·维护期间暂停
症状: 想在外部系统检修期间暂时停止外部IO调用
处理
- 批量解除受影响的流程 (从列表中多重选择)
- 检修结束后重新部署
- 自带调度器的触发器(调度·外部MQTT订阅·外部DB轮询)执行一次全部重新部署确认正常注册
场景6 — 实时调试加载队列填满
症状: 运维诊断页面的加载队列大小接近阈值,丢弃数量增加
处理
- 减少开启调试模式的流程数和加载量 (对已验证的流程关闭调试模式)
- 用触发器阶段的
*_pattern减少后续处理量 - 若系统进入自动恢复模式,负载看门狗会留下一行
DEGRADED/CRITICAL日志 — 请与管理员共享
所有场景优先顺序均为立即解除 → 查明原因 → 修改 → 验证 → 重新部署的顺序。也请一并参考变更流程(变更安全流程)。
安全·敏感信息处理
流程涉及外部API调用·邮件/短信发送·webhook接收等敏感输入输出,请遵循以下原则。
令牌·凭证
- 请勿在节点设置中直接明文输入API令牌·密码。使用
${ENV_VAR}形式的环境变量替换,从运营环境的密钥存储中注入。 - 尽可能将运营环境中
headersJSON中令牌的位置分离为环境变量,避免暴露。
// 권장
{ "Authorization": "Bearer ${MES_TOKEN}" }
// 비권장
{ "Authorization": "Bearer eyJhbGciOi..." }
保护Webhook入口点
flow_on_webhook虽是已认证用户调用的内部入口点,但对外公开时务必接线负载验证过滤器。
// flow_script_filter — 토큰 + 필수 필드 검사
if (data.token !== '${WEBHOOK_TOKEN}') return false;
if (!data.asset_id || !data.cmd_key) return false;
true
敏感数据遮蔽
实时调试面板和执行历史中会显示消息预览。请在加载前对个人信息(姓名·联系方式·账户)和令牌进行遮蔽。
// 디버그/로그 적재 직전 변환
data.email_masked = data.email
? data.email.replace(/(.{2}).+(@.+)/, '$1***$2') : null;
delete data.email;
delete data.token;
msg
外部IO的响应正文处理
flow_http_request的data.response会加载整个响应正文。若响应中包含敏感数据,请立即用后续变换节点仅提取所需的键并移除原始数据。
// 응답에서 필요한 필드만 보존
data = { id: data.response_obj.id, status: data.response_obj.status };
msg
脚本节点隔离
脚本节点在安全沙箱中执行,文件·网络·任意类访问被阻断。如需外部调用,请务必接线单独的外部集成节点。
流程指标与报警
流程运行中自动收集的指标以及查看这些指标·设置报警的位置。
按流程的KPI
在列表界面(/flow/index)的各行中显示以下KPI。
| 指标 | 含义 | 异常信号 |
|---|---|---|
| 全部执行次数 | 触发器触发后完整通过图形的次数(成功/失败合计) | 相比平常急剧减少 → 触发器已停止 / 急剧增加 → 失控 |
| 错误次数 | 图形中任何节点走FAILURE分支的次数 | 占全部5%以上需检查 |
| 错误率 | 에러 / 전체 × 100 | 超出一定阈值时应设置报警 |
| 最后执行 | 最近的触发时刻 | "5分钟以上未触发"仅在这是正常情况时才OK |
| 平均处理时间 | 一条消息通过整个图形的平均ms | 添加/删除外部IO节点时变化较大 |
| 最近处理时间趋势 | 5分钟迷你图 — 实时调试面板 | 若出现spike,检查外部系统响应 |
按节点的KPI
在编辑界面点击节点,会在节点右下角显示以下内容。
[Node Name]
처리 999 · 에러 3 · 평균 12ms · 최근 18ms
| 位置 | 显示 |
|---|---|
| 顶部指示灯 | 灰色(等待) / 绿色(处理中) / 红色(错误) |
| 底部标签 | 처리 N · 에러 M · 평균 Xms · 최근 Yms |
节点统计4种
| 计数器 | 含义 |
|---|---|
| 处理 | 进入该节点并走正常分支离开的消息数 |
| 错误 | 走FAILURE分支或发生异常的次数 |
| 平均处理时间 | 节点本身处理ms (包括外部IO) |
| 最近处理时间 | 最后一条消息处理的ms |
系统单位指标
流程引擎整体的指标可在系统 → 监控界面确认。
指标通过JMX域plantpulse.core.engine下以下五个指标公开。
| 指标 | 种类 | 含义 | 查看方法 |
|---|---|---|---|
FLOW_EXECUTOR_QUEUE_SIZE | Gauge | 流程执行池的待处理队列大小 | 若持续堆积,说明处理跟不上接入速度 |
FLOW_EXECUTOR_ACTIVE_COUNT | Gauge | 执行池的活跃线程数 | 若贴合池大小则饱和 |
FLOW_EXECUTOR_DROPPED_COUNT | Gauge (累计) | 因池饱和而被丢弃的task累计数 | 若不为0,说明消息已丢失 — 首先应查看的值 |
FLOW_DEBUG_QUEUE_SIZE | Gauge | 调试事件加载待处理队列大小 | 开启调试模式的流程越多此值越大 |
FLOW_EXECUTION_TIME | Timer | 按流程的处理时间分布 | 追踪平均·最大值 |
flow.engine.*的指标旧文档曾刊载flow.engine.queue_depth · in_flight · dispatch_lag_ms · exec_p95_ms ·
failed_per_min · script_timeout_per_min的表格,但不存在该名称的指标。
以上五个就是全部。
(flow.engine.enabled不是指标,而是engine.properties的设置键 —
是关闭整个流程引擎的主开关。)
基于指标的报警注册模式
用EQL报警或CEP → 触发器监视流程自身指标的4种模式。
模式A — 错误率超阈值
context EVERY_1_MINUTES
SELECT flow_id, count(CASE WHEN status='FAILURE' THEN 1 END) * 100.0 / count(*) AS err_pct
FROM AssetEvent.win:time(5 min)
WHERE event_type = 'FLOW_EXEC'
GROUP BY flow_id
HAVING err_pct > 5
模式B — 处理队列发生背压
当flow.engine.queue_depth > 500持续1分钟以上时报警 — 触发速度超过处理速度的情况。
模式C — 特定流程触发中断
SELECT * FROM pattern [
every a = AssetEvent(event_type='FLOW_EXEC', flow_id='my-flow')
-> ( timer:interval(15 min)
and not AssetEvent(event_type='FLOW_EXEC', flow_id=a.flow_id) )
]
正常情况下每分钟触发N件的流程若沉默15分钟以上则报警。
模式D — 外部IO超时激增
context EVERY_5_MINUTES
SELECT flow_id, node_id, count(*) AS timeout_count
FROM Log.win:time(5 min)
WHERE module = 'flow-engine' AND code = 'NODE_TIMEOUT'
GROUP BY flow_id, node_id
HAVING count(*) > 10
建议报警输出渠道
| 报警种类 | 推荐渠道 |
|---|---|
| 错误率·触发中断 (运营直接) | 邮件 + 推送通知 |
| 队列背压·超时激增 (系统) | Slack/Teams webhook |
| 自动诊断 (参考) | 仅诊断日志 |
⚠️ 报警专用流程本身绝对不要依赖它设置报警的指标 — 若报警流程本身停止,报警本身就不会来。请用系统 → 监控的外部健康检查监视报警专用流程。
消息处理语义与背压
流程引擎处理消息时的保证级别与背压(backpressure)动作。
传递保证 — At-Least-Once
流程引擎保证at-least-once传递。
| 案例 | 动作 |
|---|---|
| 正常处理 | 触发一次 → 通过图形一次 → 完成一次 |
| 处理中引擎重启 | 触发器队列中剩余的消息会在下次启动后重新处理 |
| 节点单位异常 | 仅转换到FAILURE分支。消息不会消失 |
| 外部IO超时 | 可用FAILURE分支 + flow_retry自动重试 |
可能重复 — 不是exactly-once(精确一次)。若
flow_create_*中途中断后重试,同一ID可能进入两次,请使用on_duplicate=skip或外部系统的幂等键。
幂等键模式
从流程调用外部系统时保证幂等性的模式。
// flow_script_transform — 멱등성 키 생성
msg.data.idempotency_key =
msg.metadata.asset_id + '|' +
msg.metadata.ts + '|' +
msg.data.event_type;
之后调用外部时通过请求头附加:
"headers": { "Idempotency-Key": "${data.idempotency_key}" }
处理顺序 — 按触发器单位FIFO
| 同一originator的消息 | 不同originator |
|---|---|
| 触发器单位FIFO(按到达顺序) | 并行处理(顺序无关) |
即使资产A的两个事件几乎同时发生,A的事件也会按发生顺序处理。资产A和资产B的事件可能在不同worker中并行处理。
背压(Backpressure)
处理速度跟不上触发速度时的动作。
| 情况 | 动作 |
|---|---|
| 队列深度 < 80% | 正常 — 立即enqueue新触发 |
| 队列深度80%~100% | 发生诊断WARN日志 — 队列持续接收 |
| 队列深度=100%(已满) | drop新触发 + 发生诊断ERROR日志 |
队列已满时运营人员应做的事
断路器 (外部IO)
外部IO节点(HTTP/Kafka/MQTT/邮件)在满足以下条件时会批量阻断30秒。
| 条件 | 阈值 |
|---|---|
| 最近1分钟内连续失败 | ≥10件 |
| 平均响应时间 | ≥10秒 |
阻断期间进入该节点的所有消息立即分支到FAILURE。外部系统恢复响应后自动解除阻断。这个动作的作用是隔离,使外部故障不会瘫痪整个流程引擎。
断路器开启时刻会记录在诊断日志中,标记为
code = CIRCUIT_OPENED。若外部系统快速恢复,可在执行历史中查看metadata.retry_count,手动重新执行遗漏的消息。
循环(环)防止
若在流程内产生如下循环,可能会发生无限循环。
A → B → C → A (잘못된 결선)
流程引擎用两道防线阻断循环。
| 防线 | 动作 |
|---|---|
max_depth=100 | 经过100个节点则强制终止 |
max_revisit=3 | 第4次访问同一节点瞬间丢弃消息 + 诊断ERROR |
flow_retry的SUCCESS分支 → 回到原节点的循环是有意的循环,因为metadata.retry_count也一同增加,会在max_revisit启动前正常结束。
执行历史中能看到什么
NODE_SCRIPT_TIMEOUT · FLOW_TIMEOUT · FLOW_MAX_DEPTH · TRIGGER_QUEUE_FULL ·
WEBHOOK_AUTH_FAIL · DOMAIN_DUP_KEY这样的代码曾以表格形式刊载,但**产品的任何地方都不会
生成那个字符串。**在日志中搜索该代码永远是0件。
流程引擎实际留下的是以下5种事件类型。
流程执行历史(mm_flow_log,左侧菜单Automation > 流程执行历史)中
加载的事件有五种。
| 事件类型 | 何时 | 加载条件 |
|---|---|---|
FLOW_START | 进入流程时 | 总是 |
FLOW_END | 流程处理结束(成功·超时均包括) | 总是 |
NODE_IN | 节点接收消息时 | 仅在调试模式下 |
NODE_OUT | 节点发出消息时 | 仅在调试模式下 |
NODE_ERROR | 节点处理中异常 | 总是 |
各行同时具有的字段如下。
| 字段 | 内容 |
|---|---|
level | INFO / WARN / ERROR |
flow_id · flow_node_id · node_name · node_type | 是哪个流程的哪个节点 |
msg_id | 消息单位追踪用 — 若想串联一条消息的流程,用此值绑定 |
message · error_message | 人类可读的说明、异常消息 |
duration_ns | 节点执行时间(纳秒) — 仅在NODE_OUT · NODE_ERROR中有意义 |
NODE_IN · NODE_OUT仅在调试模式下加载。平时执行历史中
只显示FLOW_START · FLOW_END是正常的。想逐个节点追踪,
请开启该流程的调试模式 — 但相应地日志量会大幅增加。
从症状出发查找
从症状而非代码出发。
| 症状 | 先查看之处 |
|---|---|
| 触发器触发了,但似乎没有节点运行 | 检查流程是否部署(开关ON) · 触发器的*_pattern是否正在过滤掉消息 |
| 处理中途中断 | 一条消息超过30秒(flow_timeout_ms)。用服务器日志的FlowExecutor timeout警告确认 |
| 同一消息反复处理 | 循环。虽会被max_revisit(3次)自动阻断,但请检查接线 |
| 图形太长无法结束 | 超过max_depth(100)。请分割为子流程 |
| 节点抛出异常 | 在执行历史中查看NODE_ERROR行的error_message |
| 外部调用失败 | 在后续节点中确认该节点的data.response_status · data.error |
限制值(100 · 3 · 30,000ms)是
ExecutionState的默认值。
集群·高可用(HA)动作
平台以集群环境安装时,流程引擎按以下规则进行分布式动作。
节点角色分离
| 角色 | 动作 |
|---|---|
| 主(Leader) | 串行处理流程图形变更(保存/部署/解除)的单一节点 |
| 工作(Worker) | 并行执行触发·消息处理的一般节点 (所有节点均兼任worker角色) |
| 调度器(Scheduler) | 负责flow_schedule cron评估 — 与主节点相同 |
主节点在平台启动时自动选举,若主节点宕机,其他节点中的一个会自动晋升为主节点。运营人员无需直接指定。
消息分配 — 按资产单位一致路由
| 分配键 | 动作 |
|---|---|
originator.id (资产/标签/订单ID) | 同一资产的消息总是路由到同一worker — 保证顺序 |
外部进入 (flow_on_webhook等) | 轮询分配 |
按此规则,防止资产A的事件在多个worker上同时处理导致顺序混乱。结果是按资产单位保证单一worker处理,资产之间并行处理同时成立。
主节点故障时的动作 — 故障转移
[Leader 다운 t=0]
↓
[다른 노드가 리더 승격 시도 t=0~3s]
↓
[새 리더 확정 t=3~5s] ← cron 스케줄·플로우 배포 변경 재개
↓
[기존 워커들은 정상 동작 유지 — 트리거 처리 영향 없음]
| 阶段 | 影响 |
|---|---|
| 0~3秒 | flow_schedule触发暂停 / 触发处理不受影响 |
| 3~5秒 | 确定新主节点,恢复调度 |
| 5秒之后 | 正常 |
cron触发按
misfire策略立即重新执行一次遗漏的触发。运行中若flow_schedule的评估单位小于5秒,可能会出现短期遗漏。
触发队列的持久性
| 项目 | 动作 |
|---|---|
| 触发队列位置 | 内存队列 + 永久存储(事务日志) |
| 节点重启时 | 队列中剩余的未处理消息会在下次启动后重新分派 |
| 处理中节点宕机 | 该消息会被重新处理,但依据at-least-once保证,外部系统需要幂等键 |
集群部署时运营人员检查清单
- 所有节点的时间同步 (NTP) — 触发时刻在节点间保持一致
- 将外部系统(MQTT/Kafka代理)配置在所有节点可访问的网络位置
flow_send_email的SMTP设置只在系统设置中注册一次 — 所有节点共享- 按节点CPU核数调整worker池大小(
flow.engine.workers)
单节点模式 (开发·小规模)
- 主·worker·调度器均在同一进程中运行
- 无故障转移 — 节点宕机时流程本身会中断
- 触发队列由内存+磁盘保证,重启后可恢复
端到端追踪
事后追踪一条消息是如何通过整个图形的方法。
自动赋予关联ID
流程引擎会对所有触发消息自动赋予关联ID(correlation ID)。
{
"type": "POST_TELEMETRY",
...,
"metadata": {
"trace_id": "tr-a1b2c3d4-e5f6-7890-...",
"span_id": "sp-01",
"parent_id": null,
...
}
}
| 字段 | 含义 |
|---|---|
metadata.trace_id | 一次触发触发中唯一的ID — 图形中所有节点通过的消息共享 |
metadata.span_id | 按节点唯一的ID — 每次经过节点都会更新 |
metadata.parent_id | 前一节点的span_id |
在执行历史界面中用trace_id搜索
在执行历史界面中将ID粘贴到trace_id输入框,会显示该消息经过的所有节点的时间顺序日志。
🔍 trace_id = tr-a1b2c3d4-...
[12:34:56.123] [trg] flow_on_tag_point 태그=MOTOR.TEMP, value=87
[12:34:56.125] [filter] flow_script_filter score>80 → TRUE
[12:34:56.126] [transform] flow_change_originator Tag → Asset
[12:34:56.130] [action] flow_create_work_order WO-20260513-001 생성
[12:34:56.241] [external] flow_send_email admin@... 발송 성공
子流程调用时trace的传播
用flow_subflow节点调用其他流程时,同一个trace_id会原样延续。也就是说主流程 + 被调用的子流程所有节点的日志都可以用同一个trace搜索。
传播到外部系统
外部HTTP调用时会自动附加X-Trace-Id请求头。
GET /api/orders HTTP/1.1
Host: erp.example.com
X-Trace-Id: tr-a1b2c3d4-e5f6-...
X-Span-Id: sp-04
外部系统接收该请求头并记录到日志中,两个系统的日志就可以用同一ID匹配。
在报警·邮件中暴露trace_id
在报警正文或邮件模板中包含${metadata.trace_id},运营人员收到报警后可立即在历史界面追踪该事件。
제목: [긴급] 모터 과열 — ${metadata.asset_id}
본문:
시각: ${metadata.ts}
값: ${data.value}°C
trace: ${metadata.trace_id}
이력 보기: https://platform.example.com/flow/log?trace_id=${metadata.trace_id}
追踪保留期
| 数据 | 保留期 |
|---|---|
| trace_id和按节点的span日志 | 7天(与执行历史相同) |
| 外部加载(长期保存) | 使用审计·历史追踪中的外部系统镜像 |
若trace_id太长导致消息大小负担较重,可用
flow_script_transform创建仅暴露最后4位数(如tr-...d4)的短格式用于报警·邮件。但搜索时需要完整ID。
图模式目录
流程图中常用的接线模式10种。
模式1 — Pipeline (简单串行)
[trigger] → [filter] → [transform] → [action]
| 特点 | 消息按顺序通过节点 | | 使用示例 | 超阈值报警 → 发送邮件 |
模式2 — Fan-out (1对N)
┌─→ [action 1]
[trigger] → [t] ─────┼─→ [action 2]
└─→ [action 3]
| 特点 | 一条消息由多个动作同时处理 | | 使用示例 | 报警发生 → 邮件 + 短信 + Slack + 创建工单 |
同一消息的副本会分配到多个节点。各分支独立处理,一个分支的失败不影响其他分支。
模式3 — Fan-in (N对1) — Merge
[trigger A] ──┐
[trigger B] ──┼─→ [flow_merge] → [aggregator] → [action]
[trigger C] ──┘
| 特点 | 将多个触发器的消息在时间窗口内汇总,一次性处理 | | 使用示例 | 将5分钟内发生的所有报警合并为一份日报 |
模式4 — Scatter-Gather (分散 → 收集)
[trigger] → [split] ─┬─→ [process] ──┐
├─→ [process] ──┼─→ [merge] → [action]
└─→ [process] ──┘
| 特点 | 按元素分割数组处理后重新合并结果 | | 使用示例 | 并行验证100件外部订单后一次性上报结果 |
模式5 — Switch (条件分支)
┌─[CRITICAL]→ [긴급 알람]
[trigger] → [flow_switch] ──┼─[WARN] → [경고 알람]
└─[NORMAL] → [통과]
| 特点 | 按条件将一条消息分流到不同路径 | | 使用示例 | 按报警优先级分离处理渠道 |
模式6 — Retry with Fallback
[risky] ─[FAILURE]─→ [flow_retry] ─[SUCCESS]─→ (다시 risky)
└[EXHAUSTED]→ [fallback action]
| 特点 | 自动重试临时故障,永久故障则替代处理 | | 使用示例 | 外部API失败时重试3次,仍失败则通知人员 |
模式7 — Circuit Breaker (活用断路器)
[trigger] → [throttle] → [external_io] ─[SUCCESS]─→ [save]
└[FAILURE]─→ [log only]
| 特点 | 用flow_throttle限制调用速度 + 断路器批量阻断 |
| 使用示例 | 保护外部系统过载 |
模式8 — Dead Letter Queue (DLQ)
[main flow] ─[FAILURE]─→ [flow_save_attributes] → 별도 자산에 적재
↑
운영자가 주기적 점검 후 수동 재처리
| 特点 | 将永久失败的消息移至单独存储 | | 使用示例 | 汇集外部ERP同步失败消息(用于手动重新处理) |
模式9 — Sliding Window Aggregation
[trigger] → [flow_throttle 60s] → [transform: 누적] → [action]
(마지막 60건만 유지)
| 特点 | 以最近N件窗口内的状态做决策 | | 使用示例 | 最近5分钟报警超过30件时进入站点紧急模式 |
模式10 — Saga (多阶段事务)
[start] → [step1] ─OK→ [step2] ─OK→ [step3] ─OK→ [complete]
│ │ │
└─FAIL→[rollback1] │
│ │
┌─FAIL→[rollback1+2]
│
└─FAIL→[rollback1+2+3]
| 特点 | 跨多个外部系统作业的部分失败通过补偿处理 | | 使用示例 | 创建工单 → 预留资产 → ERP同步 → 通知作业人员 (某一阶段失败时取消之前所有阶段) |
模式选择指南
| 需求 | 推荐模式 |
|---|---|
| 简单阈值报警 | Pipeline (1) |
| 一个事件 → 多个渠道 | Fan-out (2) |
| 多个事件 → 一个摘要 | Fan-in / Merge (3) |
| 数组批量处理 | Scatter-Gather (4) |
| 分支处理 | Switch (5) |
| 外部API可靠性 | Retry (6) |
| 保护外部系统 | Circuit (7) |
| 保存失败消息 | DLQ (8) |
| 最近N件累积决策 | Sliding (9) |
| 多个外部系统一致性 | Saga (10) |
流程测试最佳实践
为安全变更·部署流程的测试策略。
3阶段测试 — 单元 → 集成 → 模拟
| 阶段 | 工具 | 验证对象 |
|---|---|---|
| ①单元 | 编辑界面 ▶ 测试运行 | 单个节点独立动作(脚本式、外部调用响应) |
| ②集成 | 同一界面,临时负载 + 实时调试ON | 整个图形序列·分支 |
| ③模拟 | 部署 + 实际触发等待(预发布环境) | 实际消息流程·与外部系统的结合 |
测试运行负载库
可粘贴到测试运行对话框中的按触发器负载示例请参考按触发器的负载示例章节。运营环境中经常发生的边界案例也应提前准备。
边界案例示例
| 案例 | 负载 |
|---|---|
| null字段 | {"value": null} — 检查脚本是否null安全 |
| 空字符串 | {"value": ""} — 验证是否将空字符串解释为0 |
| 负数 | {"value": -1} — 阈值检查是否按预期的±符号 |
| 巨大数字 | {"value": 1e20} — 溢出/精度 |
| Unicode | {"name": "한글-Émoji-🚀"} — 外部系统编码 |
变更安全流程
1. 기존 플로우를 [내보내기] (JSON 파일 저장)
2. 새 플로우를 사본으로 만듦 (이름 끝에 `_v2`)
3. 사본의 트리거 패턴을 좁혀 일부 자산만 매칭 (예: TEST-* 사이트만)
4. 신구 동시 배포 — 새 버전 데이터 확인
5. 1주일 안정성 검증 후 신 버전을 전체 트리거 패턴으로 변경
6. 구 버전 해제 + 보관
保存回归测试场景
建议将常用场景保存为文本文件,变更后务必重新执行。
# regression-tests/alarm-to-workorder.json
{
"case": "고온 알람 → 워크오더 자동 생성",
"input": {
"type": "TAG_ALARM",
"data": { "alarm_band": "HI_HI", "value": 95.0 }, ...
},
"expected": {
"domain_changes": ["WorkOrder.CREATED"],
"notifications": ["email:admin@example.com"]
}
}
模拟外部系统 (Mock) 模式
在预发布环境中模拟外部系统响应时:
| 方法 | 说明 |
|---|---|
| 将flow_http_request的url改为模拟服务器 | 将运营URL与预发布URL变量分离 |
| 更改flow_jdbc_poll的datasource | 仅将运营DB改为预发布DB |
| 将flow_send_email的to_field改为fake@example.com | 防止误发送 |
重置计数器后进行负载测试
性能限制与调优
流程引擎的处理限制与调优技巧。
基本限制
| 项目 | 默认值 | 备注 |
|---|---|---|
一条消息的节点访问数(max_depth) | 100 | 若经过100个节点仍未结束则强制中断 |
同一节点再访问次数(max_revisit) | 3 | 防止循环无限循环 |
流程处理时间(flow_timeout_ms) | 30,000ms | 一条消息处理超过30秒则强制终止 |
| 外部IO节点超时 | 5,000ms(timeout_ms) | HTTP/Kafka/MQTT等 |
| 脚本执行超时 | 500ms(timeout_ms) | 节点单位 |
| 执行历史保留 | 7天 | 之后自动过期 |
常用窗口大小推荐
| 节点 | 推荐窗口 | 备注 |
|---|---|---|
flow_throttle | 1,000~60,000ms | 匹配外部系统API限额 |
flow_debounce | 500~5,000ms | 想让只有稳定值通过时 |
flow_merge | 5,000~60,000ms | 太短会碎片化,太长延迟增加 |
flow_retry退避 | 起始1,000ms × 2倍 | 重试5次则1·2·4·8·16秒 |
增加吞吐量的方法
- 触发器预过滤 — 用
*_pattern仅让所需消息进入(最有效) - 用子流程简化图形 — 主图形仅做分支·路由,繁重处理放到子流程
- 将外部IO接线为异步 — 用
flow_delay/flow_throttle平坦化外部API负载 - 调试模式仅用于验证阶段 — 运营稳定后关闭调试模式
- 仅接线所需类别 — 不必将Edge·外部集成等繁重节点无必要地接线到所有分支
消息大小
消息的data / metadata会经JSON序列化加载到执行历史中。巨大负载(数MB以上的响应正文等)应尽可能用变换节点仅提取所需的键。历史保留·调试显示·子流程调用中消息大小都与处理成本成正比。
审计·历史追踪
追踪流程的变更·执行历史时应确认的位置。
图形变更历史
- 自动快照 — 每次保存图形都会按版本单位加载。若要恢复到之前的状态,请委托管理员。
- 变更者/变更时刻 — 列表界面的最终修改列显示最近保存时刻。变更者按系统认证用户记录。
- 推荐记录变更原因 — 在流程元数据的说明字段中一并记录变更原因·负责人·相关工单号,可提高可追溯性。
执行事件追踪
- 消息ID追踪 — 用执行历史界面的消息ID过滤器可按时间顺序查询单条消息的整个流程(
FLOW_START→ 所有NODE_IN/NODE_OUT→FLOW_END)。 - CSV导出 — 用执行历史界面的CSV下载按钮导出当前查询结果,传递给外部审计系统。
领域变更历史
流程的动作节点(flow_create_* / flow_update_* / flow_delete_*)所做的领域变更均委托给同一个领域服务,因此也会一并记录在后台办公界面的按领域变更历史中。作业人员ID会标记为insert_user_id="flow"等,以区别于一般用户变更。
通过外部加载长期保存
若需要超过执行历史保留期(7天)的长期保存,请创建单独的流程将核心事件加载到外部系统(Kafka·外部DB等)。
[flow_on_tag_alarm]
│
▼
[flow_kafka_publish: topic=audit.alarm.events]
新流程部署检查清单
在生产环境部署新流程前应确认的项目。
图形结构
- 是否恰好接线了一个(或有意的多个)触发器节点
- 是否处理了所有动作节点的
FAILURE输出(至少flow_log) - 若有循环,是否是
flow_retry等有意的接线,是否在max_revisit保护之内 - 外部IO节点是否接线了
retry_count/retry_delay_ms或flow_retry
节点设置
- 触发器的
*_pattern是否过窄或过宽 -
*_field动态选项路径是否实际存在于消息中 - Create节点的NOT NULL字段是否全部填充(除自动补全外)
- Update节点是否确实有意进行部分更新
验证·测试
- 用正常负载通过SUCCESS分支1次
- 用异常负载(必填字段缺失等)通过FAILURE分支
- 用外部IO失败案例通过Retry → EXHAUSTED分支(如适用)
- 实时调试面板的INFO/ERROR显示是否与预期一致
-
/flow/log界面是否记录了所有步骤
运维安全
- 备份(图形JSON导出)是否已保管在外部
- 变更原因·负责人是否已记录在流程说明中
- 若有通知(Email/SMS/Push)接线,收件人是否已验证
- 部署后5~10分钟的执行历史监控计划是否已就绪
部署策略 — Canary / Blue-Green / A·B测试
安全推出新流程的3种策略。所有策略仅用平台的基本功能(触发器模式·消息dispatch·实时调试)即可实现。
策略1 — Canary (仅先应用于部分资产)
先在部分资产上应用新流程,观察一段时间后逐步扩展到全部。
[1단계 출시]
새 플로우 v2 — 트리거 패턴: asset_id LIKE 'LINE-1.%' (1개 라인만)
기존 플로우 v1 — 트리거 패턴: asset_id LIKE 'LINE-2.%' OR 'LINE-3.%' OR ...
[2단계 확대 (1주일 후 안정 확인)]
v2 패턴: 'LINE-1.%' OR 'LINE-2.%'
v1 패턴: 'LINE-3.%' OR 'LINE-4.%'
[3단계 전체 (2주 후)]
v2 패턴: '%' ← 모든 자산
v1 해제·보관
Canary推进检查清单
| 期间 | 监控项目 |
|---|---|
| 第1天 | 错误率 < 1% / 平均处理时间在原有±20%以内 |
| 1周 | 外部系统集成100%正常 / 报警发生频率适当 |
| 2周 | 累积统计 / 用户反馈 / 决定下一产线扩展 |
Canary产线应选择即使受影响也较小的产线(如维护频率高或仅夜间运转的产线)。
策略2 — Blue-Green (新旧同时运营后立即切换)
[Blue (현재)] [Green (새 버전)]
플로우 v1 — 배포됨 플로우 v2 — 배포 + 격리된 originator
모든 트리거 처리 metadata.test_mode=true 인 메시지만 처리
[전환 결정 시점]
v2 트리거 패턴을 v1 과 동일하게 변경 (1초)
v1 해제 토글 (1초)
Blue-Green的核心是在两个版本同时部署的状态下立即切换 — 发现问题时立即重新部署v1进行回滚。
v2隔离接线模式
[trigger 모든 메시지]
↓
[flow_script_filter — metadata.test_mode === true]
↓ TRUE
[새 로직]
测试消息通过/flow/{id}/run API明确指定metadata.test_mode=true发送。
策略3 — A/B测试 (性能·结果比较)
两个版本接收相同消息后执行不同动作,比较结果。
[trigger 메시지]
↓
[flow_split (메시지 복제)]
├ A 경로 → 기존 v1 액션 → [flow_save_attributes target=stats_v1]
└ B 경로 → 새 v2 액션 → [flow_save_attributes target=stats_v2]
之后在日统计界面比较stats_v1 vs stats_v2的累积结果。
A/B比较自动分析
-- EQL 으로 두 버전 비교 (예: 알람 생성 누적)
context EVERY_1_HOURS
SELECT
count(CASE WHEN metadata.flow_version='v1' THEN 1 END) AS v1_count,
count(CASE WHEN metadata.flow_version='v2' THEN 1 END) AS v2_count
FROM AssetAlarm.win:time(1 hour)
策略选择指南
| 情况 | 推荐策略 |
|---|---|
| 首次推出新自动化场景 | Canary — 在一条产线验证后扩展 |
| 现有逻辑大幅变更(结构改编) | Blue-Green — 可立即回滚 |
| 测量两种算法哪个更好 | A/B测试 |
| 简单选项值调整 | 直接变更 — 记录metadata.audit_diff后监控1周 |
回滚流程 (通用)
发现问题时立即回滚:
- 在列表界面将新版本解除开关
- (Canary/A·B的情况)将触发器模式改为0件匹配
- 监控5分钟确认原有版本是否单独正常运行
- 确认诊断日志中没有
code=FLOW_NOT_DEPLOYED - 分析原因 — 在执行历史中用
trace_id追踪失败消息
上线前检查清单 (压缩版)
新流程部署检查清单的核心内容一屏浏览:
[ ] 테스트 실행으로 정상·경계·실패 시나리오 모두 통과
[ ] 외부 IO 노드에 timeout_ms / 재시도 정책 설정됨
[ ] 자격 증명은 ${creds.*} 참조 (평문 미포함)
[ ] 트리거 패턴이 의도한 자산만 매칭
[ ] 영향받는 도메인 (자산/태그/주문) 식별됨
[ ] 운영 시간(특히 야간) 영향 검토됨
[ ] 롤백 시점·기준·담당자 결정됨
[ ] [감사 로그](#audit) 에 변경 의도 메모 작성됨
节点快速设置参考
汇集运营中常用节点关键设置的速查表。
触发器快速设置
| 节点 | 关键选项 |
|---|---|
flow_schedule | cron (例如0 */5 * * * ? = 每5分钟),或interval_ms |
flow_on_webhook | 无选项 — 通过外部POST /flow/webhook/{flow_id}触发 |
flow_on_mqtt_subscribe | broker_url、topic、client_id、username/password |
flow_jdbc_poll | dsn、sql、interval_ms |
flow_on_* (领域) | *_pattern (通配符 — MOTOR-*、SITE-?等) |
变换快速设置
| 节点 | 关键选项 |
|---|---|
flow_script_transform | language(JS/EQL)、script、timeout_ms(默认500) |
flow_change_originator | entity_type、id_field |
flow_rename_keys | mapping(例如{"old":"new"}) |
flow_template | template(${data.x} / ${metadata.y}替换) |
flow_split | path(数组位置,默认data) |
flow_to_email | subject_template、body_template |
流程控制快速设置
| 节点 | 关键选项 |
|---|---|
flow_delay | delay_ms |
flow_throttle | max_msgs、window_ms、(分支: SUCCESS/THROTTLED) |
flow_debounce | window_ms |
flow_merge | window_ms(累积到data.merged数组) |
flow_subflow | target_flow_id |
flow_retry | max_attempts(默认3)、backoff_ms(默认1000)、backoff_multiplier(默认2.0) |
flow_log | level(INFO/WARN/ERROR)、prefix |
flow_noop | (无选项) |
外部集成快速设置
| 节点 | 静态选项 | 动态选项(*_field) |
|---|---|---|
flow_http_request | method、url、headers、body_template、timeout_ms、retry_count、retry_delay_ms | url_field、method_field、body_field |
flow_kafka_publish | bootstrap_servers、topic、value_template、headers | topic_field、key_field |
flow_mqtt_publish | broker_url、topic、qos | topic_field |
flow_webhook_callback | url、method、headers | url_field |
flow_send_email | to、cc、subject、body | to_field、cc_field、subject_field、body_field |
flow_send_sms | to、text | to_field、text_field |
flow_send_push | title、body | title_field、body_field |
flow_jdbc_query | dsn、sql、params_field | — |
动作 — 集成/存储快速设置
| 节点 | 关键选项 |
|---|---|
flow_save_tag_point | tag_id_field(默认metadata.tag_id)、value_field、timestamp_field |
flow_save_attributes | entity_type_field、id_field、attributes_field |
flow_dds_publish | type、originator_field、payload_field |
flow_publish_asset_event | asset_id_field、event_type_field、severity_field、details_field |
flow_publish_asset_command | asset_id_field、tag_id、cmd_key_field、payload_field |
领域CRUD快速设置
| 节点 | 关键选项 |
|---|---|
flow_create_* | 各领域字段 + *_field动态选项。NOT NULL字段部分自动补全(自动补全项目参考领域CRUD表) |
flow_update_* | 部分更新: 仅更新输入字段,空值忽略。全量覆盖请用Delete + Create组合 |
flow_delete_* | *_id_field |
flow_start/end/pause/resume_work_order | order_id_field(默认data.order_id) |
flow_abort_work_order | + abort_code_field、abort_notes_field |
flow_update_tag_alarm_band_numeric | tag_id_field、hi_field等(仅更新输入字段) |
边缘快速设置
边缘节点全部使用相同的选项集。
| 选项 | 说明 |
|---|---|
url | 边缘REST端点(例如http://edge.local:60000/opc/server) |
method | HTTP方法(未设置时为按节点默认值) |
headers | JSON请求头(认证令牌等) |
body_template | 请求正文(未设置时原样发送data) |
timeout_ms | 5000 |
流程 REST API — 通过程序操作流程
在外部自动化工具(Ansible/GitOps/CI管道)中以代码管理流程,或外部系统想立即执行流程时使用的REST API。
认证
所有API调用需将在安全 → API认证令牌发放的令牌作为请求头附加。
Authorization: Bearer {api_token}
端点列表
1) 查询流程列表
GET /flow/list
curl -H "Authorization: Bearer ${TOKEN}" \
https://platform.example.com/flow/list
响应:
{
"data": [
{ "flow_id": "flow-abc123", "flow_name": "MES 동기화", "deployed": true,
"node_count": 12, "exec_count": 9430, "error_count": 2, "last_exec_at": 1746247200000 },
...
]
}
2) 查询单个流程
GET /flow/get/{flow_id}
响应正文包含图形节点·关系·选项全部。可直接用于备份/版本管理。
3) 创建流程
POST /flow/create
Content-Type: application/json
{
"flow_name": "신규 자동화",
"description": "외부 알림 → 워크오더 자동 생성",
"deployed": false
}
从响应中获取flow_id用于后续调用。
4) 修改流程 (元信息)
POST /flow/update
Content-Type: application/json
{
"flow_id": "flow-abc123",
"flow_name": "수정된 이름",
"description": "..."
}
5) 保存图形 (批量替换节点·关系)
POST /flow/{flow_id}/graph
Content-Type: application/json
{
"nodes": [ { "node_id": "n1", "type": "flow_on_tag_point", "options": {...}, "x": 100, "y": 100 }, ... ],
"relations": [ { "from_node_id": "n1", "to_node_id": "n2", "relation": "TRUE" }, ... ]
}
保存图形是事务性的 — 验证失败时整个图形会回滚。
6) 部署 / 解除
更改流程元信息的deployed字段会自动部署·解除。
# 배포
curl -X POST -H "Authorization: Bearer ${TOKEN}" -H "Content-Type: application/json" \
-d '{"flow_id":"flow-abc123","deployed":true}' \
https://platform.example.com/flow/update
7) 立即执行 (手动触发)
POST /flow/{flow_id}/run
Content-Type: application/json
{
"type": "WEBHOOK",
"originator": { "entity_type": "External", "id": "manual-run" },
"data": { "test": true }
}
不经过触发器节点的预过滤,从第一个节点立即执行。用于手动验证·调试很有用。
8) Dispatch (绕过触发器触发)
POST /flow/dispatch
Content-Type: application/json
{
"type": "POST_TELEMETRY",
"originator": { "entity_type": "Tag", "id": "MOTOR-001.SPEED" },
"data": { "value": 1500 },
"metadata": { "tag_id": "MOTOR-001.SPEED", "ts": 1746247200000 }
}
通过正常的触发器处理路径进入消息 — 匹配该消息的所有流程会同时触发。
9) 重置统计
POST /flow/{flow_id}/stats/reset-errors — 에러 카운트만 초기화
POST /flow/{flow_id}/stats/reset-all — 전체 카운트 초기화
10) Export / Import — 图形备份
# Export — JSON 파일로 다운로드
curl -H "Authorization: Bearer ${TOKEN}" \
https://platform.example.com/flow/${FLOW_ID}/export \
-o flow-backup.json
# Import — 같은 JSON 을 다른 환경에 등록
curl -X POST -H "Authorization: Bearer ${TOKEN}" -H "Content-Type: application/json" \
--data @flow-backup.json \
https://platform.example.com/flow/import
导入时若存在冲突ID会自动颁发新ID。外部依赖(凭证·资产ID等)导入后请另行映射。
11) 重新部署所有流程
POST /flow/redeploy
用于大规模变更后或节点重启后的批量同步。运行中调用需谨慎 — 会发生暂时性处理延迟。
12) 自我诊断
POST /flow/selftest
对流程引擎的内部组件(触发队列·调度器·节点注册表·脚本运行时)进行一次检查并返回状态。结果也会记录在诊断日志中。
13) 查询节点目录
GET /flow/catalog
返回当前注册的所有节点类型·选项模式。用于UI动态生成检查器。外部工具自动生成图形时也可参考。
用Webhook触发器从外部触发
含有flow_on_webhook触发器的流程可让外部系统通过以下URL直接触发。
POST /flow/webhook/{flow_id}
Authorization: Bearer {token} 또는 X-API-Key: {token}
Content-Type: application/json
{
"order_no": "PO-001",
"customer": "ACME",
"quantity": 1000
}
响应:
{ "status": "ACCEPTED", "trace_id": "tr-..." }
| 响应码 | 含义 |
|---|---|
| 200 | 已enqueue到消息队列(实际处理结果异步) |
| 401 | 认证令牌错误/缺失 |
| 404 | 流程ID不存在或未部署 |
| 413 | 正文超过256KB |
| 422 | 触发器节点不是flow_on_webhook |
| 429 | 超出每分钟调用限额 |
外部自动化工具集成示例
GitHub Actions — PR合并时部署流程
- name: 플로우 배포
run: |
curl -X POST \
-H "Authorization: Bearer ${{ secrets.PP_TOKEN }}" \
-H "Content-Type: application/json" \
--data @flows/mes-sync.json \
https://platform.example.com/flow/import
Ansible — 批量管理流程
- name: 모든 플로우 백업
uri:
url: "https://platform.example.com/flow/{{ item }}/export"
headers:
Authorization: "Bearer {{ pp_token }}"
dest: "/backups/{{ item }}.json"
loop: "{{ pp_flow_ids }}"
Jenkins — 每次新构建时进行selftest
stage('PlantPulse Flow Selftest') {
steps {
sh """
curl -X POST -H 'Authorization: Bearer ${PP_TOKEN}' \\
https://platform.example.com/flow/selftest \\
--fail-with-body
"""
}
}
API速率限制
| 端点 | 每分钟限额(按令牌) |
|---|---|
GET /flow/list · get · catalog | 600 |
POST /flow/create · update · delete | 60 |
POST /flow/{id}/run · dispatch · webhook/{id} | 1,000 |
POST /flow/redeploy · selftest | 10 |
超出限额时返回429 Too Many Requests响应。
OPC/PLC 工业集成模式
工业现场特有集成场景汇集。22种边缘节点和4种发布资产事件的组合运用。
模式A — OPC服务器自动注册
新产线投产时,从ERP获取产线信息后自动在边缘注册OPC服务器。
[flow_on_webhook]
↓ (data = {"line_id":"LINE-7","host":"10.0.7.10","port":4840})
[flow_change_originator]
↓ (originator → Edge:EDGE-A)
[flow_edge_opc_create]
↓ (path=/api/v1/opc, body={"opc_id":"OPC-${data.line_id}","host":"${data.host}","port":${data.port}})
[flow_edge_opc_start]
↓ (수집 시작 명령)
[flow_save_attributes]
↓ (등록 이력을 자산 메타에 기록)
边缘节点会自动从
mm_edge主表中查找api_key。运营人员无需另行进行认证设置。
模式B — PLC命令自动发布
资产报警发生时自动向PLC发布停止命令。
[flow_on_asset_alarm]
↓ where priority='ERROR' AND alarm_band='TRIP_HI'
[flow_publish_asset_command]
↓ (event_type='STOP', payload={"reason":"trip_high","triggered_by":"flow"})
[flow_log] (감사 로그)
flow_publish_asset_command立即传递到资产的命令主题(asset/{id}/cmd/STOP),边缘对PLC执行OPC Write。
模式C — OPC标签批量注册 (CSV/Excel接入)
新设备设置时一次性注册CSV中的100~1000件标签。
[flow_jdbc_poll] ← 외부 DB 에 적재된 CSV 행
↓ (한 행 = 한 메시지)
[flow_script_transform] ← 태그 정의 가공
↓
[flow_create_tag] ← 도메인 등록 (NOT NULL 자동 보강)
↓
[flow_edge_tag_create] ← 엣지에 동기화
↓
[flow_log] (성공 카운트)
注册1,000件时建议用
flow_throttle平坦化到每分钟100件以下。边缘可能会暂时响应变慢。
模式D — OPC断连自动恢复
OPC连接断开后恢复时自动重启。
[flow_on_opc_status]
↓ where prev_status='CONNECTED' AND status='DISCONNECTED'
[flow_delay 30s] ← 잠시 안정화 대기
↓
[flow_edge_opc_start] ← 수집 재시작 시도
↓ SUCCESS: 정상
└ FAILURE: [flow_retry max=3] → EXHAUSTED: [관리자 알림]
模式E — 交班开始时产线状态检查
在每次交班开始时刻对站点内所有产线进行巡检。
[flow_schedule cron='0 0 7,15,23 * * ?'] ← 7시·15시·23시
↓
[flow_edge_monitoring] ← 엣지에서 OPC 상태 전체 조회
↓ data.opcs = [...]
[flow_split] ← 배열을 행별 메시지로
↓
[flow_script_filter] ← status != 'CONNECTED' 만
↓
[flow_send_email] ← 비정상 OPC 만 한 통의 메일에 합쳐 발송
模式F — 交班时自动映射作业人员
根据排班表变更自动向工单映射作业人员。
[flow_on_entity_event type='ENTITY_UPDATED' Calendar]
↓
[flow_script_filter (시프트 시작 시각이 지금인 경우만)]
↓
[flow_jdbc_query (그 시프트의 작업자 목록 조회)]
↓ data.employees=[...]
[flow_split]
↓
[flow_update_work_order (작업자 매핑)]
模式G — 多重PLC同时调用 (Scatter-Gather)
同时调用多个PLC的数据后整合结果。
[flow_on_webhook]
↓ data.targets=['PLC-1','PLC-2','PLC-3']
[flow_split]
↓ (3개로 분할)
[flow_edge_tag_read]
↓ (병렬 호출)
[flow_merge window=5s]
↓ data.merged=[...]
[flow_script_transform (data.summary 만들기)]
↓
[flow_publish_asset_aggregation]
模式H — 自动调整报警波段
将预测分析学习到的LIMIT_MIN/MAX值自动应用到报警波段。
[flow_schedule cron='0 0 4 * * MON'] ← 매주 월요일 4시
↓
[flow_jdbc_query (forecast 결과 조회)]
↓ data.tags=[{tag_id, limit_min, limit_max}]
[flow_split]
↓
[flow_update_tag_alarm_band_numeric] ← lo=limit_min, hi=limit_max
↓ SUCCESS
[flow_log] (적용 결과 기록)
模式I — 停机自动分类
按原因分类资产停止事件,准确反映到OEE可用性。
[flow_on_asset_event event_type='SHUTDOWN']
↓
[flow_switch]
├ CASE 점심: hour_of_day(ts) BETWEEN 12 AND 13 → [flow_create_calendar (계획정지)]
├ CASE 시프트 종료: 시프트 종료 시각 ±5분 → [flow_log (정상 종료)]
├ CASE 정기점검: 마지막 정비일 +30일 경과 → [flow_create_calendar (정기점검)]
└ DEFAULT (비계획) → [flow_send_email (긴급)]
模式J — 与外部ERP双向工单同步
PlantPulse ↔ 外部ERP双向镜像。
입력 1: ERP → 플랜트펄스
[flow_on_webhook] → [flow_create_work_order]
입력 2: 플랜트펄스 → ERP
[flow_on_entity_event type='ENTITY_CREATED' WorkOrder]
→ [flow_script_filter (출처가 ERP 가 아닌 경우만)]
→ [flow_http_request (ERP REST PUT)]
为防止无限循环,请对从外部接入的消息标记
metadata.source='ERP',反方向流程中过滤该消息。
工业集成时注意事项
| 事项 | 建议 |
|---|---|
| OPC Write命令应在安全联锁后发布 | 仅在操作员确认或自动安全规则通过后 |
| 请勿将PLC轮询周期设置得太短 | 1秒以下有PLC CPU占用率急剧上升风险 |
| 边缘设备Docker容器资源限额 | OPC采集+额外容器时注意CPU/内存80%以上 |
| 试运行阶段将所有命令设为dry-run模式 | 活用flow_edge_tag_write的dry_run=true选项 |
| 应对停电 — 仅信任永久存储触发器 | flow_on_webhook等是易失性的,领域事件触发器为永久 |
外部系统集成 Cookbook
汇集常调用的外部系统的flow_http_request设置和负载示例。
Slack — Incoming Webhook通知
URL: https://hooks.slack.com/services/T0000/B0000/{secret}
Method: POST
Headers:
Content-Type: application/json
Body (body_template):
{
"text": "${data.title}",
"blocks": [
{ "type": "header", "text": { "type": "plain_text", "text": "${data.title}" } },
{ "type": "section", "text": { "type": "mrkdwn", "text": "*자산:* ${metadata.asset_id}\n*값:* ${data.value}\n*시각:* ${metadata.ts}" } }
]
}
- 成功: 200 +
ok正文 - 失败: 400 (错误的payload) / 404 (错误的secret) / 429 (超出每分钟调用限额)
- 建议: 每分钟1次以下接线
flow_throttle
Microsoft Teams — Webhook通知
URL: https://{org}.webhook.office.com/webhookb2/{id}/IncomingWebhook/{secret}
Method: POST
Headers:
Content-Type: application/json
Body:
{
"@type": "MessageCard",
"@context": "https://schema.org/extensions",
"themeColor": "FF0000",
"title": "${data.title}",
"sections": [{
"facts": [
{ "name": "자산", "value": "${metadata.asset_id}" },
{ "name": "값", "value": "${data.value}" },
{ "name": "시각", "value": "${metadata.ts}" }
]
}]
}
将
themeColor按FF0000(红)/FFA500(橙)/00C853(绿)分支以可视化优先级。
Jira — 自动创建工单
URL: https://{org}.atlassian.net/rest/api/3/issue
Method: POST
Headers:
Authorization: Basic {base64(email:apitoken)}
Content-Type: application/json
Body:
{
"fields": {
"project": { "key": "OPS" },
"summary": "[자동] ${data.title}",
"description": {
"type": "doc",
"version": 1,
"content": [{
"type": "paragraph",
"content": [{ "type": "text", "text": "자산: ${metadata.asset_id}\n값: ${data.value}" }]
}]
},
"issuetype": { "name": "Bug" },
"priority": { "name": "High" }
}
}
- 响应正文
data.response_body.key中会带有OPS-1234这样的问题键,用于后续节点 - 失败: 400(错误的字段) / 401(认证) / 403(无项目权限)
GitHub Issues — 自动创建issue
URL: https://api.github.com/repos/{owner}/{repo}/issues
Method: POST
Headers:
Authorization: Bearer {pat_token}
Accept: application/vnd.github+json
Body:
{
"title": "[자동] ${data.title}",
"body": "자산: ${metadata.asset_id}\n값: ${data.value}\n시각: ${metadata.ts}\ntrace: ${metadata.trace_id}",
"labels": ["automation", "ops"]
}
ERP (SAP S/4HANA OData) — 发布工单
URL: https://{host}/sap/opu/odata/sap/API_MAINTNOTIFICATION/MaintenanceNotification
Method: POST
Headers:
Authorization: Basic {base64(user:pass)}
X-CSRF-Token: fetch ← 별도 GET 으로 토큰 받기
Content-Type: application/json
Body:
{
"NotificationType": "M2",
"MaintenanceNotificationType": "M2",
"TechnicalObject": "${metadata.asset_id}",
"NotificationText": "${data.title}",
"MalfunctionStartDate": "${data.start_date}",
"Priority": "${data.priority}"
}
需先用GET获取CSRF令牌,然后用同一会话cookie进行POST。请连接两个
flow_http_request节点,将第一个节点响应头中的令牌传播到第二个节点的请求头。
MES (外部DB直接INSERT) — flow_jdbc_query
直接向外部MES系统的工单表插入行。
INSERT INTO mes.work_order
(order_no, asset_id, product_id, qty, status, due_date, created_by)
VALUES
(:order_no, :asset_id, :product_id, :qty, 'NEW', :due_date, 'flow')
| 节点选项 | 值 |
|---|---|
datasource_id | mes_db (在系统设置中预先注册) |
query | 上述SQL |
binds | {"order_no":"${data.order_no}","asset_id":"${metadata.asset_id}",...} |
flow_jdbc_query除SELECT外也支持INSERT/UPDATE/DELETE。事务按节点单位自动commit/rollback。
Grafana — 自动更新仪表盘通知
URL: https://{host}/api/annotations
Method: POST
Headers:
Authorization: Bearer {api_key}
Content-Type: application/json
Body:
{
"dashboardUID": "...",
"panelId": 4,
"time": ${metadata.ts},
"tags": ["alarm", "${metadata.asset_id}"],
"text": "${data.title}"
}
报警发生时会在图表上自动显示标记 — 事后分析非常有用。
Telegram — 机器人消息
URL: https://api.telegram.org/bot{token}/sendMessage
Method: POST
Body:
{
"chat_id": "-1001234567890",
"text": "🚨 *${data.title}*\n자산: `${metadata.asset_id}`\n값: ${data.value}",
"parse_mode": "Markdown"
}
REST认证模式速查表
| 认证方式 | 请求头 |
|---|---|
| API Key (请求头) | X-API-Key: {token} |
| API Key (查询) | URL末尾加上?api_key={token} |
| Bearer Token | Authorization: Bearer {token} |
| Basic Auth | Authorization: Basic {base64(user:pass)} |
| OAuth2 (Client Credentials) | 发放令牌 → Bearer请求头(用单独节点刷新) |
| HMAC签名 | 用脚本计算基于正文/时刻的签名,加入X-Signature请求头 |
响应解析模式
flow_http_request响应后在下一节点中活用data.response_body:
// 응답에서 ID 추출
data.created_id = data.response_body.id;
data.status = data.response_body.status;
// 응답 배열에서 첫 행만
data.first_item = (data.response_body.items || [])[0] || null;
// 응답이 문자열이면 JSON 파싱
if (typeof data.response_body === 'string') {
try { data.response_body = JSON.parse(data.response_body); } catch (e) {}
}
msg
外部系统调用时的安全建议
| 事项 | 建议 |
|---|---|
| 请勿将令牌以明文留在图形中 | 注册到系统设置 → 凭证存储库后以${creds.slack_webhook}形式引用 |
| 若响应正文含敏感信息则遮蔽 | 用flow_script_transform移除PII/令牌后传给下一节点 |
| 每个外部系统单独的retry策略 | 调用频率高的地方另行接线flow_retry + flow_throttle |
| 将调用结果保存为审计日志 | 用flow_save_attributes向资产添加调用历史元数据 |
数据变换 Cookbook
可直接搬到flow_script_transform节点中使用的变换模式汇集。所有示例均以JavaScript为基准,使用data / metadata简写别名。
单位转换
// 섭씨 → 화씨
data.temp_f = data.temp_c * 9 / 5 + 32;
// 바 → kPa
data.pressure_kpa = data.pressure_bar * 100;
// rpm → rad/s
data.angular_velocity = data.rpm * 2 * Math.PI / 60;
// kWh → MJ
data.energy_mj = data.energy_kwh * 3.6;
// 바이트 → MB (소수 1자리)
data.size_mb = Math.round(data.bytes / 1024 / 1024 * 10) / 10;
msg
日期·时刻转换
const d = new Date(metadata.ts);
// ISO 8601 — 2026-05-12T12:34:56.789Z
data.iso = d.toISOString();
// 사람용 — 2026-05-12 21:34:56 (KST)
data.local = d.toLocaleString('ko-KR', { hour12: false });
// 날짜만 — 2026-05-12
data.date = d.toISOString().slice(0, 10);
// 시간만 — 21:34:56
data.time = d.toTimeString().slice(0, 8);
// 분 단위로 내림 (스파크라인 키)
data.minute_key = Math.floor(metadata.ts / 60000) * 60000;
// 시간 단위로 내림
data.hour_key = Math.floor(metadata.ts / 3600000) * 3600000;
// 한국 시간 (UTC+9) 직접 더하기 (서버 시간이 UTC 인 경우)
data.kst_hour = new Date(metadata.ts + 9 * 3600000).getUTCHours();
msg
字符串标准化
// 공백 trim + 소문자화
data.normalized = (data.text || '').trim().toLowerCase();
// 한글·영문 외 제거 (특수문자/공백 정리)
data.clean = (data.text || '').replace(/[^가-힣a-zA-Z0-9]/g, '');
// camelCase → snake_case
data.snake = (data.text || '').replace(/([A-Z])/g, '_$1').toLowerCase().replace(/^_/, '');
// 전화번호 정규화 (숫자만)
data.phone_digits = (data.phone || '').replace(/\D/g, '');
// 한국 휴대전화 자동 포맷 (010-1234-5678)
const p = (data.phone || '').replace(/\D/g, '');
data.phone_formatted = p.length === 11 ? p.replace(/(\d{3})(\d{4})(\d{4})/, '$1-$2-$3') : p;
msg
嵌套对象平坦化
// data.location.address.city → data.city
function flatten(obj, prefix, out) {
out = out || {};
for (var k in obj) {
var v = obj[k];
var key = prefix ? prefix + '_' + k : k;
if (v && typeof v === 'object' && !Array.isArray(v)) flatten(v, key, out);
else out[key] = v;
}
return out;
}
data = flatten(data);
msg
平坦对象 → 嵌套
// data.user_name + data.user_email → data.user = {...}
function unflatten(obj) {
var out = {};
for (var k in obj) {
var parts = k.split('_');
var cur = out;
for (var i = 0; i < parts.length - 1; i++) {
cur[parts[i]] = cur[parts[i]] || {};
cur = cur[parts[i]];
}
cur[parts[parts.length - 1]] = obj[k];
}
return out;
}
data = unflatten(data);
msg
CSV行生成
// 외부 시스템에 CSV 한 줄 보낼 때
function csvEscape(v) {
if (v === null || v === undefined) return '';
v = String(v);
return /[,"\n]/.test(v) ? '"' + v.replace(/"/g, '""') + '"' : v;
}
data.csv_row = [
csvEscape(metadata.ts),
csvEscape(metadata.asset_id),
csvEscape(data.value),
csvEscape(data.quality)
].join(',');
msg
数组聚合
var arr = data.values || [];
// 합·평균·최소·최대
data.sum = arr.reduce(function(a, b) { return a + b; }, 0);
data.avg = arr.length ? data.sum / arr.length : 0;
data.min = arr.length ? Math.min.apply(null, arr) : null;
data.max = arr.length ? Math.max.apply(null, arr) : null;
// 중앙값
var sorted = arr.slice().sort(function(a, b) { return a - b; });
var mid = Math.floor(sorted.length / 2);
data.median = sorted.length === 0 ? null
: sorted.length % 2 ? sorted[mid]
: (sorted[mid - 1] + sorted[mid]) / 2;
msg
条件字段添加 (模式演进)
// 알람 우선순위에 따라 색상·아이콘 자동 부여
var p = data.priority || 'INFO';
data.color = { ERROR: '#d32f2f', WARN: '#f57c00', INFO: '#1976d2' }[p] || '#9e9e9e';
data.icon = { ERROR: '🔴', WARN: '🟠', INFO: '🔵' }[p] || '⚪';
data.urgency = p === 'ERROR' ? 3 : p === 'WARN' ? 2 : 1;
msg
负载部分遮蔽
function maskEmail(e) {
if (!e || e.indexOf('@') < 0) return e;
var parts = e.split('@');
return parts[0].slice(0, 2) + '***@' + parts[1];
}
function maskPhone(p) {
return (p || '').replace(/(\d{3})\d{4}(\d{4})/, '$1-****-$2');
}
data.user_email = maskEmail(data.user_email);
data.user_phone = maskPhone(data.user_phone);
msg
多语言(i18n)模板
// 메시지 본문을 사용자 로케일에 따라 분기
const tpl = {
'ko-KR': '🚨 ${asset} 온도 ${value}°C 임계 초과',
'en-US': '🚨 ${asset} temperature ${value}°C exceeds threshold',
'ja-JP': '🚨 ${asset} 温度 ${value}°C 閾値超過'
};
const locale = metadata.user_locale || 'ko-KR';
const template = tpl[locale] || tpl['ko-KR'];
data.notification = template
.replace('${asset}', metadata.asset_id)
.replace('${value}', data.value);
msg
JSON Path安全查询
// data.deep.nested.field 처럼 깊은 경로를 안전하게 조회
function get(obj, path, dflt) {
var keys = path.split('.');
var cur = obj;
for (var i = 0; i < keys.length; i++) {
if (cur === null || cur === undefined) return dflt;
cur = cur[keys[i]];
}
return cur === undefined ? dflt : cur;
}
data.city = get(data, 'location.address.city', 'Unknown');
msg
消息合并 (flow_merge后处理)
// flow_merge 가 data.merged 배열로 모은 메시지들을 1건으로 합침
var items = data.merged || [];
data.summary = {
count: items.length,
first_ts: items[0]?.metadata?.ts,
last_ts: items[items.length - 1]?.metadata?.ts,
assets: Array.from(new Set(items.map(function(m) { return m.metadata?.asset_id; }))),
max_value: Math.max.apply(null, items.map(function(m) { return m.data?.value || 0 }))
};
delete data.merged;
msg
可复用脚本集
可直接搬到flow_script_filter / flow_script_transform / flow_switch节点中使用的脚本库。所有脚本均以JavaScript为基准,使用data / metadata简写别名。
变换脚本
单位转换·标注
// 섭씨 → 화씨 + 등급 라벨
data.temp_f = data.temp * 1.8 + 32;
data.grade = data.temp > 80 ? 'HIGH' : data.temp < 0 ? 'LOW' : 'OK';
msg
时刻·班次·星期元数据补充
const d = new Date(metadata.ts);
const h = d.getHours();
metadata.shift = (h >= 6 && h < 14) ? 'A' : (h < 22) ? 'B' : 'C';
metadata.weekday = ['SUN','MON','TUE','WED','THU','FRI','SAT'][d.getDay()];
metadata.is_weekend = (d.getDay() === 0 || d.getDay() === 6);
msg
平均值±3σ报警波段计算
const m = data.mean, s = data.stddev || 1;
data.tag_id = originator.id + '.TEMP';
data.hi = m + 3*s; data.lo = m - 3*s;
data.hi_hi = m + 4*s; data.lo_lo = m - 4*s;
data.use_alarm = true;
msg
MES PO → 工单映射
const po = data;
data = {
master_id: 'WO-MES-' + po.po_no,
asset_id: po.line_id || 'UNASSIGNED',
title: po.product_name + ' (' + po.qty + ')',
due_date: po.delivery_date,
qty: po.qty,
product_id: po.product_code
};
msg
外部负载规范化 (统一各种字段名)
// 외부 시스템에 따라 키 이름이 다른 경우 — 한 줄로 정규화
data.value = data.value ?? data.val ?? data.v;
data.timestamp = data.timestamp ?? data.ts ?? metadata.ts;
data.tag_id = data.tag_id ?? data.tagId ?? data.id;
msg
累积阈值违反次数 (Stateful模式)
// 자산 attribute 와 함께 사용 — 상태는 메시지 자체에 적재
data.consecutive_fail = (data.consecutive_fail ?? 0) + (data.passed ? 0 : 1);
data.alert = data.consecutive_fail >= 5;
msg
消息统计摘要 (Merge结果加工)
// flow_merge 후 data.merged 배열을 받아 요약
const arr = data.merged || [];
data.count = arr.length;
data.values = arr.map(m => m.data.value).filter(v => v != null);
data.mean = data.values.reduce((a,b)=>a+b, 0) / (data.values.length || 1);
data.max = Math.max(...data.values);
data.min = Math.min(...data.values);
delete data.merged;
msg
按时间段应用阈值
const h = new Date(metadata.ts).getHours();
const threshold = (h >= 8 && h < 18) ? 90 : 70; // 주간 90, 야간 70
data.alarm = data.value > threshold;
data.threshold = threshold;
msg
过滤器脚本
优先级白名单
['ERROR', 'CRITICAL'].includes(data.priority)
时间窗口 (仅工作日)
const h = new Date(metadata.ts).getHours();
h >= 8 && h < 20
资产·标签ID模式
/^MOTOR-.*$/.test(originator.id)
// 또는
metadata.tag_id && metadata.tag_id.startsWith('LINE-A.')
阈值 + 稳定性 (连续N次)
data.value > 100 && (data.consecutive_count ?? 0) >= 5
营业时间·节假日检查
const d = new Date(metadata.ts);
const h = d.getHours();
const wd = d.getDay();
// 평일 09-18시만 통과
wd >= 1 && wd <= 5 && h >= 9 && h < 18
仅特定站点
['SITE-01', 'SITE-02'].includes(metadata.site_id)
检查所有必填字段是否存在
data.value != null && data.tag_id && metadata.ts
阻断自我发布消息 (双向同步)
metadata.source !== 'flow'
Switch case脚本
按优先级等级
// case "Critical"
['CRITICAL', 'EMERGENCY'].includes(data.priority)
// case "High"
data.priority === 'ERROR'
// case "Normal"
['WARN', 'INFO'].includes(data.priority)
按资产状态
// case "Running"
data.status === 'RUN'
// case "Stopped"
['STOP', 'IDLE', 'PAUSED'].includes(data.status)
// case "Faulted"
data.status === 'FAULT' || data.error_count > 0
按工单生命周期
// case "Start"
data.event_type === 'START_REQUEST'
// case "End"
data.event_type === 'COMPLETE' && data.qty_done >= data.qty_planned
// case "Abort"
data.event_type === 'CANCEL'
以上脚本均汇集了运营环境中常用的模式。可直接接线到图形中运行,只需根据领域调整阈值·字段名即可。
术语表
流程引擎术语
| 术语 | 含义 |
|---|---|
| entity_type | 消息主体实体种类 — Asset、Tag、Site、Order、Customer、Product、Employee、Calendar等 |
| originator | 消息所指的主体实体 (entity_type + id)。例如: Asset/MOTOR-001 |
| type | 消息分类标签。用于识别触发器接收到何种事件 |
| data | 消息的正文(负载) — 变换·动作节点主要读写 |
| metadata | 消息的上下文(时刻·站点·班次·标签ID等) — 即使转换也会保留到流程结束 |
| relation | 节点输出连线的标签。SUCCESS/FAILURE/TRUE/FALSE/MATCH/NO_MATCH/DEFAULT/THROTTLED/EXHAUSTED (均为大写) |
*_field动态选项 | 从消息负载路径(如data.tag_id)读取值而非静态值的输入方式 |
| 通配符模式 | 触发器*_pattern选项的通配符表达 — *匹配0个及以上字符,?匹配恰好1个字符 |
| 快照 | 保存图形时自动加载的版本单位备份。事故时用于恢复到以前的状态 |
| 部分更新(fetch+merge) | Update节点先查询现有记录后仅合并输入字段的方式。空值被忽略 |
| SKIPPED | 未匹配触发器模式的消息不流向后续节点,也从计数器中排除的处理 |
| EXHAUSTED | flow_retry节点超过最大重试次数时触发的分支 |
| THROTTLED | flow_throttle节点阻断超出窗口内限额消息时触发的分支 |
领域·缩写词典
按快速参考排列了工厂运行环境中常用的缩写,方便在本手册中查找。
| 缩写 | 含义 | 在本手册中 |
|---|---|---|
| MES | Manufacturing Execution System — 工单·生产实绩管理 | 外部集成对象 (HTTP/MQTT/外部DB) |
| ERP | Enterprise Resource Planning — 企业资源·计划系统 | 外部集成对象 |
| SCADA | Supervisory Control and Data Acquisition — 监控·控制 | 外部集成对象 / flow_publish_asset_command的出口 |
| OPC | Open Platform Communications — 工业通信标准 | 边缘类别(flow_edge_opc_*)的对象 |
| OEE | Overall Equipment Effectiveness — 稼动率 × 性能 × 质量 | flow_on_oee_event触发器 |
| RAM | Reliability·Availability·Maintainability | flow_on_ram_event触发器 |
| EMS | Energy Management System | flow_on_ems_event触发器 |
| EQL | 领域事件规则表达式 | flow_script_filter/flow_script_transform的language: "EQL"选项 |
| CEP | 复合事件处理 (Complex Event Processing) | 基于EQL的规则引擎。报警必须仅通过CEP路径产生才能维持一致性 |
| PO | Purchase Order — 采购·生产订单 | MES集成时转换为工单的单位 |
| WO | Work Order (工单) | flow_*_work_order节点组 |
| CMMS | Computerized Maintenance Management System | 外部集成对象 |
| MTTF/MTTR | Mean Time To Failure / To Repair | RAM事件的核心指标 |
| HMI | Human–Machine Interface | SCADA等的运行界面 |
版本说明
本手册当前所涉及的主要功能组的引入时间点。若您正在运行较早版本,部分功能可能有所不同。
V2026.05 — 领域自动化扩展
- 新设边缘(Edge)类别11种 — OPC服务器/标签CRUD + 监控自动化
- 工单状态转换5种 —
flow_start/end/pause/resume/abort_work_order - 标签报警波段部分更新2种 — 数值型/布尔型报警波段仅更新输入字段
- 新增**
flow_retry流程控制** — 退避 + 最大尝试次数 + EXHAUSTED分支 - 新触发器5种 —
flow_on_asset_health_status/_connection_status/_oee_event/_ram_event/_ems_event - Update节点部分更新(fetch+merge) — 查询现有记录后仅合并输入字段
- Relations统一为大写 —
SUCCESS/FAILURE/...均为大写,现有图形自动转换 - SKIPPED处理 — 未匹配触发器模式的消息从计数器中排除
- 全部重新部署(redeploy)运维工具 — 批量重新加载活跃流程
- 全部计数重置 — 批量重置按节点+按流程的统计
V2026.03 — 发布稳定性
- 保存图形时加载自动快照
- 实时调试面板 (检查器下方2秒刷新) + 节点指示灯·duration显示
- import/export (图形JSON文件)
- 引入负载看门狗 — 负载持续时发出一行诊断日志
更早版本
- M1 — 引擎·UI骨架, 图形CRUD, 可视化画布
- M2 — 领域动作节点(资产·标签·工单·生产领域等)
- M3 — 触发器多样化, 过滤器/变换/外部集成8种, 脚本引擎
- M4 — 调试加载·执行历史界面, 统计, 部署/解除, 导入/导出
- M5 — 领域触发器集成, 场景批量验证
本手册涉及的所有功能和动作均以上述V2026.05时点为准。
相关界面
- CEP (复合事件处理): 基于EQL的领域事件规则引擎
- 报警: 报警发生·历史
- 数据点: 标签数据查询·分析
- 开发者指南: 流程引擎: 引擎架构·节点接口·DB模式