フロー (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時間以内に最初のフロー作成)
- コア概念 — 5分
- クイックスタート — 5分
- エンドツーエンドチュートリアル — 30分 (ステップ 1~10 に従う)
- 画面構成 + 編集画面 — 10分
- 事例フロー 1·2·3 に従う — 10分
→ 最初のフロー配備完了。その後必要なノードをノードカタログから検索。
🧑🏭 現場運用者 (自動化シナリオ作成)
- トリガー別ペイロード例 — 実データ構造の理解
- ノードオプション詳細リファレンス — よく使うノードオプション熟知
- データ変換 Cookbook — 一般的な変換パターン複写使用
- グラフパターンカタログ — 結線パターン選択
- 事例フロー (17種) — シナリオ別完成品参照
🛠 システム管理者 (運用・チューニング・障害対応)
- フロー指標とアラーム — どの指標を見るべきか
- 実行履歴で何を見るか — 診断ログ解釈
- クラスター・HA 動作 — マルチノード環境理解
- エンドツーエンドトレース — 問題メッセージ追跡
- パフォーマンス限界とチューニング + 緊急対応手順
🔌 開発者・統合エンジニア (外部システム連携)
- 外部システム統合 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) | フロー開始点となるノード。ドメインイベント(タグポイント・アラーム・資産イベントなど)、外部進入(ウェブフック・MQTT・外部DB)、時間(スケジュール)の3種類があります |
クイックスタート
フローは以下の4ステップで最も早く習得できます。
- リスト画面で
새 플로우ボタンを押して空のフローを作成 (名前と説明を入力). - 編集画面で左パレットのトリガーノード(例:
flow_on_tag_alarm) → フィルター → アクション(例:flow_send_email)を順に配置し、ノード間をワイヤーで接続します。 - 各ノードをクリックして右インスペクターでオプションを入力し、右上の保存 → テスト実行で一度動作を確認します。
- 上部の配備トグルをオンにするとトリガーイベント受信時に自動実行され、ライブデバッグパネルと実行履歴画面で結果確認できます。
自動化シナリオの結線パターンは、本ドキュメントの事例フローと活用例を参照してください。最初から最後までフロー作成する段階別チュートリアルはエンドツーエンドチュートリアルで追跡できます。
エンドツーエンドチュートリアル
本番環境で実際に使用可能な自動化フロー1つを最初から最後まで作成する段階別チュートリアル。シナリオ: モーターの温度が 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 のprefix部分(例: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、後続ノード灰色 — 正常。
- すべての分岐検証後に全体カウント初期化で統計ウィンドウをリセット
- 上部の配備トグル ON
これで実際の運用環境でモーター温度が 80°C を超えると、即座に自動的にメンテナンス作業指示が発行され、担当者にメールが送られます。
ステップ 10 — 運用監視
配備直後 5~10分間、以下を確認するとセーフです。
| 場所 | 確認項目 |
|---|---|
| リスト画面 | 当該フロー行の実行回数が正常範囲で増加しているか (暴走していないか) |
| ライブデバッグパネル | ノード失敗がないか |
| 実行履歴画面 | 失敗メッセージがあれば NODE_ERROR の原因確認 |
| 受信メール | 意図しない頻度でアラート届いていないか |
次のステップ
このフローを発展させたい場合:
- Retry 結線追加 — メール/プッシュ送信失敗時のバックオフ再試行 (事例 8 参照)
- 自動帯域調整 — 6時間平均 ± 3σ で閾値自動補正 (事例 6)
- ダウンタイム累積 — 過熱履歴を資産 attribute に累積 (事例 14)
- 品質ライン隔離 — 過熱が連続発生したらライン自動停止 (事例 15)
画面構成
フローは以下3画面で構成されます。
| 画面 | 用途 |
|---|---|
| リスト | 登録されたフロー一覧・検索・一括配備・インポート |
| 編集 | ビジュアルキャンバスでノード配置・接続・設定 |
| 実行履歴 | ノード単位実行ログ・タイムライン照会 |
リスト画面
上部検索・作成エリアとフロー一覧テーブルで構成されます。
上部ツール
| 項目 | 説明 |
|---|---|
| ステータスフィルター | 전체 / 배포 / 해제 |
| 名前・説明検索 | キーワードでフローをフィルター |
| 新規フロー | 空のフローを作成 (名前・説明を入力) |
| インポート | エクスポート受け取った 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)から値を抽出。値がなければ静的値にフォールバック |
スクリプトノード(フィルター・変換・スイッチ)はインスペクター内で直接コードエディターで編集できて、別ウィンドウで展開して大きな画面で作成することもできます。
コードエディターフォントは可読性強化モノスペースフォントスタック(Cascadia Code · JetBrains Mono · Consolas · Menlo 優先)で表示され、日本語コメントも安定して配置されます。
カウント初期化
ノード設定パネル右上には2つの初期化ボタンがアイコンで表示されます。各ボタンにマウスを合わせるとツールチップが案内します。
| アイコンボタン | 動作 |
|---|---|
| 🩹 (絆創膏) — エラーカウント初期化 | このフロー内のすべてのノードの累積エラーカウントのみ を 0 にリセット |
| 🔄 (回転矢印) — 全体カウント初期化 | 処理・エラー・処理時間すべて + フロー単位統計をすべて 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 ウェブフック受信 |
KAFKA_INBOUND / MQTT_INBOUND | 外部トピック受信 |
TIMER | スケジュール発火 |
トリガー別ペイロード例
スクリプトノード作成時に、どのフィールドにアクセスできるか正確に知る必要があります。以下は各トリガーが生成する実際のメッセージ JSON 例です。すべてのトリガーメッセージには共通で type、originator、data、metadata が含まれます。
flow_on_tag_point — タグポイント受け取り
タグ1つに値が入るたびに発火。最も一般的なトリガーです。
{
"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 3タイプで統合発火。単一ノード1つですべての3種ライフサイクルを受け取ります。
{
"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 ヘッダーはmetadataに一部のみ(remote_addr/method/request_id)表示されます。
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 クエリーを定期実行して、各行ごとにメッセージ1件ずつ発火。
{
"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 というノードはありませんアラームトリガーは対象に応じて2つに分かれます — タグアラームは**flow_on_tag_alarm、
資産アラームはflow_on_asset_alarm**です。前のドキュメントの例が flow_on_alarm を
使っていたので、そのまま従うとパレットに見つからないノード名になります。
すべてのトリガーノードは
*_patternオプション(glob:*、?)でメッセージ単位事前フィルタリングができます。パターン未マッチメッセージは後続ノードに送信されず、実行カウントも増加しません(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内に値が1階層多く入っているメッセージを 後ノードで${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 | デバッグログ (レベル / プレフィックス) |
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]─▶ [에러 핸들러 / 알림]
ノードオプション詳細リファレンス
複雑なノードのオプションを1行表ではなく、オプション別デフォルト・例・失敗処理まで詳細に説明します。
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 | * | メッセージ事前フィルタリング 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 キーまたは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 | — | CC |
bcc / bcc_field | string | — | BCC |
subject / subject_field | string | (必須1個) | 件名 — テンプレート置換 |
body / body_field | text | (必須1個) | 本文 (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 | 1ポーリング当たり最大行 (これ超過しても安全) |
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 | (必須1個) | 発行本文 — 未設定時 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/デフォルト) |
{컬럼명}_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 カラムにデフォルト値を自動で埋めます。
| ドメイン | 自動填充カラム |
|---|---|
| 作業指示 | 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 を自動 lookup して 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 ノードは2種類の表現式をサポート。
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
スイッチ例 (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 | 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 | (必須1個) | 受信者ユーザー ID (カンマ区切り)またはペイロードパス |
title / title_field | string | (必須1個) | 通知タイトル (50文字以内推奨) |
body / body_field | text | (必須1個) | 本文 (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)
複数アラームを一度に1通にまとめる (5分単位):
[flow_on_asset_alarm]
↓
[flow_merge window=300s]
↓ (data.merged 배열)
[flow_script_transform — 요약 만들기]
↓ data.push_body="알람 N건: A자산, B자산, ..."
[flow_send_push]
ユーザー不在時自動エスカレーション
プッシュ未確認30分後に SMS またはメールで自動エスカレーション。
[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: 外部ウェブフックで資産コマンド発行
外部システムが送信 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でルーティングされるので、別フィルター不要で安全に結線。
事例 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分に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 3種イベントを30秒単位でまとめて1回のレポートメッセージで送信。
[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: エネルギー閾値超過 — ライン一時停止推奨
ラインの時間当たりエネルギー消費が予算超過したら運用者に SMS と共に推奨メッセージ。
[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 も新規割り当て、ワイヤーリンク自動リマップ
- import 直後フロー解除状態で適用され、運用者検証後直接配備必要
エラーハンドリング
ノード単位
- ノード処理中例外発生したら自動的に
FAILURErelation でルーティング。 FAILURE出力に接続ノードなければ、メッセージは drop されてエラーログのみ。- トリガーノードの
*_pattern未マッチメッセージはSKIPPED処理され、後続ノードに送信されず、処理/エラーカウンターにも含まれません。 - すべてのエラーはライブデバッグパネルと実行履歴画面に表示。
フロー単位
| 項目 | デフォルト | 説明 |
|---|---|---|
| max_depth | 100 | 1メッセージ処理中の累積ノード訪問数制限 |
| max_revisit | 3 | 同じノード再訪回数制限 (サイクル無限ループ防止) |
| flow_timeout_ms | 30,000 | 処理時間超過時に強制終了 |
外部 IO ノード再試行
HTTP·Kafka·外部DB など外部 IO ノードは、retry_count / retry_delay_ms オプションで、ノード内部で即座再試行可能で、すべての再試行失敗時 FAILURE relation でルーティング。より精緻なバックオフやEXHAUSTED 分岐処理が必要な場合は flow_retry ノード別に使用。
運用診断
管理者は、システムメニューでフロー エンジンのディスパッチキュー状態(ワーカー稼働有無、待機キュー大きさ、累積処理/失敗/ドロップ件数、最後エラー)確認できます。
サーバーはバックグラウンドでフロー ワーカープールの負荷を定期点検し、負荷が一定時間以上続くと運用ログに1行で状態変化(HEALTHY → DEGRADED → CRITICAL)を記録。正常復帰すると回復ログ1行さらに残ります。
権限
フローは別権限モデル導入せず、既存システム認証/権限をそのまま使用。
| 機能 | 必要権限 |
|---|---|
| リスト照会 / 実行履歴照会 | すべての認証ユーザー |
| フロー作成·編集·配備·削除 | ADMIN |
| インポート/エクスポート / すべて再配置 | ADMIN |
運用パターン (Recipes)
よく使われる結線形態をまとめたライブラリ。各パターンはそのまま複写して新フローの出発点として活用できます。
パターン 1 — 処理 + 通知分岐
処理結果に応じて SUCCESS は適用、FAILURE は通知で二重分岐。
[Action Node]
├── SUCCESS → [flow_log] / [flow_save_attributes] / ...
└── FAILURE → [flow_send_email] / [flow_send_push]
パターン 2 — 安全な再試行
外部 IO が一時的失敗に強くなるようにバックオフ + EXHAUSTED ハンドラー構成。
[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
高頻度入力を一定時間モーデ1メッセージに変換。
[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 — スロットル + デバウンス結合
分当たり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 |
| ウェブフック受信 | 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ステップ — 隔離テスト
問題ノードを新規臨時フローに単独移動、テスト実行で1メッセージのみ発火させて結果を隔離検証。正常動作なら、元フロー結線・直前段階メッセージ形状が疑わしい。
6ステップ — カウンター初期化後に再現
全体カウント初期化後、1メッセージのみ発火させてきれいな統計で再現すると、問題がより明確。
よくある問題
| 症状 | 原因・対応 |
|---|---|
| トリガーが発火していない | フロー配備状態か、トリガーノードの *_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. フロー1つにトリガー複数置けますか?
A. はい。同じフロー内にトリガーノード複数置くとすべて進入点になって各自独立発火。flow_merge と結合するとイベント種類複数を1後続処理で集約可能。
Q. トリガーなしフロー作れますか?
A. はい。保存時に警告表示されますが保存できます。テスト実行または別フローの flow_subflow 呼び出しでのみ発火させる「ライブラリー」形式で使用可能。
Q. 1メッセージが複数ノード同時に通過させられますか? A. はい。1ノード出力ポートに複数ワイヤー接続すると同じメッセージが複数後続ノードに同時分岐送信。
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. 同じペイロードが2度処理されています。どう防げばいいですか?
A. トリガーステージで *_pattern で絞るか、flow_throttle/flow_debounce で頻度制限。外部入力(ウェブフック·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 など)を1度は発火させて、実行履歴で結果確認。
- 運用反映 — 検証フロー グラフ JSON を
내보내기→ 本番環境で가져오기→ 運用者が検証後に配備トグル。 - 問題発生時ロールバック — 即座に解除トグルで無効化。自動スナップショット適用されているので、開発者に依頼して以前バージョンに戻します。
カウンター·統計活用
- リスト画面の処理トレンド スパークラインが急に平坦化したら、トリガー未発火または SKIPPED 処理の可能性確認。
- エラー比率ドーナツチャートが異常なら、ノード設定の
*_fieldパスを疑ってみてください。メッセージペイロードにキーなければ頻繁に失敗。 - 運用検証後に全体カウント初期化で統計ウィンドウをリセットして、平時ベースラインを新規測定。
緊急対応手順
運用中に問題発生したときの神速対応ステップ別手順。
シナリオ 1 — 特定フロー暴走(メッセージ暴増)
症状: 1フロー実行件数が平時比数十~数百倍暴増、エラーも同増
対応
- 即座に当該フロー解除トグルで無効化 (リスト画面で1秒)
- ライブデバッグパネル·実行履歴で、どのトリガーが暴走起こしているか確認
- トリガーノード
*_patternオプション狭める、または直後にflow_throttle/flow_debounce結線 - 必要なら
flow_check_existence_fieldまたはflow_msg_type_filterでメッセージタイプ限定 - 修正後テスト実行で再現 → 正常確認後に配備再開
シナリオ 2 — 外部システム障害で一括失敗
症状: HTTP/Kafka/外部DB など外部 IO ノードのエラーが同時多発
対応
- 影響広ければ関連フロー一括解除(リスト画面で一括選択後一括解除)
- 外部システム復旧確認
flow_retry結線がない外部 IO ノードなら結線追加- 外部システム応答が遅くなってたら
timeout_ms調整 - 運用環境に応じてすべて再配置後に配備再開
シナリオ 3 — サイクルで無限ループ
症状: 単一メッセージが同ノード繰り返し通過、処理時間累積
対応
max_revisit保護で最大3回までのみ再訪許可で自動遮断されますが、運用安全のため解除後に点検- グラフでサイクルを視覚的に追跡 (編集画面でワイヤー順追い)
- 意図されたループなら(
flow_retry)max_attemptsが適切か確認 - 意図されないサイクルなら結線削除または
flow_check_relationで分岐追加 - 修正後全体カウント初期化 → 再配置
シナリオ 4 — データ破損疑い (不正な自動更新)
症状: 自動化でドメイン データが意図と違う形に更新
対応
- 即座に当該フロー解除
- 外部保管中の直前グラフ JSON バックアップまたは自動スナップショットで以前バージョン確認 (管理者協力)
- グラフ分析: 意図しない Update ノード結線·誤った
*_fieldパス·スクリプト変換エラー確認 - 影響受けたドメインデータは別手順でバックオフィス画面で修正
- 修正グラフをテスト実行で検証後に再配置
シナリオ 5 — システム点検·保守中一時中止
症状: 外部システム点検時間中、外部 IO 呼び出し一時停止したい
対応
- 影響フロー一括解除(リストで複数選択)
- 点検終了後に配備再開
- 独自スケジューラー持つトリガー(スケジュール·外部 MQTT 購読·外部DB ポーリング)はすべて再配置1回実行で正常登録確認
シナリオ 6 — ライブデバッグ適用キュー満杯
症状: 運用診断ページ適用キューサイズが閾値近く、ドロップ件数増加
対応
- デバッグモードがオンのフロー数と適用量削減 (検証終フロー はデバッグモード OFF)
- トリガーステージの
*_patternで後続処理量自体削減 - システム自動回復モード入れば、負荷ウォッチドッグに
DEGRADED/CRITICALログ1行残ります — 管理者に共有。
すべてのシナリオで優先順位は即座解除 → 原因把握 → 修正 → 検証 → 再配置順序。変更手順も参照。
セキュリティ·機密情報処理
フロー は外部 API 呼び出し·メール/SMS 送信·ウェブフック受信など機密入出力を扱うので、以下原則守ってください。
トークン·認証情報
- API トークン·パスワードをノード設定に平文で直接入力しないでください。
${ENV_VAR}形式の環境変数置換を使用して運用環境の秘密ストレージから注入受け取ってください。 - 本番環境でトークン非公開のため
headersJSON のトークン位置を可能なら環境変数で分離。
// 권장
{ "Authorization": "Bearer ${MES_TOKEN}" }
// 비권장
{ "Authorization": "Bearer eyJhbGciOi..." }
ウェブフック進入点保護
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 が表示。
| 指標 | 意味 | 異常信号 |
|---|---|---|
| 全体実行件数 | トリガー発火してグラフ1度通過した件数 (成功/失敗計) | 平時比急減 → トリガー死亡 / 急増 → 暴走 |
| エラー件数 | グラフ中どのノードでも FAILURE 分岐落ちした件数 | 全体比 5%以上なら点検 |
| エラー率 | 에러 / 전체 × 100 | 一定閾値超過時にアラーム立ててください |
| 最終実行 | 最近発火時刻 | 「5分以上発火なし」が正常の場合のみ OK |
| 平均処理時間 | メッセージ1件がグラフ全体通過する平均 ms | 外部 IO ノード追加/削除時に変動大 |
| 最近処理時間推移 | 5分スパークライン — ライブデバッグパネル | spike 発生時は外部システム応答点検 |
ノード単位 KPI
編集画面でノードクリックすると、ノード右下に以下が表示。
[Node Name]
처리 999 · 에러 3 · 평균 12ms · 최근 18ms
| 位置 | 表示 |
|---|---|
| 上端点灯 | 灰色(待機) / 緑色(処理中) / 赤色(エラー) |
| 下端ラベル | 처리 N · 에러 M · 평균 Xms · 최근 Yms |
ノード統計 4種
| カウンター | 意味 |
|---|---|
| 処理 | そのノードに入ってきて正常分岐で出て行ったメッセージ数 |
| エラー | FAILURE 分岐または例外発生数 |
| 平均処理時間 | ノード自体処理 ms (外部 IO 含む) |
| 最近処理時間 | 最後1件処理 ms |
システム単位メトリック
フロー エンジン全体のメトリックはシステム → モニタリング画面で確認可能。
メトリックは JMX ドメイン plantpulse.core.engine 下に以下5つで公開。
| メトリック | 種類 | 意味 | 見方 |
|---|---|---|---|
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 を表で載せていますが、そのような名前のメトリックは
存在しません。 上の5つがすべて。
(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 伝達を保証。
| ケース | 動作 |
|---|---|
| 正常処理 | 1度発火 → 1度グラフ通過 → 1度完了 |
| 処理中エンジン再起動 | トリガーキューに残ったメッセージは次起動後に再処理 |
| ノード単位例外 | FAILURE 分岐のみ転移。メッセージが消失しない |
| 外部 IO タイムアウト | FAILURE 分岐 + flow_retry で自動再試行可能 |
重複可能性 — 正確に1度(exactly-once)ではありません。
flow_create_*が中断再試行されたら同一 ID が2度入ってくる可能性があり、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 のイベント2件がほぼ同時発生しても、A のイベント は発生順で処理。資産 A と資産 B のイベント は別ワーカーで並列処理される可能性。
バックプレッシャー(Backpressure)
処理速度が発火速度に追いつかない時の動作。
| 状況 | 動作 |
|---|---|
| キュー深さ < 80% | 正常 — 新規トリガーすぐに enqueue |
| キュー深さ 80%~100% | 診断 WARN ログ発生 — キューは引き続き受容 |
| キュー深さ = 100% (満杯) | 新規トリガーをdrop + 診断 ERROR ログ発生 |
キューが満杯時に運用者がすること
- システム → モニタリングで
flow.engine.queue_depth確認 - リスト画面で実行件数暴増したフロー特定
- そのフロー トリガーの事前フィルター(
*_pattern)狭めて負荷低減 - または、そのフロー配備を一時解除して非常処理
サーキットブレーカー (外部 IO)
外部 IO ノード(HTTP/Kafka/MQTT/メール)は、以下条件充たしたら30秒間一括ブロック。
| 条件 | 閾値 |
|---|---|
| 最近1分内連続失敗 | ≥ 10件 |
| 平均応答時間 | ≥ 10秒 |
ブロック中にそのノード入ってくるメッセージはすぐに FAILURE で分岐。外部システム応答回復したら自動ブロック解除。この動作は、外部障害がフロー エンジン自体を麻痺させないように隔離する役割。
サーキット開いた時点は診断ログに
code = CIRCUIT_OPENEDで記録。外部システムが早く回復したら実行履歴のmetadata.retry_countを見て、喪失メッセージを手動再実行できます。
サイクル(循環)防止
フロー内に以下のようなサイクル生じると無限ループ可能性。
A → B → C → A (잘못된 결선)
フロー エンジンは2つの防御線でサイクルをブロック。
| 防御線 | 動作 |
|---|---|
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 > フロー実行履歴)に適用される
イベント は5種類。
| イベントタイプ | いつ | 適用条件 |
|---|---|---|
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 | メッセージ単位追跡用 — 1メッセージの流れつなぎ見たい時この値でまとめ |
message · error_message | 人間向け説明、例外メッセージ |
duration_ns | ノード実行時間 (ナノ秒) — NODE_OUT · NODE_ERROR でのみ意味あり |
NODE_IN · NODE_OUT はデバッグモード時のみ適用。平時実行履歴に
FLOW_START · FLOW_END のみ見えるのは正常。ノード1個ずつ追い見たい場合は
そのフロー のデバッグモードをオン — その代わりログ量大幅増加。
症状で探す
コード の代わりに症状から出発。
| 症状 | 先に見る場所 |
|---|---|
| トリガー は発火しているが、ノード何も動いていないような感じ | フロー配備(トグル ON)されているか · トリガーの *_pattern がメッセージ濾し取っていないか |
| 処理が途中で止まる | 1メッセージ30秒(flow_timeout_ms)超過。サーバーログの FlowExecutor timeout 警告で確認 |
| 同じメッセージが反復処理される | サイクル。max_revisit(3回)で自動遮断されますが、結線点検 |
| グラフ長くて終わらない | max_depth(100)超過。サブフロー で分割 |
| ノードが例外throw する | 実行履歴で NODE_ERROR 行の error_message を見ます |
| 外部呼び出し失敗 | 当該ノード data.response_status · data.error を後ノードで確認 |
上限値(100 · 3 · 30,000ms)は
ExecutionStateのデフォルト値。
クラスター·高可用性(HA)動作
プラットフォーム がクラスター環境で導入された場合、フロー エンジンは以下ルールで分散動作。
ノード役割分離
| 役割 | 動作 |
|---|---|
| リーダー(Leader) | フロー グラフ変更(保存/配備/解除)を直列処理する単一ノード |
| ワーカー(Worker) | トリガー発火·メッセージ処理を並列実行する一般ノード (すべてのノードがワーカー役も兼ねます) |
| スケジューラー(Scheduler) | flow_schedule cron 評価担当 — リーダーと同一ノード |
リーダーノード はプラットフォーム起動時に自動選出され、リーダーがダウンすると別ノード が自動的にリーダーに昇格。運用者が直接指定する必要ありません。
メッセージ分配 — 資産単位一貫ルーティング
| 分配キー | 動作 |
|---|---|
originator.id (資産/タグ/注文 ID) | 同じ資産メッセージは常に同じワーカーへルーティング — シーケンス保証 |
外部受入 (flow_on_webhook など) | ラウンドロビン分配 |
このルールで、資産 A のイベント が複数ワーカーで同時処理されて順序がぐじゃぐじゃになる事態を防止。結果的に、資産単位では単一ワーカー処理保証、資産間では並列処理が同時に成立。
リーダー障害時動作 — フェイルオーバー
[Leader 다운 t=0]
↓
[다른 노드가 리더 승격 시도 t=0~3s]
↓
[새 리더 확정 t=3~5s] ← cron 스케줄·플로우 배포 변경 재개
↓
[기존 워커들은 정상 동작 유지 — 트리거 처리 영향 없음]
| ステージ | 影響 |
|---|---|
| 0~3秒 | flow_schedule 発火一時停止 / トリガー処理は影響なし |
| 3~5秒 | 新リーダー確定、スケジュール 再開 |
| 5秒以降 | 正常 |
cron 発火は**
misfireポリシーに従って、1度失火した発火を即座に再実行**。運用中はflow_scheduleの評価単位が5秒未満なら短期失火が見える可能性。
トリガーキューの永続性
| 項目 | 動作 |
|---|---|
| トリガーキュー位置 | インメモリキュー + 永続ストレージ (トランザクションログ) |
| ノード再起動時 | キューに残っていた未処理メッセージは次起動後に再ディスパッチ |
| 処理中ノードダウン | そのメッセージは再処理されますがat-least-once 保証で外部システムに멱等性キー必要 |
クラスター配備時運用者チェックリスト
- すべてのノードの時間同期 (NTP) — トリガー発火時刻がノード間一致
- 外部システム(MQTT/Kafka ブローカー)をすべてのノードが到達可能なネットワーク位置に配置
flow_send_emailの SMTP 設定はシステム設定に1度のみ登録 — すべてのノードが共有- ワーカープール サイズ(
flow.engine.workers)をノード別 CPU コア数に合わせ調整
単一ノード モード (開発·小規模)
- リーダー·ワーカー·スケジューラー すべて1プロセスで動作
- フェイルオーバーなし — ノード ダウン したらフロー自体中断
- トリガーキューはインメモリ + ディスク で保証され、再起動後復旧
エンドツーエンドトレース
1メッセージがグラフ全体をどう通過したか事後追跡する方法。
自動相関 ID 付与
フロー エンジンはすべてのトリガー発火メッセージに**相関 ID(correlation ID)**を自動付与。
{
"type": "POST_TELEMETRY",
...,
"metadata": {
"trace_id": "tr-a1b2c3d4-e5f6-7890-...",
"span_id": "sp-01",
"parent_id": null,
...
}
}
| フィールド | 意味 |
|---|---|
metadata.trace_id | 1トリガー発火全体に固有 ID — グラフのすべてのノード通過メッセージが共有 |
metadata.span_id | ノード別固有 ID — ノード通過ごとに更新 |
metadata.parent_id | 直前ノード span_id |
実行履歴画面で trace_id で検索
実行履歴画面で trace_id 入力欄に 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 to N)
┌─→ [action 1]
[trigger] → [t] ─────┼─→ [action 2]
└─→ [action 3]
| 特徴 | 1メッセージを複数アクションが同時処理 | | 使用例 | アラーム発生 → メール + SMS + Slack + ワークオーダー作成 |
同じメッセージの複写が複数ノードに分配。各分岐は独立処理され、1分岐失敗が他分岐に影響なし。
パターン 3 — Fan-in (N to 1) — Merge
[trigger A] ──┐
[trigger B] ──┼─→ [flow_merge] → [aggregator] → [action]
[trigger C] ──┘
| 特徴 | 複数トリガーのメッセージを時間ウィンドウ内に集めて1度処理 | | 使用例 | 5分中に発生したすべてのアラームをまとめて日報1通 |
パターン 4 — Scatter-Gather (分散 → 集約)
[trigger] → [split] ─┬─→ [process] ──┐
├─→ [process] ──┼─→ [merge] → [action]
└─→ [process] ──┘
| 特徴 | 配列を要素別分割処理後、結果再集約 | | 使用例 | 100件の外部注文を並列検証後、結果一括報告 |
パターン 5 — Switch (条件分岐)
┌─[CRITICAL]→ [긴급 알람]
[trigger] → [flow_switch] ──┼─[WARN] → [경고 알람]
└─[NORMAL] → [통과]
| 特徴 | 1メッセージを条件別に異なるルートへ | | 使用例 | アラーム優先度別処理チャネル分離 |
パターン 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 同期 → 作業者通知 (1ステージ失敗時に前ステージ すべて取消) |
パターン選択ガイド
| 要件 | 推奨パターン |
|---|---|
| シンプル閾値アラーム | Pipeline (1) |
| 1事件 → 複数チャネル | Fan-out (2) |
| 複数事件 → 1要約 | Fan-in / Merge (3) |
| 配列一括処理 | Scatter-Gather (4) |
| 分岐処理 | Switch (5) |
| 外部API 信頼性 | Retry (6) |
| 外部システム保護 | Circuit (7) |
| 失敗メッセージ保管 | DLQ (8) |
| 最近N件累積意思決定 | Sliding (9) |
| 複数外部システム 一貫性 | Saga (10) |
フロー テスト ベストプラクティス
フローを安全に変更·配備するためのテスト戦略。
3ステージテスト — 単位 → 統合 → シミュレーション
| ステージ | ツール | 検証対象 |
|---|---|---|
| ① 単位 | 編集画面 ▶ テスト実行 | 1ノード 単独動作 (スクリプト式、外部呼び出し応答) |
| ② 統合 | 同画面、臨時ペイロード + ライブデバッグ 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 へ | 誤送信防止 |
カウンター初期化後に負荷テスト
- 新規バージョンフローカウント初期化 (全体)
- 5分間、正常トリガー発火させる
- リスト画面で処理/エラー カウント、平均処理時間確認
- p95処理時間が既存比 20%以上延びたら、原因分析後にロールバック検討
パフォーマンス限界とチューニング
フロー エンジンの処理限界とチューニング 手法。
基本限界
| 項目 | デフォルト | 備考 |
|---|---|---|
1メッセージのノード訪問数 (max_depth) | 100 | ノード100個通過する間に終わらなければ強制中断 |
同じノード再訪回数 (max_revisit) | 3 | サイクル無限ループ防止 |
フロー処理時間 (flow_timeout_ms) | 30,000ms | 1メッセージ処理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 | 短いと断片化、長いと latency 増加 |
flow_retry バックオフ | 開始 1,000ms × 2倍 | 5回再試行なら 1·2·4·8·16秒 |
処理量を増やす方法
- トリガー事前フィルター —
*_patternで必要メッセージのみ進入(最効果的) - サブフローでグラフ単純化 — メイン グラフは分岐·ルーティングのみ、重い処理はサブフローへ
- 外部 IO を非同期で結線 —
flow_delay/flow_throttleで外部 API 負荷平坦化 - デバッグモードは検証ステージのみ — 運用安定後にデバッグモード OFF
- 必要カテゴリのみ結線 — 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]
新規フロー配備チェックリスト
本番環境に新フロー配備する前に確認する項目。
グラフ構造
- トリガーノードが正確に1つ(または意図された複数)結線されているか
- すべてのアクションノード
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 エクスポート)が外部に保管されているか
- 変更理由·担当者がフロー説明に記録されているか
- 通知(メール/SMS/プッシュ)結線があれば、受信者が検証されているか
- 配備直後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 の核は2バージョンが同時配備された状態で即座切り替え — 問題発見時に v1 を即座に再配備してロールバック。
v2 隔離結線パターン
[trigger 모든 메시지]
↓
[flow_script_filter — metadata.test_mode === true]
↓ TRUE
[새 로직]
テスト用メッセージは/flow/{id}/runAPI で metadata.test_mode=true を明示して送信。
戦略 3 — A/B テスト (パフォーマンス·結果比較)
2バージョンが同じメッセージを受けて異なるアクション 後、結果を比較。
[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 — 1ラインで検証後に拡大 |
| 既存ロジック大きな変更 (構造改編) | Blue-Green — 即座ロールバック可能 |
| 2アルゴリズム中どちらがいいか測定 | A/B テスト |
| シンプルなオプション値調整 | 直接変更 — metadata.audit_diff メモ後1週間モニタリング |
ロールバック手順 (共通)
問題発見時に即座ロールバック:
- リスト画面で新バージョン解除トグル
- (Canary/A·B の場合)トリガー パターンを0件マッチに変更
- 既存バージョンが単独動作しているか5分モニタリング
- 診断ログで
code=FLOW_NOT_DEPLOYEDがないか確認 - 原因分析 — 実行履歴で
trace_idで失敗メッセージ追跡
出荷前 チェックリスト (圧縮版)
新規フロー配備チェックリストの核のみ1画面:
[ ] 테스트 실행으로 정상·경계·실패 시나리오 모두 통과
[ ] 외부 IO 노드에 timeout_ms / 재시도 정책 설정됨
[ ] 자격 증명은 ${creds.*} 참조 (평문 미포함)
[ ] 트리거 패턴이 의도한 자산만 매칭
[ ] 영향받는 도메인 (자산/태그/주문) 식별됨
[ ] 운영 시간(특히 야간) 영향 검토됨
[ ] 롤백 시점·기준·담당자 결정됨
[ ] [감사 로그](#audit) 에 변경 의도 메모 작성됨
ノード高速設定リファレンス
運用中によく使うノードの核設定のみを1表にまとめたチートシート。
トリガー高速設定
| ノード | 核オプション |
|---|---|
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 (glob — 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
Import 時に衝突 ID あれば、新 ID で自動発行。外部依存性(認証情報·資産 ID など)は import 後に別途マッピング。
11) すべてのフロー 再配備
POST /flow/redeploy
大量変更後または、ノード再起動後に一括同期に使用。運用中呼び出しは注意 — 一時的処理遅延発生。
12) 自己診断
POST /flow/selftest
フロー エンジンの内部コンポーネント(トリガー キュー·ディスパッチャー·ノード レジストリー·スクリプト ランタイム)を1度点検して状態返却。結果は診断ログにも記録。
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 Rate Limit
| エンドポイント | 分当たり上限 (トークン当たり) |
|---|---|
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を lookup。運用者が別途認証設定する必要ありません。
パターン 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 (不正なペイロード) / 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 — イシュー 自動作成
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}"
}
CSRF トークンを先に GET で取得後、同セッション クッキー で POST する必要。
flow_http_requestノード 2つ接続して、最初ノード応答ヘッダーのトークンを2番目ノード のヘッダーで伝播。
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 トークン | 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'
スイッチ 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'
上スクリプトはすべて運用環境でよく使うパターンを集めたもの。グラフに直接結線しても動作し、ドメインに合わせて閾値·フィールド名のみ調整すればOK。
用語辞典
フロー エンジン用語
| 用語 | 意味 |
|---|---|
| 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)から値を読む入力方式 |
| glob パターン | トリガー *_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 ファイル)
- 負荷ウォッチドッグ導入 — 負荷継続時に診断ログ1行発行
それより以前
- M1 — エンジン·UI 骨組み、グラフ CRUD、ビジュアル キャンバス
- M2 — ドメイン アクション ノード(資産·タグ·作業指示·生産ドメイン など)
- M3 — トリガー多様化、フィルター/変換/外部連携 8種、スクリプト エンジン
- M4 — デバッグ適用·実行履歴画面、統計、配備/解除、インポート/エクスポート
- M5 — ドメイン トリガー統合、シナリオ一括検証
本マニュアルが扱うすべての機能と動作は上 V2026.05 時点を基準にしています。
関連画面
- CEP (複合イベント処理): EQL ベースのドメイン イベント ルール エンジン
- アラーム: アラーム発生·履歴
- データポイント: タグ データ照会·分析
- 開発者ガイド: フロー エンジン: エンジンアーキテクチャ·ノード インターフェース·DB スキーマ