← FlyteCONTENT HISTORY

Update to Flyte

Snapshot Sep 30, 2026 · 22:59 UTC · version 1.0.1

Collection source: not recorded for this historical snapshot.

WHAT CHANGED · RULE-BASED ANALYSIS

First saved snapshot

No earlier snapshot is available to establish a change.

Compare saved observations

Download comparison JSON
Full technical diff · 0 changed fields
Full snapshot data
{
  "description": "Entry-point orchestrator for porting Flyte 1 (flytekit) code to Flyte 2 (flyte). Explains the v1 to v2 shift, the terminology mapping, a recommended migration strategy, hybrid v1/v2 pipelines, and routes to sibling migration skills. Use when the user wants to migrate, port, or upgrade Flyte 1 (flytekit) code to Flyte 2. Trigger words are migrate, flytekit, v1 to v2, port, upgrade, convert workflow.",
  "included_files": [],
  "name": "flyte-migrate",
  "skill_md_contents": "---\nname: flyte-migrate\ndescription: Entry-point orchestrator for porting Flyte 1 (flytekit) code to Flyte 2 (flyte). Explains the v1 to v2 shift, the terminology mapping, a recommended migration strategy, hybrid v1/v2 pipelines, and routes to sibling migration skills. Use when the user wants to migrate, port, or upgrade Flyte 1 (flytekit) code to Flyte 2. Trigger words are migrate, flytekit, v1 to v2, port, upgrade, convert workflow.\n---\n\n# Flyte 1 to 2 Migration Skill\n\nThis is the entry point for migrating a Flyte 1 (`flytekit`) codebase to Flyte 2 (`flyte`). Flyte 2 is a fundamental shift: there is no `@workflow` decorator, everything is a `@env.task`, orchestration runs as real Python at runtime, and parallelism is expressed with `asyncio`. This skill explains the overall shift, gives a recommended migration strategy, covers hybrid v1/v2 pipelines during the transition, and routes to the sibling skills that handle each theme in depth.\n\n## Grounding References\n\n| Resource | URL |\n|---|---|\n| Migration guide | https://www.union.ai/docs/v2/flyte/user-guide/migration/flyte-2/ |\n| Official docs | https://www.union.ai/docs/v2/flyte |\n| Docs index (LLMs) | https://www.union.ai/docs/v2/flyte/llms.txt |\n| SDK API reference | https://www.union.ai/docs/v2/union/api-reference/flyte-sdk/ |\n| Example code | https://github.com/unionai/unionai-examples |\n| Flyte MCP tools | Available via `flyte-mcp` server |\n\n## The overall v1 to v2 shift\n\nTwo conceptual shifts motivate almost every change — **pure Python execution** and the **asynchronous model** — after which most migrations come down to a couple of mechanical moves.\n\n- **`flytekit` (package) becomes `flyte`.** Imports change from `import flytekit` to `import flyte`.\n- **`pyflyte` (CLI) becomes `flyte`.** The command-line tool was renamed.\n- **Everything is a task.** In Flyte 1, `@workflow` functions were constrained to a DSL subset of Python that compiled to a static DAG. In Flyte 2 there is **no `@workflow` decorator**: everything is a `@env.task`, and a \"workflow\" is simply a task that calls other tasks. Loops, conditionals, and `try`/`except` work anywhere.\n- **Async is the parallelism model.** Flyte 2 is built on `asyncio`, with the Flyte orchestrator acting as the event loop, scheduling awaited tasks across distributed infrastructure. `await` signals where a task can be scheduled in parallel, and `asyncio.gather` tells the orchestrator that a set of tasks are independent.\n\n### Simplified API mapping\n\n| Use case | Flyte 1 | Flyte 2 |\n| --- | --- | --- |\n| Environment management | `N/A` | `TaskEnvironment` |\n| Perform basic computation | `@task` | `@env.task` |\n| Combine tasks into a workflow | `@workflow` | `@env.task` |\n| Create dynamic workflows | `@dynamic` | `@env.task` |\n| Fanout parallelism | `flytekit.map` | Python `for` loop with `asyncio.gather` |\n| Conditional execution | `flytekit.conditional` | Python `if-elif-else` |\n| Catching workflow failures | `@workflow(on_failure=...)` | Python `try-except` |\n\n## Terminology and concept mapping\n\nSeveral Flyte 1 concepts were renamed or reshaped in Flyte 2. The table below maps the ones you'll meet most often.\n\n| Flyte 1 | Flyte 2 | Notes |\n|---|---|---|\n| `flytekit` (package) | `flyte` (package) | The Python SDK was renamed; imports change from `import flytekit` to `import flyte`. |\n| `pyflyte` (CLI) | `flyte` (CLI) | The command-line tool was renamed. |\n| `@task` / `@workflow` / `@dynamic` | `@env.task` | A single task decorator off a `flyte.TaskEnvironment`. Workflows and dynamic tasks are no longer distinct constructs: everything is a task, and orchestration is plain Python. |\n| `map_task()` | `flyte.map()` | Plus `asyncio.gather()` for async fan-out. |\n| `conditional()` | native `if` / `elif` / `else` | Branching is now ordinary Python control flow, not a DSL. |\n| `ImageSpec` | `flyte.Image` | Container image definition. |\n| `current_context()` | `flyte.ctx()` | Runtime context access. |\n| `FlyteFile` / `FlyteDirectory` | `flyte.io.File` / `flyte.io.Dir` | Offloaded file and directory references. |\n| `StructuredDataset` | `flyte.io.DataFrame` | Offloaded tabular data. |\n| `LaunchPlan` | `flyte.Trigger` | Scheduling and parameterized entry points. |\n| `CronSchedule` | `flyte.Cron` | Cron-based scheduling, used with a `flyte.Trigger`. |\n| Decks (`enable_deck=True`) | Reports (`report=True`) | Custom HTML rendered in the UI during/after a run. |\n\n## The two mechanical changes behind (almost) every migration\n\nMost of a migration comes down to two moves.\n\n### 1. Move task configuration into a `TaskEnvironment`\n\nInstead of configuring the image, resources, and caching on each task decorator, configure them once on a `flyte.TaskEnvironment` and share it across tasks:\n\n```python\nenv = flyte.TaskEnvironment(\n    name=\"training\",\n    image=flyte.Image.from_debian_base().with_pip_packages(\"scikit-learn\", \"pandas\"),\n    resources=flyte.Resources(cpu=\"2\", memory=\"4Gi\"),\n    cache=\"auto\",\n)\n```\n\n### 2. Replace `@task` / `@workflow` / `@dynamic` with `@env.task`\n\nEvery decorated function becomes an `@env.task`. There is no separate workflow or dynamic construct: a \"workflow\" is simply a task that calls other tasks, and orchestration is plain Python. The `env` in `@env.task` is just the variable you assigned your `TaskEnvironment` to — name it whatever you like.\n\n## Package imports\n\nThe package is renamed from `flytekit` to `flyte`, and the workflow/dynamic/map_task imports disappear.\n\n### Flyte 1\n\n```python\nimport flytekit\nfrom flytekit import task, workflow, dynamic, map_task\nfrom flytekit import ImageSpec, Resources, Secret\nfrom flytekit import current_context, LaunchPlan, CronSchedule\n```\n\n### Flyte 2\n\n```python\nimport flyte\nfrom flyte import TaskEnvironment, Resources, Secret\nfrom flyte import Image, Trigger, Cron\n```\n\n## Before and after: pure Python execution\n\n### Flyte 1\n\n```python\nimport flytekit\n\nimage = flytekit.ImageSpec(\n    name=\"hello-world-image\",\n    packages=[\"requests\"],\n)\n\n@flytekit.task(container_image=image)\ndef mean(data: list[float]) -> float:\n    return sum(list) / len(list)\n\n@flytekit.workflow\ndef main(data: list[float]) -> float:\n    output = mean(data)\n\n    # ❌ performing trivial operations in a workflow is not allowed\n    # output = output / 100\n\n    # ❌ if/else is not allowed\n    # if output < 0:\n    #     raise ValueError(\"Output cannot be negative\")\n\n    return output\n```\n\n### Flyte 2\n\n```python\nimport flyte\n\nenv = flyte.TaskEnvironment(\n    \"hello_world\",\n    image=flyte.Image.from_debian_base().with_pip_packages(\"requests\"),\n)\n\n@env.task\ndef mean(data: list[float]) -> float:\n    return sum(data) / len(data)\n\n@env.task\ndef main(data: list[float]) -> float:\n    output = mean(data)\n\n    # ✅ performing trivial operations in a workflow is allowed\n    output = output / 100\n\n    # ✅ if/else is allowed\n    if output < 0:\n        raise ValueError(\"Output cannot be negative\")\n\n    return output\n```\n\n## Quick reference: minimal Flyte 2 module\n\n```python\nimport asyncio\nimport flyte\n\n# 1. Define an image\nimage = (\n    flyte.Image.from_debian_base(python_version=(3, 11))\n    .with_pip_packages(\"pandas\", \"numpy\")\n)\n\n# 2. Create a TaskEnvironment\nenv = flyte.TaskEnvironment(\n    name=\"my_env\",\n    image=image,\n    resources=flyte.Resources(cpu=\"1\", memory=\"2Gi\"),\n)\n\n# 3. Define tasks\n@env.task\nasync def process(x: int) -> int:\n    return x * 2\n\n# 4. Define the entrypoint task\n@env.task\nasync def main(items: list[int]) -> list[int]:\n    results = await asyncio.gather(*[process(x) for x in items])\n    return list(results)\n\n# 5. Run it\nif __name__ == \"__main__\":\n    flyte.init_from_config()\n    run = flyte.run(main, items=[1, 2, 3, 4, 5])\n    print(run.url)\n    run.wait()\n```\n\n```bash\n# CLI\nflyte run my_module.py main --items '[1,2,3,4,5]'   # remote (default)\nflyte run --local my_module.py main --items '[1,2,3,4,5]'\nflyte deploy my_module.py my_env\n```\n\n## Recommended migration strategy\n\nMigrations rarely happen all at once. Work incrementally and lean on hybrid pipelines while the transition is in progress.\n\n1. **Assess the codebase.** Inventory every `@task`, `@workflow`, `@dynamic`, and `map_task`; the images (`ImageSpec`), resources, and secrets; the control-flow constructs (`conditional`, `on_failure`, `>>`); the data types (`FlyteFile`, `FlyteDirectory`, `StructuredDataset`); and any schedules (`LaunchPlan`, `CronSchedule`).\n2. **Establish the `TaskEnvironment`(s).** Group tasks by their image/resource/cache needs and define a `flyte.TaskEnvironment` for each group. This is mechanical change #1 and unblocks everything else.\n3. **Port leaf tasks first, then orchestration.** Convert atomic compute tasks (`@task` → `@env.task`), then rebuild the `@workflow`/`@dynamic` orchestration as plain-Python driver tasks that call them.\n4. **Migrate control flow and I/O.** Replace `conditional()` with `if`/`elif`/`else`, `on_failure` with `try`/`except`, `map_task` with `flyte.map` / `asyncio.gather`, and the `FlyteFile`/`FlyteDirectory`/`StructuredDataset` types with their `flyte.io` equivalents.\n5. **Update config, CLI, and schedules.** Swap `pyflyte` for `flyte`, migrate config files, and convert `LaunchPlan`/`CronSchedule` to `flyte.Trigger`/`flyte.Cron`.\n6. **Run hybrid during the transition.** Keep unported v1 workflows callable via bridge tasks (see below) until every piece is on v2.\n\n### Sibling skills to route to\n\nMigrate by theme. Start with tasks and workflows, then jump to whatever the workload needs:\n\n- **`flyte-migrate-tasks-workflows`** — the structural shift: `@task`/`@workflow` → `@env.task`, sequential ordering, nested \"subworkflows\", and the `@task` → `TaskEnvironment` parameter mapping.\n- **`flyte-migrate-config`** — moving image/resources/cache to the `TaskEnvironment`, GPUs, secrets, caching, scheduling with triggers, and the `pyflyte` → `flyte` command/config-file changes.\n- **`flyte-migrate-control-flow`** — `conditional()` and `@dynamic` become plain Python `if`/loops, `on_failure` becomes `try`/`except`, and `map_task` → `flyte.map` / `asyncio.gather`.\n- **`flyte-migrate-data-io`** — `FlyteFile`/`FlyteDirectory` → `flyte.io.File`/`Dir`, `StructuredDataset` → `flyte.io.DataFrame`, dataclasses, and ETL patterns.\n- **`flyte-migrate-ml`** — small-model training, hyperparameter optimization, deep learning, batch inference, and end-to-end pipelines.\n\n## Hybrid v1 and v2 pipelines\n\nFor a while you'll have Flyte 1 and Flyte 2 workloads running side by side, and you'll want them to call each other: a Flyte 1 workflow that kicks off a newly ported Flyte 2 task, or a Flyte 2 task that triggers a workflow that hasn't been migrated yet.\n\nYou can bridge the two in both directions. The idea is the same each way: one task installs **both** SDKs, authenticates to the **other** control plane, fetches the entity it wants to run, and launches it. Keep the bridging task lightweight and focused on orchestration.\n\n### Running a Flyte 2 task from a Flyte 1 workflow\n\nThe bridge is a single Flyte 1 task that runs the Flyte 2 client. Give it an image with **both** `flytekit` and `flyte` installed, provide a Flyte 2 API key as a secret, authenticate inside the task with `flyte.init_from_api_key()`, fetch the deployed task with `flyte.remote.Task.get(...)`, and run it.\n\n```python\nimport flytekit\nfrom flytekit import task, workflow, ImageSpec, Secret, current_context\n\n# The bridge image needs BOTH the v1 (flytekit) and v2 (flyte) SDKs.\nbridge_image = ImageSpec(\n    name=\"v1-to-v2-bridge\",\n    packages=[\"flytekit\", \"flyte\"],\n)\n\n@task(\n    container_image=bridge_image,\n    secret_requests=[Secret(group=\"flyte\", key=\"flyte_api_key\")],\n)\ndef launch_v2_from_v1(x: int) -> str:\n    import flyte\n    import flyte.remote\n\n    # Authenticate to the Flyte 2 control plane with the API key.\n    # Option A: read the mounted secret and pass it explicitly.\n    api_key = current_context().secrets.get(group=\"flyte\", key=\"flyte_api_key\")\n    flyte.init_from_api_key(api_key=api_key)\n\n    # Option B: if FLYTE_API_KEY is set as an env var, no argument is needed:\n    #     flyte.init_from_api_key()\n\n    # Fetch the deployed Flyte 2 task and run it.\n    remote_v2_task = flyte.remote.Task.get(\n        \"my_v2_env.process\",\n        auto_version=\"latest\",\n    )\n    run = flyte.run(remote_v2_task, x=x)\n    run.wait()  # optional: block until the v2 run finishes\n    return run.url\n\n@workflow\ndef main(x: int) -> str:\n    return launch_v2_from_v1(x=x)\n```\n\nThe referenced Flyte 2 task (`my_v2_env.process` above) must be **deployed** before the bridge runs. Use `flyte.init_from_api_key()` here — do **not** use `flyte.init_from_config()`, which reads a `config.yaml` that has no API-key field.\n\n### Running a Flyte 1 workflow from a Flyte 2 task\n\nThe reverse works the same way: a Flyte 2 task installs the Flyte 1 client and uses `FlyteRemote` to launch a Flyte 1 workflow.\n\n```python\nimport flyte\n\nenv = flyte.TaskEnvironment(\n    name=\"v2_to_v1_bridge\",\n    # The image needs the Flyte 1 client installed.\n    image=flyte.Image.from_debian_base().with_pip_packages(\"flytekit\"),\n    # Supply credentials for the Flyte 1 control plane (config or API key).\n    secrets=[flyte.Secret(key=\"v1_client_secret\", as_env_var=\"V1_CLIENT_SECRET\")],\n)\n\n@env.task\nasync def launch_v1_from_v2(x: int) -> str:\n    from flytekit.remote import FlyteRemote\n    from flytekit.configuration import Config\n\n    # Point the client at your Flyte 1 cluster.\n    remote = FlyteRemote(\n        config=Config.for_endpoint(endpoint=\"my-v1-cluster.example.com\"),\n        default_project=\"flytesnacks\",\n        default_domain=\"development\",\n    )\n\n    # Fetch the deployed Flyte 1 workflow and execute it.\n    wf = remote.fetch_workflow(name=\"my_v1_module.main\", version=\"v1.2.3\")\n    execution = remote.execute(wf, inputs={\"x\": x}, wait=True)\n    return execution.id.name\n```\n\n### Hybrid considerations\n\n- **Both SDKs in one image.** The bridging task installs `flytekit` and `flyte` together. Pin versions and watch for dependency conflicts; keep the bridge image minimal.\n- **Deploy the callee first.** For v1→v2, the Flyte 2 task must be deployed (`flyte deploy`) before `flyte.remote.Task.get()` can resolve it. For v2→v1, the Flyte 1 workflow must be registered on its cluster.\n- **Wait vs. fire-and-forget.** Both `run.wait()` (v2) and `execute(..., wait=True)` (v1) block until the launched run finishes. Omit them to launch and return immediately.\n- **Credentials cross a boundary.** The bridge authenticates to a *different* control plane than the one it runs on. Store the API key or client credentials as a secret — never hard-code them.\n- **Keep the bridge lightweight.** Like any orchestrating task, it should mostly launch and assemble results rather than do heavy compute.\n\n## Gotchas\n\nFlyte 2 lets each Python task act as its own engine, launching sub-tasks and assembling their outputs. That flexibility warrants some caveats.\n\n### Common gotchas\n\n- **`flyte.map` returns a generator.** Wrap it in `list()` to materialize results, unlike `map_task` which returned a list directly.\n- **`memory`, not `mem`.** The `Resources` parameter was renamed, and there are no separate `requests`/`limits` — a single value serves as both.\n- **GPUs use a `\"T4:1\"` string.** Type and count are combined; the separate `accelerator=` argument is gone.\n- **Image, resources, and cache live on the `TaskEnvironment`.** Set them once at the env level instead of repeating them on every task decorator.\n- **`current_context()` is gone.** Read secrets from environment variables and use `flyte.ctx()` for runtime context.\n- **The `>>` ordering operator is gone.** Sequential (sync) calls and sequential `await`s are naturally ordered.\n- **Retries no longer have a platform cap.** In Flyte 1 the control plane capped attempts at 3; in Flyte 2 total attempts equal `retries + 1`. Audit any large `retries` values before deploying.\n- **You can only `await` async tasks.** Call a sync task from an async context with `.aio()`.\n- **Pick an entrypoint task name.** There's no `@workflow`, so the top-level task is just a task (commonly `main`); run it with `flyte run module.py main`.\n- **Type annotations are more lenient.** Flyte 2 will pickle untyped I/O rather than rejecting it at registration.\n- **Keep orchestration lightweight.** A task that calls other tasks acts as a driver pod. Avoid heavy CPU work in it.\n\n## Anti-Patterns\n\n1. **Don't introduce non-determinism into orchestration.** When a task launches another task, a new Action ID is determined as a hash of the inputs and task definition — consistent hashing is what makes recovery and replay work. Branching on `datetime.now()` or other non-deterministic values breaks that guarantee: on retry, a *different* downstream task may get kicked off. If non-determinism is unavoidable, decorate sub-task functions with `@trace` for fine-grained checkpointing and observability.\n2. **Don't do heavy compute in a driver task.** When a task runs other tasks and assembles their outputs, it becomes a driver pod (work that Flyte Propeller did in v1). A CPU-bound function between two `await`s makes the driver pod hang and slows downstream kickoff. Keep parent tasks focused on orchestration:\n\n```python\n@env.task\nasync def t_main():\n    await t1()\n    local_cpu_intensive_function()  # ❌ blocks the driver pod between t1 and t2\n    await t2()\n```\n\n3. **Don't rely on global state across tasks.** Each task runs in its own isolated container; globals are not carried across task containers. Any state that must persist has to be reconstructable through repeated deterministic execution.\n4. **Don't materialize huge in-memory I/O between tasks.** Outputs are materialized in the parent pod's memory, so passing a 1 GB `list[float]` requires the pod to hold all of it, risking OOM. Use `flyte.io.File`, `flyte.io.Dir`, and `flyte.io.DataFrame` — they're materialized only as pointers to offloaded data, so their memory footprint stays low.\n5. **Don't skip type hints at the \"workflow\" level.** The top-level task now runs at runtime, so the system can't guarantee type safety across the DAG the way the v1 DSL did. Use Python type hints and a type checker like `mypy` at all levels, including the top-most task.\n"
}

SHA-256 of public snapshot: b393f27b57f96a6a975ce60cd602154c59fa53a3053dd457107300ab899454e5