← 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": "Migrates Flyte 1 data types and offloaded I/O to Flyte 2. Use when migrating Flyte 1 data types and I/O (files, directories, dataframes, dataclasses) to Flyte 2, converting FlyteFile, FlyteDirectory, or StructuredDataset to flyte.io.File, flyte.io.Dir, and flyte.io.DataFrame. Trigger words are FlyteFile, FlyteDirectory, StructuredDataset, DataFrame, dataclass, Pydantic, type, I/O, and serialization.",
  "included_files": [],
  "name": "flyte-migrate-data-io",
  "skill_md_contents": "---\nname: flyte-migrate-data-io\ndescription: Migrates Flyte 1 data types and offloaded I/O to Flyte 2. Use when migrating Flyte 1 data types and I/O (files, directories, dataframes, dataclasses) to Flyte 2, converting FlyteFile, FlyteDirectory, or StructuredDataset to flyte.io.File, flyte.io.Dir, and flyte.io.DataFrame. Trigger words are FlyteFile, FlyteDirectory, StructuredDataset, DataFrame, dataclass, Pydantic, type, I/O, and serialization.\n---\n\n# Flyte 1 to 2 Migration: Data Types and I/O\n\nFlyte 2 renames the offloaded-data types and makes their I/O `async`, but the mental model is the same: pass lightweight references to large data between tasks, not the materialized bytes. `FlyteFile`, `FlyteDirectory`, and `StructuredDataset` become `flyte.io.File`, `flyte.io.Dir`, and `flyte.io.DataFrame`. Plain dataclasses and Pydantic `BaseModel`s work directly as task I/O with no JSON mixin.\n\n## Grounding References\n\n| Resource | URL |\n|---|---|\n| Migration guide | https://www.union.ai/docs/v2/flyte/user-guide/migration/flyte-2/data-io/ |\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## Type Mapping\n\n| Flyte 1 | Flyte 2 | Notes |\n|---|---|---|\n| `flytekit.types.file.FlyteFile` | `flyte.io.File` | I/O is `async` |\n| `flytekit.types.directory.FlyteDirectory` | `flyte.io.Dir` | I/O is `async` |\n| `flytekit.types.structured.StructuredDataset` | `flyte.io.DataFrame` | build with `from_df`, read with `open(...).all()` |\n| `@dataclass_json` + `@dataclass` | plain `@dataclass` | no mixin needed |\n| Pydantic `BaseModel` (+ config) | plain Pydantic `BaseModel` | works directly as task I/O |\n\n## Offloaded Data: The Mental Model\n\n`File`, `Dir`, and `DataFrame` are lightweight references (pointers) to data offloaded in blob storage — not the materialized bytes. In Flyte 2 the read/write operations are `async`: upload with `await File.from_local(local_path)`, read with `async with f.open(\"rb\") as fh: await fh.read()`, build a frame with `flyte.io.DataFrame.from_df(df)` (sync constructor), and read it with `await fdf.open(pandas.DataFrame).all()`.\n\n## Files and Directories\n\n`FlyteFile` and `FlyteDirectory` become `flyte.io.File` and `flyte.io.Dir` — the way you pass model artifacts and datasets between tasks. Use `await File.from_local(...)` to upload and `async with file.open(...)` to read.\n\n### Flyte 1\n\n```python\nimport os\n\nfrom flytekit import task, workflow, current_context\nfrom flytekit.types.file import FlyteFile\n\n@task\ndef write_file(content: str) -> FlyteFile:\n    path = os.path.join(current_context().working_directory, \"out.txt\")\n    with open(path, \"w\") as f:\n        f.write(content)\n    return FlyteFile(path=path)\n\n@task\ndef read_file(f: FlyteFile) -> str:\n    with open(f.download()) as fh:\n        return fh.read()\n\n@workflow\ndef main(content: str) -> str:\n    f = write_file(content=content)\n    return read_file(f=f)\n```\n\n### Flyte 2\n\n```python\nimport flyte\nfrom flyte.io import File\n\nenv = flyte.TaskEnvironment(name=\"files\")\n\n@env.task\nasync def write_file(content: str) -> File:\n    with open(\"out.txt\", \"w\") as f:\n        f.write(content)\n    # File.from_local uploads the file to blob storage and returns a reference\n    # (a lightweight pointer, not the materialized bytes).\n    return await File.from_local(\"out.txt\")\n\n@env.task\nasync def read_file(f: File) -> str:\n    async with f.open(\"rb\") as fh:\n        return (await fh.read()).decode(\"utf-8\")\n\n@env.task\nasync def main(content: str) -> str:\n    f = await write_file(content)\n    return await read_file(f)\n```\n\nDirectories follow the same pattern: import `Dir` from `flyte.io` and use its `async` upload/read methods in place of `FlyteDirectory`. See [Files and directories](https://www.union.ai/docs/v2/flyte/user-guide/task-programming/files-and-directories) for more.\n\n## DataFrames\n\n`StructuredDataset` becomes `flyte.io.DataFrame`. Construct one with `flyte.io.DataFrame.from_df(df)` and read it back with `await df.open(pandas.DataFrame).all()`.\n\n### Flyte 1\n\n```python\nimport pandas as pd\nfrom flytekit import task, workflow\nfrom flytekit.types.structured import StructuredDataset\n\n@task\ndef make_df() -> StructuredDataset:\n    df = pd.DataFrame({\"employee_id\": [1, 2, 3], \"salary\": [50000, 60000, 70000]})\n    return StructuredDataset(dataframe=df)\n\n@task\ndef total_payroll(sd: StructuredDataset) -> float:\n    df = sd.open(pd.DataFrame).all()\n    return float(df[\"salary\"].sum())\n\n@workflow\ndef main() -> float:\n    return total_payroll(sd=make_df())\n```\n\n### Flyte 2\n\n```python\nimport pandas as pd\nimport flyte\nimport flyte.io\n\nenv = flyte.TaskEnvironment(\n    name=\"dataframe\",\n    image=flyte.Image.from_debian_base().with_pip_packages(\"pandas\", \"pyarrow\"),\n)\n\n@env.task\nasync def make_df() -> flyte.io.DataFrame:\n    df = pd.DataFrame({\"employee_id\": [1, 2, 3], \"salary\": [50000, 60000, 70000]})\n    # StructuredDataset becomes flyte.io.DataFrame.\n    return flyte.io.DataFrame.from_df(df)\n\n@env.task\nasync def total_payroll(fdf: flyte.io.DataFrame) -> float:\n    df = await fdf.open(pd.DataFrame).all()\n    return float(df[\"salary\"].sum())\n\n@env.task\nasync def main() -> float:\n    return await total_payroll(await make_df())\n```\n\nAdd the dataframe dependencies (for example `pandas` and `pyarrow`) to the `TaskEnvironment` image. See [DataFrames](https://www.union.ai/docs/v2/flyte/user-guide/task-programming/dataframes) for more.\n\n## Dataclasses and Structured Types\n\nFlyte 1 required a `@dataclass_json` mixin for dataclass I/O. In Flyte 2, plain dataclasses (and Pydantic `BaseModel`s) work directly as task inputs and outputs — handy for passing around a training config.\n\n### Flyte 1\n\n```python\nfrom dataclasses import dataclass\n\nfrom dataclasses_json import dataclass_json\nfrom flytekit import task, workflow\n\n@dataclass_json\n@dataclass\nclass TrainingConfig:\n    learning_rate: float\n    n_estimators: int\n    max_depth: int = 6\n\n@task\ndef make_config(learning_rate: float, n_estimators: int) -> TrainingConfig:\n    return TrainingConfig(learning_rate=learning_rate, n_estimators=n_estimators)\n\n@task\ndef train(config: TrainingConfig) -> str:\n    return (\n        f\"trained with lr={config.learning_rate}, \"\n        f\"n_estimators={config.n_estimators}, max_depth={config.max_depth}\"\n    )\n\n@workflow\ndef main(learning_rate: float, n_estimators: int) -> str:\n    config = make_config(learning_rate=learning_rate, n_estimators=n_estimators)\n    return train(config=config)\n```\n\n### Flyte 2\n\n```python\nfrom dataclasses import dataclass\n\nimport flyte\n\nenv = flyte.TaskEnvironment(name=\"dataclasses\")\n\n# Plain dataclasses work directly as task I/O -- no @dataclass_json mixin needed.\n# Pydantic BaseModels work the same way.\n@dataclass\nclass TrainingConfig:\n    learning_rate: float\n    n_estimators: int\n    max_depth: int = 6\n\n@env.task\ndef make_config(learning_rate: float, n_estimators: int) -> TrainingConfig:\n    return TrainingConfig(learning_rate=learning_rate, n_estimators=n_estimators)\n\n@env.task\ndef train(config: TrainingConfig) -> str:\n    return (\n        f\"trained with lr={config.learning_rate}, \"\n        f\"n_estimators={config.n_estimators}, max_depth={config.max_depth}\"\n    )\n\n@env.task\ndef main(learning_rate: float, n_estimators: int) -> str:\n    config = make_config(learning_rate, n_estimators)\n    return train(config)\n```\n\n## Data ETL: Putting It Together\n\nExtract, clean, aggregate, and write out a feature table. `StructuredDataset` becomes `flyte.io.DataFrame`, and the tasks become `async`.\n\n### Flyte 1\n\n```python\nimport pandas as pd\nfrom flytekit import task, workflow\nfrom flytekit.types.structured import StructuredDataset\n\n@task\ndef extract() -> pd.DataFrame:\n    # Read raw transaction records (stand-in for a real source).\n    return pd.DataFrame(\n        {\n            \"user_id\": [1, 1, 2, 3, 3, 3],\n            \"amount\": [10.0, 5.0, 20.0, 7.5, 2.5, 1.0],\n        }\n    )\n\n@task\ndef transform(df: pd.DataFrame) -> StructuredDataset:\n    # Clean and aggregate into a per-user feature table.\n    df = df[df[\"amount\"] > 0]\n    agg = df.groupby(\"user_id\", as_index=False)[\"amount\"].sum()\n    return StructuredDataset(dataframe=agg)\n\n@task\ndef load(sd: StructuredDataset) -> int:\n    df = sd.open(pd.DataFrame).all()\n    return len(df)\n\n@workflow\ndef main() -> int:\n    raw = extract()\n    features = transform(df=raw)\n    return load(sd=features)\n```\n\n### Flyte 2\n\n```python\nimport pandas as pd\nimport flyte\nimport flyte.io\n\nenv = flyte.TaskEnvironment(\n    name=\"data_etl\",\n    image=flyte.Image.from_debian_base().with_pip_packages(\"pandas\", \"pyarrow\"),\n)\n\n@env.task\nasync def extract() -> pd.DataFrame:\n    # Read raw transaction records (stand-in for a real source).\n    return pd.DataFrame(\n        {\n            \"user_id\": [1, 1, 2, 3, 3, 3],\n            \"amount\": [10.0, 5.0, 20.0, 7.5, 2.5, 1.0],\n        }\n    )\n\n@env.task\nasync def transform(df: pd.DataFrame) -> flyte.io.DataFrame:\n    # Clean and aggregate into a per-user feature table.\n    df = df[df[\"amount\"] > 0]\n    agg = df.groupby(\"user_id\", as_index=False)[\"amount\"].sum()\n    # StructuredDataset becomes flyte.io.DataFrame.\n    return flyte.io.DataFrame.from_df(agg)\n\n@env.task\nasync def load(sd: flyte.io.DataFrame) -> int:\n    df = await sd.open(pd.DataFrame).all()\n    return len(df)\n\n@env.task\nasync def main() -> int:\n    raw = await extract()\n    features = await transform(raw)\n    return await load(features)\n```\n\n## Anti-Patterns\n\n1. **Don't call the offloaded-data I/O synchronously** — `File.from_local`, `file.open(...).read()`, and `DataFrame.open(...).all()` are `async` in Flyte 2; `await` them inside `async` tasks.\n2. **Don't keep the `@dataclass_json` mixin** — plain `@dataclass` and Pydantic `BaseModel`s serialize as task I/O directly; drop `dataclasses_json`.\n3. **Don't return `StructuredDataset(dataframe=df)`** — use `flyte.io.DataFrame.from_df(df)` instead.\n4. **Don't materialize large data into task outputs** — return `File`, `Dir`, or `DataFrame` references, not the raw bytes or full frames.\n5. **Don't forget the dataframe dependencies** — add `pandas` and `pyarrow` (or your engine) to the `TaskEnvironment` image so DataFrame I/O works remotely.\n6. **Don't import from `flytekit.types.*`** — import `File` and `Dir` from `flyte.io`, and use `flyte.io.DataFrame`.\n"
}

SHA-256 of public snapshot: a8b5d9016934f294e6cf4726d4bf903590585964ca4bd16979a422ceff38c9a2