フロー構築
データの取り込み・変換・出力をステップで組み合わせて自動化します。複雑な統合ツールではなく、シンプルな UI で AI とのデータ連携に特化しています。
仕組み
- ステップを上から順に実行します。データは**行の集まり(list[dict])**として各ステップを流れます。
- ステップは大きく 取得(source)→ 加工(transform)→ 出力(load / action) に分かれます。
- 外部サービスの認証情報は 接続管理(ApiConnection) に保存し、ステップから参照します(接続は作成した本人のみが利用できます。管理者も他人の接続は使えません)。
- ステップの利用にロール制限はありません。シェルコマンド・rsync・Python 関数を含め、すべてのステップを全ユーザーが利用できます。
- 値の差し込みにはテンプレートが使えます(Jinja / GitHub Actions と同じ見た目): 今の行の列は
{{ 列名 }}/{{ order.customer.name }}/{{ items[0].id }}、前のステップの出力は{{ steps["ステップ名"][0].列名 }}、実行情報は{{ sys.today }}/{{ sys.now }}/{{ sys.count }}/{{ workflow.name }}/{{ workflow.run_id }}/{{ env.NAME }}/{{ loop.index }}。steps/workflow/env/sys/loop/dataは予約語で、行の列より優先されます。旧形式({{ $json.列名 }}/{{ $node["…"].json }}/{{ $today }}など)も引き続き同じ意味で使えます。 - テンプレートが式 1 個だけ(地の文なし)のとき、「値をセット」では list / dict が文字列化されず構造のまま列に入ります。文と混ぜると JSON 文字列になります。
- ステップ名を変更しても、他ステップの
steps["…"]参照は自動で追随します。 - AI 作成: エディタの「AI 作成」パネルに言葉で指示すると、パイプライン全体(ステップと接続線)を AI が下書き・修正します。変更案を承認するとエディタに反映され、保存で確定します。
共通の設定(全ステップ)
各ステップには、失敗時の挙動を制御する共通オプションがあります。
| 項目 | 既定 | 説明 |
|---|---|---|
| リトライ | off | 失敗時に再試行する |
| 最大リトライ回数 | 3 | 再試行の回数 |
| リトライ間隔(ms) | 1000 | 再試行までの待ち時間 |
| エラー時の挙動 | stop | stop(中断)/ continue(無視して継続)/ continue_with_error_output(エラー行を出力して継続) |
オペレータ一覧
ファイル連携
| 表示名 | キー | 役割 |
|---|---|---|
| ファイル取得 | source_file | CSV / Excel / JSON / JSONL / PDF / TSV / DOCX / PPTX などを読込(フォルダ・glob・接続 URI・仮想パス対応) |
| ファイル保存 | load_file | CSV / TSV / JSON / JSONL / MD / Excel / gzip で書き出し |
| ファイル差分取得 | source_file_landed | 追加・変更されたファイルだけを取得(差分トリガ)。ファイルごとのスナップショット比較で判定するため、時計ズレや mtime を保持したコピーでも取りこぼしません。「削除も検知」を有効にすると消えたファイルを削除イベントとして出力(「RAG に保存」の upsert が削除を反映)。既定は指定フォルダの直下だけを見る。「サブフォルダも含める」で配下も対象。接続先フォルダを「ローカルへ同期」した同期フォルダは着荷フォルダとして扱われ、「更新」で取り込んだファイルごとに file.uploaded イベントが出る(イベント起動の pipeline がそのまま使える)。マイファイル / 共有 / 接続フォルダ(SFTP / SMB / S3 等)に対応 |
| 画像読み込み (VLM) | source_image_ai | 画像を VLM(失敗時 OCR)で説明文に変換 |
| rsync 取得 / 送信 | source_rsync / load_rsync | リモートと rsync で同期 |
API 連携
| 表示名 | キー | 役割 |
|---|---|---|
| REST API 取得 | source_rest | REST をページング・認証つきで取得(保存済みの「REST API」接続を利用可) |
| REST API 送信 | action_http_request | 各行ごとに HTTP リクエストを送信し、応答を行に付与 |
| カスタム MCP 実行 | action_mcp_call | カスタム MCP サーバーの tool を各行ごとに実行し、結果を行に付与(入力 0 行なら 1 回実行)。自分が作成したか、タグ共有され稼働中のサーバーが対象 |
| Webhook 受信 | source_webhook | Webhook 受信を有効化。Idempotency-Key / X-Delivery-Id 等の配送 ID(idempotency_header で指定可)が同じ再送は実行を重複させません |
| Webhook 送信 | action_webhook | Webhook URL へ署名・リトライつき送信 |
| メール送信 | action_email | SMTP でメール送信 |
| チャット通知送信 | action_notify | Slack / Teams / Chatwork / LINE WORKS / LINE / Discord / Telegram へ送信(送信先はプロバイダ選択式) |
SQL・データ蓄積
| 表示名 | キー | 役割 |
|---|---|---|
| SQL 取得 (SELECT) | source_sql | SELECT / WITH のみ。クエリ直書き or テーブル+条件。差分取得対応 |
| SQL 保存 (INSERT) | load_sql | append / replace / merge。テーブル自動作成・バッチ投入(PostgreSQL / MySQL / MariaDB 対応) |
| SQL 更新 (UPDATE) | update_sql | 主キーで行単位 UPDATE(PostgreSQL / MySQL / MariaDB 対応) |
| SQL 実行 (任意 SQL) | execute_sql | 任意の DDL / DML を実行 |
| BigQuery 取得 / 書込 | source_bigquery / load_bigquery | BigQuery をクエリ / ストリーム挿入(SA JSON 認証) |
| ClickHouse 取得 / 書込 | source_clickhouse / load_clickhouse | ClickHouse の読み書き |
| データセットから取得 | source_dataset | 内蔵のデータセットを読み込み(別ワークフローや AI 分析の保存結果を入力に) |
| データセットに保存 | load_postgres_dataset | 内蔵のデータセットに保存(ダッシュボード や BI キーで参照) |
| RAG に保存 | load_to_rag | 行をまとめて 1 文書としてナレッジに埋め込み |
| RAG から検索 | source_rag_retrieve | ナレッジから類似チャンクを取得 |
データ処理
| 表示名 | キー | 役割 |
|---|---|---|
| 行の絞り込み (WHERE) | transform_filter | 条件で行を残す / 除く |
| 列の編集 (SET) | transform_set | 列の設定 / 改名 / 削除 / 指定列だけ残す |
| 列名の変更 (AS) | transform_rename_keys | 全行のキーを一括改名 |
| 並べ替え (ORDER BY) | transform_sort | 昇順 / 降順 / ランダム(重複除去・件数制限も) |
| 件数制限 (LIMIT) | transform_limit | 先頭 / 末尾 N 件 |
| 重複削除 (DISTINCT) | transform_dedup | 指定列で重複除去 |
| 行に分割 (UNNEST) | transform_split | 1 行を複数行に展開 |
| 集計・グループ化 (GROUP BY) | transform_aggregate | グループ化 + 集約関数 |
| データ結合 (JOIN) | transform_merge | append / inner / left / outer 結合 |
| データ比較 (EXCEPT) | transform_compare | 2 つのデータの差分検出 |
| 型変換 (CAST) | transform_cast | integer / float / boolean / date / datetime / string |
| データ検証 (CHECK) | transform_validate | ルール検証(flag / reject) |
ユーティリティ(AI・コード)
| 表示名 | キー | 役割 |
|---|---|---|
| AI エージェント | ai_agent | map / filter / reduce / classify を LLM で実行 |
| Python 関数 | transform_code | サンドボックスで transform(data) を実行 |
| シェルコマンド | transform_execute_command | コマンド実行し標準出力を行に付与 |
| 日時操作 | transform_datetime | フォーマット / 加減算 / 抽出 / 差分 |
| 暗号化・ハッシュ | transform_crypto | hash / base64 / URL encode / hmac |
| Markdown / HTML | transform_markdown | Markdown ↔ HTML 変換 |
| XML 処理 | transform_xml | XML ↔ dict 変換 |
条件分岐・繰り返し
| 表示名 | キー | 役割 |
|---|---|---|
| 条件分岐 (IF) | flow_if | true / false に振り分け |
| 分岐 (SWITCH) | flow_switch | 名前つき出力へ多分岐 |
| ループ (FOR EACH) | flow_loop | バッチ分割で逐次 / 並列処理 |
| 待機 (WAIT) | transform_wait | 指定秒スリープ |
| 終了 (BREAK) | flow_break | 条件成立でループを抜ける |
外部システム連携
業務システム(ERP / SFA)・監視・製造 IoT と連携します。接続情報は接続管理で管理します。
業務システム(ERP / SFA / CRM)
専用ステップは持たず、「REST API 取得」「REST API 送信」+ 保存済みの「REST API」接続で繋ぎます(認証ヘッダ・OAuth2 client credentials・ページングに対応)。設定例:
| サービス | 接続の設定 | 取得ステップの設定 |
|---|---|---|
| kintone | Base URL https://<sub>.cybozu.com/k/v1、認証 api_key、ヘッダ名 X-Cybozu-API-Token | URL records.json、クエリ app=<ID>、data_path=records、ページング offset(limit 500。1 万件超は kintone 側で絞り込む) |
| Salesforce | Base URL https://<org>.my.salesforce.com、認証 oauth2(token URL /services/oauth2/token、接続アプリの client id / secret) | URL /services/data/v59.0/query、クエリ q=<SOQL>、data_path=records、ページング link(next_path=nextRecordsUrl) |
書き込みは「REST API 送信」で行単位に POST / PATCH します。
監視・ログ
| 表示名 | キー | 役割 |
|---|---|---|
| Elasticsearch 取得 / 書込 | source_elasticsearch / load_elasticsearch | クエリ検索 / _bulk 書込 |
| ログ取得 | source_logs | アプリケーションログを取得 |
| 各種メトリクス | source_pipeline_metrics ほか | パイプライン / チャット / システムの統計 |
製造・IoT(FA 連携) — 工場ラインの PLC / センサに直結します。
| 表示名 | キー | 役割 |
|---|---|---|
| OPC-UA 取得 / 書込 | source_opcua / load_opcua | OPC-UA サーバのノード値を読み書き(PLC データ取得の標準) |
| Modbus 取得 / 書込 | source_modbus / load_modbus | Modbus TCP のレジスタ / コイルを読み書き |
| MQTT 受信 / 送信 | source_mqtt / load_mqtt | MQTT ブローカを購読(有界収集)/ publish |
| SLMP / 三菱 取得 / 書込 | source_slmp / load_slmp | 三菱 PLC(MC プロトコル 3E / 4E)のデバイス読み書き ※実験的 |
| FINS / オムロン 取得 / 書込 | source_fins / load_fins | オムロン PLC(FINS/UDP)のメモリエリア読み書き ※実験的 |
S3 / GCS / SFTP / FTP / SMB / Box / Google Drive / OneDrive / Dropbox / Azure Blob / SharePoint などのファイル転送は、独立オペレータではなく ファイル取得 / 保存(source_file / load_file)に接続 URI を指定して扱います。接続の作り方はファイルを参照してください。
主要オペレータのパラメータ
AI エージェント(ai_agent)
| 項目 | 既定 | 説明 |
|---|---|---|
| 操作 | map | map(各行を変換)/ filter(残す行を判定)/ reduce(グループ集約)/ classify(分類) |
| モデル | 文脈の既定 | 使う LLM |
| プロンプト / システムプロンプト | — | 指示文 |
| 出力列 | ai_output(classify は _label) | 結果を入れる列 |
| 分類カテゴリ | — | classify のラベル候補。multi_label で複数付与 |
| temperature ほか | 0.7(classify は 0.1) | チャットと同じ生成パラメータ |
| 行データの自動添付 | オン | 行 JSON をプロンプトに自動添付(オフなら {{ 列名 }} 差し込みのみ) |
| 添付する列 | 全列 | 添付する列の絞り込み(トークン削減・列の除外) |
Python 関数(transform_code)
transform(data) を定義します。data は行の配列です。mode は all(全行を一度に)/ each(1 行ずつ)。サンドボックスで実行され、ファイル・ネットワーク・危険な属性アクセスは禁止、timeout_seconds(既定 30)で打ち切られます。
ファイル取得 / 保存(source_file / load_file)
取得は file_path(フォルダ末尾 /・glob・接続 URI・仮想パス myfile/・shared/・external/<label>)、format(未指定なら拡張子で判定)、encoding、delimiter、max_rows など。保存は file_path(必須)、format(既定 csv)、append、compress(.gz)。
SQL 取得 / 保存(source_sql / load_sql)
取得は接続(connectionId)+ query(SELECT / WITH)または table + columns / where / order_by / limit、差分取得(incremental)。差分取得は 2 方式: カーソル(cursor_column の前回最大値より大きい行だけを DB に絞らせる。lookback で基準を少し戻して遅れて届いた行を拾い直せる。更新日時や連番のある表向け、削除は取れない)と スナップショット(primary_key と行のハッシュを台帳に持ち、追加・変更・削除を突き合わせる。emit_deleted で消えた行を削除イベントとして出力し、「RAG に保存」の upsert が反映する。更新日時の無い表や削除を伝えたい表向け、SQL_SNAPSHOT_MAX_ROWS 以下の行数)。スナップショットの台帳はファイル差分取得と同じ仕組みで、失敗した実行では確定しません。保存は table(必須)、write_disposition(append / replace / merge)、primary_key(merge 用)、batch_size(既定 1000)。
RAG に保存 / 検索(load_to_rag / source_rag_retrieve)
保存は bot_id(必須)、mode(append / replace_all)、chunk_size(600)、chunk_overlap(100)。検索は bot_id(必須)、query、top_k(5)、similarity_threshold。
条件の演算子(filter / IF / break 共通)
== != > >= < <= / contains not_contains starts_with ends_with / regex / exists not_exists / is_empty is_not_empty。match で AND(all)/ OR(any)、case_sensitive の切り替え。
比較は値の型を見て行います。両辺が数値なら数値("1,000" > 999 は真)、日付なら日付(2026/09/04 == 2026-09-04 は真)、それ以外は文字列(既定は大文字小文字を無視)。片側だけが数値・日付の場合は従来どおり文字列として比べ、その旨を判定理由に残します。type(number / date / text / bool)で型を固定でき、value には {{ $json.列名 }} などの差し込み、field には a.b.c のネスト指定が使えます。実行履歴の各ステップに「条件の判定」(先頭 10 行の判定根拠)が表示されます。
集約関数(transform_aggregate)
count / sum / avg / min / max / count_unique / concat / first / last / collect。group_by 空で全体集約。
ステップ実行(プレビュー)の挙動
編集画面の「ステップ実行」は上流から順に実行して結果を表示しますが、外部に書き込む・送信するステップ(SQL 保存、ファイル保存、メール、Webhook 送信、通知、SaaS 登録など)は実行しません。該当ステップは「書込のためプレビューでは実行しません」と表示され、入力行がそのまま次のステップに渡ります。実際の書込は「実行」で行ってください。プレビューは PIPELINE_PREVIEW_TIMEOUT_SEC(既定 120 秒)で打ち切られます。
共通の詳細設定(全ステップ)
各ステップの設定下部「詳細設定」で、失敗時の再試行(回数・待ち時間)、失敗時の扱い(停止 / 続行 / エラー行を付けて続行)、入力 0 行でも実行するか、を設定できます。外部に書き込むステップでは二重書込を防ぐため再試行は行われません。「続行」で握った失敗があると、その実行では差分取得のカーソルやファイルのスナップショットは確定せず、次回に持ち越されます。
無人運転の挙動
- 停止中の取りこぼしを取り直す: 外部ディレクトリ監視は、サーバ停止中に置かれたファイルを起動時に走査して発火します。スケジュール実行も、停止中に逃した直近 1 回分を起動時に実行します。同じ内容・同じ配送を二重に実行しないよう、実行には冪等キーが付きます。
- 中断された実行は再実行: 再起動などで中断した実行は自動で待機列に戻ります(
PIPELINE_RUN_MAX_ATTEMPTS、既定 2 回まで)。 - 連続失敗は自動停止: 同じワークフローが連続で失敗すると(
PIPELINE_AUTO_PAUSE_AFTER_FAILURES、既定 5 回)自動起動(イベント / スケジュール / Webhook)を一時停止し、一覧に「自動停止中」と表示します。原因を直して手動実行すると解除されます。 - 差分は根拠つきで残る: ファイル着荷の差分(追加 / 更新 / 削除と前後の size・更新時刻)、SQL 差分取得のカーソル範囲、RAG 同期で削除したキーは実行履歴の「差分の根拠」で確認できます。処理済みの記録(カーソル・スナップショット)は、実行が成功しエラーのあるステップが無い場合にだけ確定します。失敗した分は次回に持ち越されます。
- 同サイズ・同時刻の上書き: ファイル着荷の
content_hashを有効にすると、内容のハッシュで変更を検知します(既定は無効)。
実行先(worker pool)— 重いワークフローを別プロセス / 別マシンへ
ワークフローの実行は本体に内蔵された常駐 worker が行いますが、AI 推論や大量データ処理など重いものは、別プロセスや別マシンの worker に逃がせます。仕組みは「待機列(DB)から worker が取りに来る」方式で、worker にポートはありません。
- 各ワークフローには 実行先(pool 名) があり、既定は
default(内蔵 worker)。変更できるのは管理者だけです(資源の配分のため。管理画面 > システム > 「ワークフローの実行先」で、稼働中の worker を見ながら付け替えます)。 - 外部 worker は本体と同じ
.env(同じ DB)でdigitalbase worker --pool <名前>を起動します(.envのPIPELINE_WORKER_POOLでも可)。名前が一致したワークフローの実行だけを取ります。 - その pool の worker が居ないと実行は「待機中」のまま止まります(管理画面に警告が出ます)。worker の生存は heartbeat(
PIPELINE_WORKER_HEARTBEAT_SEC)で判定し、別マシンの worker が落ちた実行も自動で待機列に戻ります。 - 外部 worker は SIGTERM で「実行中のものを終えてから」停止します。
代表的なパイプライン例
- 社内 DB をナレッジ化: SQL 取得 → AI エージェント(要約) → RAG に保存
- ファイル取込 → 蓄積: ファイル取得(CSV glob) → 型変換 → データ検証 → データセットに保存
- API → BigQuery: REST API 取得 → 絞り込み → 集計 → BigQuery 書込
- 差分監視 → 画像読取 → 通知: ファイル差分取得 → VLM画像取得 → チャット通知送信
- CRM 連携: REST API 取得(kintone) → データ比較(差分) → 条件分岐 → REST API 送信(Salesforce)
- 監視 → 分類 → メール: メトリクス取得 → AI エージェント(障害分類) → 絞り込み → メール送信