← 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": "Guides correct types, I/O, and serialization for common data (Pandas, Arrow, Parquet, images, audio, HF datasets), including data locality and storage best practices. Use when the user needs help with type annotations, data serialization, file I/O between tasks, custom type transformers, DataFrame handling, or choosing the right Flyte type for their data. Trigger words: \"type\", \"serialize\", \"deserialize\", \"DataFrame\", \"File\", \"Directory\", \"custom type\", \"data format\", \"Parquet\", \"Arrow\", \"Pandas\", \"type transformer\".",
"included_files": [],
"name": "flyte-sdk-types",
"skill_md_contents": "---\nname: flyte-sdk-types\ndescription: 'Guides correct types, I/O, and serialization for common data (Pandas, Arrow, Parquet, images, audio, HF datasets), including data locality and storage best practices. Use when the user needs help with type annotations, data serialization, file I/O between tasks, custom type transformers, DataFrame handling, or choosing the right Flyte type for their data. Trigger words: \"type\", \"serialize\", \"deserialize\", \"DataFrame\", \"File\", \"Directory\", \"custom type\", \"data format\", \"Parquet\", \"Arrow\", \"Pandas\", \"type transformer\".'\n---\n\n# Flyte 2 SDK Types Skill\n\nGuide correct type annotations, I/O patterns, and serialization for Flyte 2 workflows.\n\n## Grounding References\n\n| Resource | URL |\n|---|---|\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| CLI API reference | https://www.union.ai/docs/v2/union/api-reference/flyte-cli/ |\n| flyte-sdk source | https://github.com/flyteorg/flyte-sdk |\n| Example code | https://github.com/unionai/unionai-examples |\n| Flyte MCP tools | Available via the `flyte-cluster` and `flyte-docs` MCP servers |\n\n**Ground unfamiliar APIs in real examples.** When unsure of a current Flyte 2 API, or for a pattern not shown below, and the `flyte-docs` search tools are available, search them first — by exact symbol (`TaskEnvironment`, `flyte.io.File`, `map_task`), since matching is literal substring, not semantic — then adapt a real example rather than inventing one, and cite the file or section you pulled it from. (Flyte 2 is not `flytekit`; priors are often wrong.)\n\n## Type System Overview\n\nFlyte 2 uses **Python type hints** for serialization. Every task input/output must have a type annotation. Flyte's type transformer system handles conversion between Python types and remote storage (S3, GCS, etc.).\n\n### Supported Native Types\n\n| Python Type | Flyte Type | Remote Transport |\n|---|---|---|\n| `int`, `float`, `bool`, `str` | Literal | Inline (JSON) |\n| `list[T]`, `dict[K, V]` | Collection / Map | Inline (JSON) for small, blob for large |\n| `flyte.io.File` | Blob | Uploaded to metadata bucket |\n| `flyte.io.Dir` | Blob (directory) | Uploaded to metadata bucket |\n| `flyte.io.DataFrame` | DataFrame | Parquet in metadata bucket |\n| `dataclass` | Struct | JSON in metadata bucket |\n| `pydantic.BaseModel` | Struct | JSON in metadata bucket |\n| `datetime`, `timedelta` | DateTime / Duration | Inline |\n\n## flyte.io.File — Single Files\n\nUse `flyte.io.File` for any single file that flows between tasks. Files are **automatically uploaded** to the metadata bucket at runtime.\n\n```python\nimport flyte\nimport flyte.io\n\n@env.task\nasync def download_url(url: str) -> flyte.io.File:\n \"\"\"Download a file and return as flyte.io.File.\"\"\"\n import urllib.request\n local_path = f\"/tmp/{url.split('/')[-1]}\"\n urllib.request.urlretrieve(url, local_path)\n return flyte.io.File(path=local_path)\n\n@env.task\nasync def process(file: flyte.io.File) -> dict:\n \"\"\"Read a file — flyte.io.File downloads it automatically.\"\"\"\n # file.path gives the local path (already downloaded)\n with open(file.path) as f:\n content = f.read()\n return {\"lines\": len(content.splitlines())}\n\n@env.task\nasync def main(url: str) -> dict:\n downloaded = await download_url(url)\n return await process(downloaded) # type: flyte.io.File flows as remote reference\n```\n\n### flyte.io.File best practices\n\n- **Always pass `flyte.io.File` between tasks** — never pass file paths as strings. Flyte serializes the reference to the remote blob.\n- **`file.path`** gives the local download path inside the task container.\n- **Don't hardcode paths** — let Flyte manage the download/upload lifecycle.\n- **Compression**: Flyte infers format from extension (`.parquet`, `.csv`, `.json`, `.pt`, `.png`, etc.).\n\n## flyte.io.Dir — Directories\n\nUse `flyte.io.Dir` for a collection of files (e.g., model checkpoints, output artifacts).\n\n```python\nimport flyte\nimport flyte.io\n\n@env.task\nasync def train(checkpoint_dir: flyte.io.Dir) -> flyte.io.Dir:\n \"\"\"Train and save checkpoints to a directory.\"\"\"\n # Write checkpoints\n for i in range(10):\n path = f\"{checkpoint_dir.path}/checkpoint_{i}.pt\"\n save_model(path)\n return checkpoint_dir\n\n@env.task\nasync def evaluate(checkpoints: flyte.io.Dir) -> dict:\n \"\"\"Load checkpoints from a directory.\"\"\"\n # List all files in the directory\n files = list(checkpoints.path.glob(\"*.pt\"))\n ...\n```\n\n## flyte.io.DataFrame — Polars DataFrames\n\nFlyte 2 has built-in support for Polars DataFrames. They are **passed by reference** (Parquet in the metadata bucket), not inline.\n\n```python\nimport flyte\nimport flyte.io\n\n@env.task\nasync def load_csv(url: str) -> flyte.io.DataFrame:\n \"\"\"Load a CSV and return as Polars DataFrame.\"\"\"\n import polars as pl\n df = pl.read_csv(url)\n return flyte.io.DataFrame(df)\n\n@env.task\nasync def clean(df: flyte.io.DataFrame) -> flyte.io.DataFrame:\n \"\"\"Clean the DataFrame.\"\"\"\n inner = df.to_polars() # Get the underlying Polars DataFrame\n cleaned = inner.drop_nulls()\n return flyte.io.DataFrame(cleaned)\n\n@env.task\nasync def save_parquet(df: flyte.io.DataFrame, path: str) -> flyte.io.File:\n \"\"\"Save DataFrame to Parquet.\"\"\"\n inner = df.to_polars()\n inner.write_parquet(path)\n return flyte.io.File(path=path)\n\n@env.task\nasync def main(url: str) -> flyte.io.DataFrame:\n raw = await load_csv(url)\n cleaned = await clean(raw)\n return cleaned # Flows as Parquet reference\n```\n\n### Polars DataFrame patterns\n\n```python\n# Convert to Polars for manipulation\ninner_df = df.to_polars()\n\n# Convert from Polars\ndf = flyte.io.DataFrame(inner_df)\n\n# Common operations\ndf = flyte.io.DataFrame(inner_df.filter(pl.col(\"age\") > 18))\ndf = flyte.io.DataFrame(inner_df.group_by(\"category\").agg(pl.col(\"value\").mean()))\n\n# Check shape\nprint(df.shape) # (rows, cols)\nprint(df.schema) # column names and types\n```\n\n### Eager vs Lazy DataFrames\n\n```python\n# Eager (loaded into memory)\ndf = flyte.io.DataFrame(inner_df)\n\n# Lazy (streaming, for large datasets)\ndf = flyte.io.DataFrame(inner_df.lazy())\n\n# Materialize lazy to eager\neager = df.to_polars() # materializes\n```\n\n## Dataclass and Pydantic Models\n\nFor structured data, use Python dataclasses or Pydantic models. Flyte serializes them to JSON.\n\n```python\nfrom dataclasses import dataclass\nfrom pydantic import BaseModel\nimport flyte\nimport flyte.io\n\n@dataclass\nclass TrainingConfig:\n learning_rate: float\n batch_size: int\n epochs: int\n\nclass PredictionOutput(BaseModel):\n predictions: list[float]\n confidence: list[float]\n model_version: str\n\n@env.task\nasync def train(config: TrainingConfig) -> flyte.io.File:\n # config.learning_rate, config.batch_size, etc.\n ...\n\n@env.task\nasync def predict(model: flyte.io.File, data: flyte.io.DataFrame) -> PredictionOutput:\n return PredictionOutput(\n predictions=[0.5, 0.8, 0.3],\n confidence=[0.9, 0.7, 0.95],\n model_version=\"v1.0\",\n )\n```\n\n## Custom Type Transformers\n\nExtend Flyte's type system to support custom types (e.g., PIL Images, HuggingFace datasets).\n\n### PIL Image transformer\n\n```python\nfrom PIL import Image\nimport flyte\nimport flyte.io\nfrom flyte.types import TypeTransformer\n\nclass PILImageTransformer(TypeTransformer[Image.Image]):\n _type = Image.Image\n\n def get_type(self, input: Image.Image) -> type:\n return Image.Image\n\n def save(self, img: Image.Image, path: str) -> None:\n img.save(path)\n\n def load(self, path: str) -> Image.Image:\n return Image.open(path)\n\n# Register the transformer\nflyte.types.TypeEngine.register(PILImageTransformer())\n\n# Now use it in tasks\n@env.task\nasync def process_image(img: Image.Image) -> flyte.io.File:\n # img is a PIL Image, already downloaded\n ...\n```\n\n### HuggingFace Dataset transformer\n\n```python\nfrom datasets import Dataset\nimport flyte\n\nclass HFDatasetTransformer(TypeTransformer[Dataset]):\n _type = Dataset\n\n def get_type(self, input: Dataset) -> type:\n return Dataset\n\n def save(self, ds: Dataset, path: str) -> None:\n ds.save_to_disk(path)\n\n def load(self, path: str) -> Dataset:\n return Dataset.load_from_disk(path)\n\nflyte.types.TypeEngine.register(HFDatasetTransformer())\n```\n\n## Data I/O Patterns by Domain\n\n### ETL / Data Engineering\n\n```python\n# CSV → Parquet conversion\n@env.task(cache=\"auto\")\nasync def csv_to_parquet(csv_file: flyte.io.File, output_path: str) -> flyte.io.File:\n import polars as pl\n df = pl.read_csv(csv_file.path)\n df.write_parquet(output_path)\n return flyte.io.File(path=output_path)\n\n# JsonlFile — batched JSONL reading\n@env.task\nasync def process_jsonl(path: str) -> int:\n from flyte.extend import JsonlFile\n jf = JsonlFile(path)\n count = 0\n async for record in jf.stream():\n process(record)\n count += 1\n return count\n\n# JsonlDir — batched JSONL directory\n@env.task\nasync def process_jsonl_dir(dir_path: str) -> dict:\n from flyte.extend import JsonlDir\n jd = JsonlDir(dir_path)\n total = 0\n async for batch in jd.stream_batches():\n total += len(batch)\n return {\"records\": total}\n```\n\n### Image Processing\n\n```python\nfrom PIL import Image\nimport flyte\nimport flyte.io\n\n@env.task\nasync def resize_image(input_file: flyte.io.File, size: tuple[int, int]) -> flyte.io.File:\n img = Image.open(input_file.path)\n resized = img.resize(size)\n output_path = f\"/tmp/resized_{size[0]}x{size[1]}.png\"\n resized.save(output_path)\n return flyte.io.File(path=output_path)\n\n@env.task\nasync def batch_resize(files: list[flyte.io.File], size: tuple[int, int]) -> list:\n import asyncio\n return await asyncio.gather(*(resize_image(f, size) for f in files))\n```\n\n### Audio Processing\n\n```python\nimport librosa\nimport numpy as np\nimport flyte\nimport flyte.io\n\n@env.task\nasync def extract_features(audio_file: flyte.io.File) -> flyte.io.File:\n \"\"\"Extract MFCC features from audio file.\"\"\"\n y, sr = librosa.load(audio_file.path, sr=None)\n mfccs = librosa.feature.mfcc(y=y, sr=sr)\n # Save as numpy array\n output_path = \"/tmp/mfccs.npz\"\n np.savez(output_path, mfccs=mfccs, sr=sr)\n return flyte.io.File(path=output_path)\n```\n\n### HuggingFace Datasets\n\n```python\nfrom datasets import load_dataset, Dataset\nimport flyte\nimport flyte.io\n\n@env_task\nasync def load_dataset_from_hub(dataset_name: str, split: str = \"train\") -> flyte.io.File:\n \"\"\"Load a HF dataset and save locally.\"\"\"\n ds = load_dataset(dataset_name, split=split)\n path = f\"/tmp/{dataset_name.replace('/', '_')}_{split}\"\n ds.save_to_disk(path)\n return flyte.io.File(path=path)\n\n@env.task\nasync def process_dataset(file: flyte.io.File) -> flyte.io.DataFrame:\n \"\"\"Convert HF dataset to Flyte DataFrame.\"\"\"\n ds = Dataset.load_from_disk(file.path)\n df = ds.to_pandas()\n return flyte.io.DataFrame(df)\n```\n\n## Data Locality and Storage Best Practices\n\n### How data flows between tasks\n\n1. **By reference (default)** — large data (DataFrames, Files, Directories) is uploaded to the metadata bucket. Tasks receive a remote reference and download on demand.\n2. **Inline (small data)** — primitives (int, float, str, bool) and small collections are passed inline as JSON.\n\n### Choosing the right transport\n\n| Data type | Transport | Max size |\n|---|---|---|\n| `int`, `float`, `bool`, `str` | Inline | None |\n| `list`, `dict` (small) | Inline | ~10 MB |\n| `flyte.io.File` | Reference (S3/GCS) | Unlimited |\n| `flyte.io.Dir` | Reference (S3/GCS) | Unlimited |\n| `flyte.io.DataFrame` (Polars) | Reference (Parquet) | Unlimited |\n| `dataclass` / `BaseModel` | Inline (JSON) | ~10 MB |\n\n### Storage best practices\n\n1. **Use `flyte.io.File` for files** — don't pass strings. Flyte manages the upload/download.\n2. **Use `flyte.io.DataFrame` for tabular data** — stored as Parquet, efficient for ETL pipelines.\n3. **Use `flyte.io.Dir` for collections** — model checkpoints, output artifacts.\n4. **Set `cache=\"auto\"`** on idempotent tasks (ETL, transforms) to avoid re-processing.\n5. **Use `raw_data_path`** for per-run customization: `flyte.with_runcontext(raw_data_path=\"s3://my-bucket/{run_id}/\")`\n6. **Large data should never be inline** — if your dataclass exceeds ~10 MB, switch to `flyte.io.File` or `flyte.io.Dir`.\n\n## Inline I/O Threshold\n\nControl when data is passed inline vs by reference:\n\n```python\n@env.task(inline_output_limit=\"10MB\") # data > 10MB goes by reference\nasync def process(data: dict) -> dict:\n ...\n```\n\nDefault threshold is generous. For ML training outputs or large DataFrames, set it lower to avoid overhead.\n\n## Common Type Mistakes\n\n1. **Missing type hints** — Flyte 2 requires type annotations on all task inputs/outputs. No type hint = serialization error.\n2. **Passing `flyte.io.File` as a string** — always use `flyte.io.File(path=...)` objects. The path string alone won't serialize.\n3. **Using Pandas instead of Polars** — Flyte 2's native DataFrame is Polars. Use `df.to_polars()` to get the underlying DataFrame.\n4. **Not registering custom transformers** — if you register a custom type transformer, it must be registered before task execution.\n5. **Forgetting `.path` on flyte.io.File** — inside a task, `file` is a FlyteFile object, not a string. Use `file.path` for the local path.\n"
}SHA-256 of public snapshot: 4e87c168f5cb30bfe6da0d3ae46fc6dc25c63b5497ef19fb5daf216f5fb78731