Skip to main content

flowx-migrate

End-to-end migration of a source orchestrator's pipelines (Azure Data Factory, Apache Airflow) to Databricks Lakeflow Jobs. Orchestrates discover, convert, and package phases in sequence.

インストールへ移動

ソース情報

リポジトリ
databricks-solutions/flowx
ソースの最終更新活動
2026年8月10日 18:17
検出された SKILL.md の言語
英語
スター
4
フォーク
1

インストール方法

デフォルトでは、最初にソースを確認する Prompt が選択されています。直接コマンドに切り替えるか、ローカルコピーをダウンロードすることもできます。

ソースファイルを確認

インストールを決める前に、SKILL.md と SkillsMP に表示されている付属ファイルをお読みください。

ファイルエクスプローラー
2 ファイル

SKILL.md を表示中

SKILL.md
ソースの指示 · 読み取り専用プレビュー
name
flowx-migrate
description
End-to-end migration of a source orchestrator's pipelines (Azure Data Factory, Apache Airflow) to Databricks Lakeflow Jobs. Orchestrates discover, convert, and package phases in sequence.
triggers
["migrate pipelines","migrate ADF","migrate airflow","ADF to Databricks","airflow to Databricks","migrate to Lakeflow","migrate data factory"]
# End-to-End Source to Databricks Migration Orchestrate the complete migration of a source orchestrator's pipelines to Databricks Lakeflow Jobs via Declarative Automation Bundles. This skill runs all three phases in sequence: discover, convert, package. ## Context This is the top-level orchestration skill. It runs the full migration pipeline: 1. **Discover** — Parse the source's definitions into a typed inventory 2. **Convert** — Convert the source's tasks to Databricks IR (deterministic + agentic) 3. **Package** — Generate Databricks Declarative Automation Bundles for deployment Each phase builds on the output of the previous phase. The user is shown a summary and asked to confirm before proceeding to the next phase. ## Step 0 — Identify the source (required) Ask which orchestrator the user is migrating **from**, or infer it from the input: **Azure Data Factory / Fabric DF** (`--source adf`) or **Apache Airflow** (`--source airflow`). There is no default. Pass `--source <name>` to the discover and convert phases (package is source-independent). The discover/convert skills route to the matching `sources/<source>.md` guide for source-specific detail; the invocations below show ADF but apply to any source by swapping `--source` and the source path (`--adf-source-path` / `--airflow-source-path`, both aliases of `--source-path`). ## How to run this skill — MCP tools or venv CLI This skill orchestrates all three phases. Run the **`setup`** skill first if you haven't. There are two execution paths: ### MCP tools (Databricks Genie Code, or a local stdio registration) In Genie Code this is the **only** path — the phases run on the deployed `mcp-flowx` app, so there is **no venv, no `bootstrap.sh`, and no `.migration-venv`**. Run **no** `python3`/`$PY`/`bash` commands on this path; the `"$PY" -m …` snippets in the steps below are the **local-CLI fallback only**. Everything goes through the single **`flowx`** tool, one `command` per step: The hosted server **cannot read your workspace/volume files**, so pass the ADF JSON **inline** as `adf_definitions` (a mapping of relative path → ARM JSON content mirroring the Git-export layout — `pipeline/…`, `dataset/…`, `linkedService/…`, `trigger/…`). You read those files and supply them. The **recommended Genie path is a single `migrate` call**, which avoids re-sending the payload per phase: ``` flowx(command="migrate", parameters={ "source": "adf", "adf_definitions": {"pipeline/Foo.json": {...}, "linkedService/Bar.json": {...}, ...}, "output_dir": ..., "catalog": ..., "schema": ..., "pipeline": "<optional>"}) ``` **`migrate` is interactive when configuration options exist.** After translating, if the pipeline raises any configuration options (e.g. how to handle an `activity_and_notify` motif, metadata-driven bulk-copy consolidation, non-Databricks task compute), it **does not package** — it returns the **full option schema once**: `{"status": "needs_input", "pending_options": [{"pipeline_name", "options": [{option_id, prompt, rationale, choices, free_text, default, show_when}, …]}, …], "report_path", "output_dir"}`. You drive the whole chain locally — no per-answer round trip: 1. **Ask an option only when its `show_when` is satisfied** — every clause `{option_id, in:[values]}` must match an answer you've already collected (empty `show_when` = always ask). So `notify_slack_url` surfaces only after `notify_destination=slack`; the metadata-driven `access`/`size`/`lookup_tool` chain only after `metadata_driven_consolidate=consolidate`. Present each option's `prompt`/`rationale`/`choices`; honor the `default`. 2. **Validate** each answer against `choices` (a `free_text` option accepts any value); collect picks as `"option_id=value"` strings. Run any data action inline (e.g. the lookup query when `metadata_driven_lookup_tool=have`). 3. When **every applicable option** is answered, call `migrate` **once** with the same parameters plus `"answers": ["option_id=value", …]`. It applies them and packages (`"status": "completed"`). To accept all defaults and skip the prompts, pass `"interactive": false`. (Re-calling `migrate` with `answers` reuses the existing report and skips re-running discover/convert.) > **Large factories (hundreds–thousands of pipelines): do not inline.** Inline `adf_definitions` > passes through your context window and is capped (~5 MB). Instead point the server at the source by > reference: either stage the ADF export to a **UC Volume** and pass > `"adf_volume_path": "/Volumes/cat/sch/adf_export"` (read via the SDK Files API), or pass > `"adf_workspace_path": "/Workspace/Shared/adf_export"` for an ADF Git folder already in the workspace > (read via the SDK Workspace API). For output, pass `"output_volume_path": "/Volumes/cat/sch/dab"` > **or** `"output_workspace_path": "/Workspace/Shared/dab"` so the generated bundle is written to that > target via the SDK (returned as `bundle_uploaded` instead of inline `bundle`). Grant the > `mcp-flowx` app's service principal read on the source and write on the output target. For step-by-step control, run the commands in order (the app reuses `output_dir` across calls, so only `discover` needs the source input). `source` ("adf" | "airflow") is required for discover/convert and for `inputs discover`/`inputs convert`; for Airflow, swap `adf_definitions` for `airflow_source_path`. `merge_agentic` is ADF-only. `package` and `inputs package` are source-independent: ``` flowx(command="inputs", parameters={"phase": "discover", "source": "adf"}) # source req for discover/convert flowx(command="discover", parameters={"source": "adf", "adf_definitions": {...}, "output_dir": ..., "pipeline": ...}) flowx(command="convert", parameters={"source": "adf", "output_dir": ..., "pipeline": ...}) flowx(command="merge_agentic", parameters={"source": "adf", "report_path": ..., "agentic_results_dir": ..., "output_path": ...}) # ADF only, if agentic results flowx(command="inspect", parameters={"report_path": ...}) flowx(command="apply_answers", parameters={"report_path": ..., "answers": [...], "output_dir": ...}) flowx(command="package", parameters={"output_dir": ..., "catalog": ..., "schema": ...}) flowx(command="record_results", parameters={...}) / flowx(command="install_dashboard", parameters={...}) ``` For an Airflow report with eligible leaf placeholders, use the `flowx-resolve-airflow-gaps` skill between convert and package. It calls `resolve_agentic` with `action="prepare"`, stages one or more provider candidates, and applies only the gap fingerprints the user explicitly accepts. Package must then receive `<output_dir>/.work/translation_report.agentic.json` as `report_path`. The server's `output_dir` is ephemeral and not reachable from your workspace, so **have `migrate`/ `package` write the DAB to the target via the SDK** — pass `"output_volume_path": "/Volumes/…"` or `"output_workspace_path": "/Workspace/…"` and the bundle is uploaded there (returned as `bundle_uploaded`). Only when neither is set is the bundle returned inline as `bundle = {"files": {relpath: text,…}, "truncated": [...]}` (small bundles), which you must then persist yourself. Either way the user ends up with the DAB to validate and deploy. Each call returns a structured result (summaries / file trees); use those in place of reading files. Wherever a step below shows `"$PY" -m flowx.adapter <cmd> …`, call `flowx(command="<cmd>", parameters={...})` instead. > `databricks bundle validate` / `deploy` of the *generated* bundle is still a user-driven CLI step > (web terminal / local / CI-CD); present the bundle for review. ### venv CLI (local, no MCP server) Ensure the venv exists (`setup` Path B / `bootstrap.sh`), then run the commands below with the venv interpreter (from the marker file `<plugin_dir>/.migration-venv`) and `src/` on `PYTHONPATH` (use `$PY` anywhere a command shows `python3`): ```bash export PYTHONPATH="<plugin_dir>/src" PY="$(cat <plugin_dir>/.migration-venv)" "$PY" -m flowx.adapter inputs discover --source adf # or --source airflow ``` If Python or pip is missing, `bootstrap.sh` prints a warning telling the user what to install — relay it and stop until they have Python 3.12+ and pip. ## Workflow Follow these steps in order: ### Step 0 — Gather phase inputs via the adapter Before invoking discover, run the adapter inputs subcommand once per phase so the agent surfaces the matching free-text prompts: ```bash "$PY" -m flowx.adapter inputs discover --source adf # or --source airflow "$PY" -m flowx.adapter inputs convert --source adf # or --source airflow "$PY" -m flowx.adapter inputs package # source-independent ``` Each response carries the options for that phase plus their descriptions and defaults. Collect answers from the user (or accept the defaults) and thread the values into the downstream CLI calls. All phases share **one** migration `<output_dir>` (default `./flowx_output`). ### Step 1 — Gather inputs Ask the user for all required inputs upfront: | Parameter | Description | Required | Default | |---|---|---|---| | ADF source path | UC volume path or local directory with ADF JSON files | Yes | — | | Output directory | Single shared root for all flowx output (bundle + `metadata/`) | No | `./flowx_output` | | Target catalog | Unity Catalog catalog for tables/volumes | No | `main` | | Target schema | Schema within the catalog | No | `default` | | Bundle name | Name for the generated DABs project | No | derived from pipelines | Example prompt: > To migrate your ADF pipelines, I need: > 1. Where are your ADF JSON exports? (UC volume path like `/Volumes/main/default/adf_export` or local directory) > 2. Where should I write the output? (default: `./flowx_output/`) > 3. What target catalog and schema? (default: `main.default`) ### Step 2 — Phase 1: Discover Invoke the `flowx:flowx-discover` skill with `--source <source>`, the source path, and `--output-dir <output_dir>` (the shared migration dir). Discover routes to its `sources/<source>.md` guide and writes `<output_dir>/metadata/{inventory.json, profile_report.csv}` (plus `<pipeline>.arm.json` for ADF). Wait for discover to complete and present the inventory summary: ``` Phase 1: Discover — Complete ========================== Pipelines parsed: 12 Total activities: 47 Deterministic: 35 (74.5%) Agentic: 10 (21.3%) Unsupported: 2 ( 4.3%) Coverage: 95.7% ``` ### Step 3 — Checkpoint: confirm proceed Ask the user to review the inventory and confirm before continuing: > The discover phase found 47 activities across 12 pipelines. 95.7% have a translation path (74.5% deterministic, 21.3% agentic). 2 activities are unsupported and will need manual handling. > > Proceed to the translation phase? (yes/no) If the user says no, explain the options: - Re-run discover with a different source directory - Review `<output_dir>/metadata/inventory.json` (and `profile_report.csv`) to understand unsupported activities and pipeline complexity - Manually classify activities before proceeding If the user says yes, proceed to step 4. ### Step 4 — Phase 2: Convert Invoke the `flowx:flowx-convert` skill with: - `--source <source>`: the same source discover used - Source path: the original source path (same one discover used) - Output dir: the same shared `<output_dir>` (convert writes its report to `<output_dir>/.work/`) Wait for the translation to complete and present the summary: ``` Phase 2: Convert — Complete ============================= Deterministic translated: 35 (74.5%) Agentic translated: 8 (17.0%) Failed: 4 ( 8.5%) Overall coverage: 91.5% ``` ### Step 5 — Present translation details Show the user: 1. What was translated deterministically (bulk — just counts by type) 2. What was translated via agentic skills (list each with the skill used) 3. What failed and why (list each with the failure reason) For failures, suggest: - Manual notebook creation - Retry with additional context - Skip and add placeholder ### Step 5.1 — Gather just-in-time translation configuration Run `inspect` **once** to get the full option schema (every option carries a `show_when` condition), then drive the chain locally — ask an option only when its `show_when` clauses are all satisfied by the answers collected so far; never re-run `inspect` per follow-up. When the user opts to consolidate a metadata-driven motif and the agent has a database tool, run the lookup query to get CSV rows; otherwise prompt the user for a CSV file path or literal CSV string. Apply everything in **one** `modify` call (all `--answer OPTION_ID=VALUE` flags, plus `--lookup-csv "<csv-or-path>"` when needed — no intermediate JSON file). #### Legacy flow details Before bundle generation, run the adapter inspect CLI on the translation report to surface any pipeline-modifier options the IR raises: ```bash "$PY" -m flowx.adapter inspect <output_dir>/.work/translation_report.json ``` `inspect` returns the full option tree at once. Ask each option only when its `show_when` clauses (`{option_id, in:[values]}`) are all satisfied by the answers collected so far (empty = always); present its rationale, choices, and affected task keys. Then apply all collected answers in one `modify` call as repeatable `--answer OPTION_ID=VALUE` flags: ```bash "$PY" -m flowx.adapter modify \ <output_dir>/.work/translation_report.json \
GitHubで見る
この SKILL.md は非常に大きいため、SkillsMP では最初のセクションだけを表示しています。 GitHubで見る