← FlyteCONTENT HISTORYWHAT CHANGED · RULE-BASED ANALYSIS
Update to Flyte
Snapshot Sep 30, 2026 · 22:59 UTC · version 1.0.1
Collection source: not recorded for this historical snapshot.
First saved snapshot
No earlier snapshot is available to establish a change.
Compare saved observations
Download comparison JSONFull technical diff · 0 changed fields
Full snapshot data
{
"description": "Migrates Flyte 1 branching, dynamic workflows, failure handling, and fan-out to native Flyte 2 Python. Use when migrating Flyte 1 branching, dynamic workflows, failure handling, or map_task/fan-out to Flyte 2. Trigger words: conditional, @dynamic, map_task, on_failure, branching, parallelism, fan-out, flyte.map, asyncio.gather.",
"included_files": [],
"name": "flyte-migrate-control-flow",
"skill_md_contents": "---\nname: flyte-migrate-control-flow\ndescription: \"Migrates Flyte 1 branching, dynamic workflows, failure handling, and fan-out to native Flyte 2 Python. Use when migrating Flyte 1 branching, dynamic workflows, failure handling, or map_task/fan-out to Flyte 2. Trigger words: conditional, @dynamic, map_task, on_failure, branching, parallelism, fan-out, flyte.map, asyncio.gather.\"\n---\n\n# Flyte 1 to 2 Migration: Control Flow and Parallelism\n\nFlyte 1 expressed branching, dynamic fan-out, and failure handling through DSL constructs (`conditional()`, `@dynamic`, `@workflow(on_failure=...)`) and `map_task`. In Flyte 2 these are all ordinary Python, because orchestration runs as real Python at runtime. Native `if`/`elif`/`else` replaces the conditional DSL, plain task loops replace `@dynamic`, `try`/`except` replaces `on_failure`, and `flyte.map` / `asyncio.gather` replace `map_task`.\n\n## Grounding References\n\n| Resource | URL |\n|---|---|\n| Migration guide (Control flow) | https://www.union.ai/docs/v2/flyte/user-guide/migration/flyte-2/control-flow/ |\n| Migration guide (Parallelism) | https://www.union.ai/docs/v2/flyte/user-guide/migration/flyte-2/parallelism/ |\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## Conditional Execution\n\nThe `conditional()` DSL becomes ordinary Python `if` / `elif` / `else` — for example, choosing a model based on dataset size.\n\n### Flyte 1\n\n```python\nfrom flytekit import task, workflow, conditional\n\n@task\ndef train_gradient_boosting(n_rows: int) -> str:\n return f\"trained gradient boosting on {n_rows} rows\"\n\n@task\ndef train_logistic_regression(n_rows: int) -> str:\n return f\"trained logistic regression on {n_rows} rows\"\n\n@workflow\ndef main(n_rows: int) -> str:\n # Pick the model based on dataset size.\n return (\n conditional(\"model_choice\")\n .if_(n_rows > 10_000)\n .then(train_gradient_boosting(n_rows=n_rows))\n .else_()\n .then(train_logistic_regression(n_rows=n_rows))\n )\n```\n\n### Flyte 2\n\n```python\nimport flyte\n\nenv = flyte.TaskEnvironment(name=\"conditional\")\n\n@env.task\ndef train_gradient_boosting(n_rows: int) -> str:\n return f\"trained gradient boosting on {n_rows} rows\"\n\n@env.task\ndef train_logistic_regression(n_rows: int) -> str:\n return f\"trained logistic regression on {n_rows} rows\"\n\n# Branching is now ordinary Python control flow -- no conditional() DSL.\n@env.task\ndef main(n_rows: int) -> str:\n if n_rows > 10_000:\n return train_gradient_boosting(n_rows)\n return train_logistic_regression(n_rows)\n```\n\n## Dynamic Workflows\n\n`@dynamic` existed so a task could generate a variable number of subtask calls at runtime (e.g. one per data partition discovered at runtime). In Flyte 2 every task can do this natively, so `@dynamic` simply disappears — loop over runtime data in an ordinary `@env.task`.\n\n### Flyte 1\n\n```python\nfrom flytekit import task, workflow, dynamic\n\n@task\ndef list_partitions(n: int) -> list[int]:\n return list(range(n))\n\n@task\ndef process_partition(partition_id: int) -> int:\n # Aggregate one data partition.\n return partition_id * 2\n\n@dynamic\ndef process_all(partitions: list[int]) -> list[int]:\n results = []\n for partition_id in partitions:\n results.append(process_partition(partition_id=partition_id))\n return results\n\n@workflow\ndef main(n: int) -> list[int]:\n partitions = list_partitions(n=n)\n return process_all(partitions=partitions)\n```\n\n### Flyte 2\n\n```python\nimport flyte\n\nenv = flyte.TaskEnvironment(name=\"dynamic\")\n\n@env.task\ndef process_partition(partition_id: int) -> int:\n # Aggregate one data partition.\n return partition_id * 2\n\n# No @dynamic decorator needed: a plain task can loop over runtime data (e.g. a\n# variable number of partitions discovered at runtime) and call other tasks.\n@env.task\ndef main(n: int) -> list[int]:\n return [process_partition(partition_id) for partition_id in range(n)]\n```\n\n## Error Handling\n\nFlyte 1's `@workflow(on_failure=...)` handler becomes ordinary Python `try` / `except` — catch a failed training run, run cleanup, and recover or re-raise.\n\n### Flyte 1\n\n```python\nfrom flytekit import task, workflow\n\n@task\ndef train_fold(max_depth: int) -> float:\n if max_depth <= 0:\n raise ValueError(\"max_depth must be positive\")\n # Return validation accuracy for this hyperparameter.\n return 0.90 + 0.001 * max_depth\n\n@task\ndef notify_failure() -> None:\n print(\"training run failed -- sending alert\")\n\n# The on_failure handler runs if any node in the workflow fails. There is no\n# try/except inside a Flyte 1 workflow.\n@workflow(on_failure=notify_failure)\ndef main(max_depth: int) -> float:\n return train_fold(max_depth=max_depth)\n```\n\n### Flyte 2\n\n```python\nimport flyte\n\nenv = flyte.TaskEnvironment(name=\"error_handling\")\n\n@env.task\nasync def train_fold(max_depth: int) -> float:\n if max_depth <= 0:\n raise ValueError(\"max_depth must be positive\")\n return 0.90 + 0.001 * max_depth\n\n# Failure handling is ordinary Python try/except -- no on_failure handler.\n@env.task\nasync def main(max_depth: int) -> float:\n try:\n return await train_fold(max_depth)\n except ValueError as e:\n print(f\"invalid hyperparameter ({e}); falling back to a safe default\")\n # Recover with a safe default instead of failing the whole run.\n return await train_fold(max_depth=6)\n```\n\nFlyte 2 also exposes typed errors, so you can catch a specific failure and retry with more resources — a common need for memory-hungry training jobs:\n\n```python\ntry:\n return await train_fold(sample_size)\nexcept flyte.errors.OOMError:\n # Retry the same task with a larger memory request.\n return await train_fold.override(\n resources=flyte.Resources(memory=\"16Gi\")\n )(sample_size)\n```\n\n## Fan-out: map_task\n\n`map_task()` becomes `flyte.map()`, a near drop-in replacement. The one catch: `flyte.map` returns a generator, so wrap it in `list()`. For new code, the idiomatic approach is Python `async`/`await` with `asyncio.gather()`, which gives finer control over concurrency and error handling.\n\n### Flyte 1\n\n```python\nfrom functools import partial\n\nfrom flytekit import task, workflow, map_task\n\n@task\ndef get_shards(n: int) -> list[int]:\n return list(range(n))\n\n@task\ndef score_shard(shard_id: int, model_version: int) -> int:\n # Score one shard of records with the given model version.\n return shard_id * model_version\n\n@workflow\ndef main(n: int, model_version: int) -> list[int]:\n shards = get_shards(n=n)\n return map_task(\n partial(score_shard, model_version=model_version),\n concurrency=10,\n )(shard_id=shards)\n```\n\n### Flyte 2 (flyte.map)\n\n```python\nimport flyte\nfrom functools import partial\n\nenv = flyte.TaskEnvironment(name=\"map_task\")\n\n@env.task\ndef score_shard(shard_id: int, model_version: int) -> int:\n # Score one shard of records with the given model version.\n return shard_id * model_version\n\n@env.task\ndef main(n: int, model_version: int) -> list[int]:\n bound = partial(score_shard, model_version=model_version)\n # flyte.map is a drop-in for map_task, but it returns a generator, so wrap\n # it in list() to materialize the results.\n return list(flyte.map(bound, range(n), concurrency=10))\n```\n\n### Flyte 2 (asyncio.gather)\n\n```python\nimport asyncio\n\nimport flyte\n\nenv = flyte.TaskEnvironment(name=\"map_task\")\n\n@env.task\nasync def score_shard_async(shard_id: int, model_version: int) -> int:\n return shard_id * model_version\n\n@env.task\nasync def main_async(n: int, model_version: int) -> list[int]:\n # asyncio.gather is the idiomatic Flyte 2 way to fan out.\n coros = [score_shard_async(i, model_version) for i in range(n)]\n return list(await asyncio.gather(*coros))\n```\n\n### Choosing flyte.map vs asyncio.gather\n\n| Feature | `flyte.map` (sync) | `asyncio.gather` (async) |\n|---|---|---|\n| Syntax | `list(flyte.map(fn, items))` | `await asyncio.gather(*tasks)` |\n| Concurrency limit | Built-in `concurrency=N` | Use `asyncio.Semaphore` |\n| Streaming / as-completed | No | Yes, via `asyncio.as_completed()` |\n| Error handling | `return_exceptions=True` | Check return type |\n\nUse `flyte.map` for the smallest change from Flyte 1 `map_task`, or when stuck in synchronous code. Use `asyncio.gather` for new code where you want streaming results or fine-grained concurrency control.\n\n### Concurrency Control and Error Handling\n\n`map_task`'s `concurrency` and `min_success_ratio` become an `asyncio.Semaphore` and `return_exceptions=True`:\n\n```python\nimport asyncio\n\n@env.task\nasync def main(items: list[int], max_concurrent: int = 5) -> list[str]:\n sem = asyncio.Semaphore(max_concurrent)\n\n async def process_with_limit(item: int) -> str:\n async with sem:\n return await process_item(item)\n\n tasks = [process_with_limit(i) for i in items]\n results = await asyncio.gather(*tasks, return_exceptions=True)\n\n return [r for r in results if not isinstance(r, Exception)]\n```\n\n## Data Backfills\n\nReprocessing a range of dates is a textbook `@dynamic` use case in Flyte 1, because the number of days is only known at runtime. In Flyte 2 it's a plain task that builds the date range and fans the days out with `asyncio.gather`.\n\n### Flyte 1\n\n```python\nfrom datetime import date, timedelta\n\nfrom flytekit import task, workflow, dynamic\n\n@task\ndef process_day(day: str) -> int:\n # Reprocess a single day's partition; return the row count.\n return len(day)\n\n# @dynamic is needed because the number of days is only known at runtime.\n@dynamic\ndef backfill(start: str, days: int) -> list[int]:\n base = date.fromisoformat(start)\n results = []\n for i in range(days):\n day = (base + timedelta(days=i)).isoformat()\n results.append(process_day(day=day))\n return results\n\n@workflow\ndef main(start: str, days: int) -> list[int]:\n return backfill(start=start, days=days)\n```\n\n### Flyte 2\n\n```python\nimport asyncio\nfrom datetime import date, timedelta\n\nimport flyte\n\nenv = flyte.TaskEnvironment(name=\"data_backfill\")\n\n@env.task\nasync def process_day(day: str) -> int:\n # Reprocess a single day's partition; return the row count.\n return len(day)\n\n# A plain task builds the date range at runtime and fans the days out in\n# parallel with asyncio.gather -- no @dynamic and no map_task needed.\n@env.task\nasync def main(start: str, days: int) -> list[int]:\n base = date.fromisoformat(start)\n coros = [\n process_day((base + timedelta(days=i)).isoformat())\n for i in range(days)\n ]\n return list(await asyncio.gather(*coros))\n```\n\n## Anti-Patterns\n\n1. **Don't import `conditional`, `dynamic`, or `map_task` from `flytekit`** — none exist in Flyte 2. Branching is native `if`/`elif`/`else`, dynamic fan-out is a plain task loop, and `map_task` becomes `flyte.map`.\n2. **Don't keep the `conditional().if_().then().else_()` DSL** — rewrite it as ordinary Python control flow inside an `@env.task`.\n3. **Don't reach for `@dynamic`** — every Flyte 2 task can loop over runtime data and call other tasks, so drop the decorator entirely.\n4. **Don't pass `on_failure=...` to `@workflow`** — there is no workflow decorator in Flyte 2; handle failures with ordinary `try`/`except` inside a task.\n5. **Don't forget to `list()` a `flyte.map` result** — it returns a generator, not a materialized list.\n6. **Don't forget to `await` async fan-out** — `asyncio.gather(*coros)` returns a coroutine; without `await` you get a coroutine object instead of results.\n7. **Don't drop concurrency limits** — port `concurrency=N` to `flyte.map(..., concurrency=N)` or an `asyncio.Semaphore`, and `min_success_ratio` to `return_exceptions=True` with filtering.\n8. **Don't use Union-only features** — avoid `ReusePolicy` and other Union-specific APIs.\n"
}SHA-256 of public snapshot: c0f0c5da912314bf6a8d465647ebe77a2d10d7dccc54c62de8e524b4c65a9ff3