← 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": "Suggests performance improvements (task granularity, caching, resource requests, data format changes) using observed run metadata when available. Use when the user wants to optimize workflow performance, debug slow tasks, configure caching, tune resources, or improve throughput. Trigger words: \"optimize\", \"performance\", \"slow\", \"cache\", \"caching\", \"resource\", \"throughput\", \"latency\", \"speed up\", \"bottleneck\", \"profiling\", \"metadata\".",
  "included_files": [],
  "name": "flyte-sdk-optimize",
  "skill_md_contents": "---\nname: flyte-sdk-optimize\ndescription: 'Suggests performance improvements (task granularity, caching, resource requests, data format changes) using observed run metadata when available. Use when the user wants to optimize workflow performance, debug slow tasks, configure caching, tune resources, or improve throughput. Trigger words: \"optimize\", \"performance\", \"slow\", \"cache\", \"caching\", \"resource\", \"throughput\", \"latency\", \"speed up\", \"bottleneck\", \"profiling\", \"metadata\".'\n---\n\n# Flyte 2 SDK Optimize Skill\n\nOptimize Flyte 2 workflows for performance, cost, and reliability.\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## Optimization Strategy Overview\n\nPerformance optimization in Flyte follows a hierarchy:\n\n1. **Reduce container overhead** — use traces for lightweight ops\n2. **Parallelize work** — use `flyte.map` for fan-out\n3. **Cache results** — use `cache=\"auto\"` for idempotent tasks\n4. **Tune resources** — set appropriate CPU/memory/GPU\n5. **Optimize data transfer** — choose efficient formats, reduce inline I/O\n6. **Use reusable containers** — shared environments reduce image pull time\n\n## Caching\n\n### Enable automatic caching\n\n```python\n@env.task(cache=\"auto\")  # versioned by function body + inputs\nasync def preprocess(data: list[str]) -> flyte.io.File:\n    ...\n```\n\n### Cache key strategies\n\n```python\n@env.task(cache=\"auto\")  # default: function body + inputs\nasync def task_a(data: str) -> flyte.io.File:\n    ...\n\n@env.task(cache=\"override\", salt=\"v2\")  # add salt for cache key variation\nasync def task_b(data: str) -> flyte.io.File:\n    ...\n\n@env.task(cache=\"disable\")  # always re-run\nasync def task_c(data: str) -> flyte.io.File:\n    ...\n```\n\n### Content-based caching for DataFrames\n\n```python\n@env.task(cache=\"auto\")\nasync def transform(df: flyte.io.DataFrame) -> flyte.io.DataFrame:\n    \"\"\"Cache key includes DataFrame content hash.\"\"\"\n    ...\n```\n\n### Ignoring specific inputs in cache key\n\n```python\n@env.task(cache=\"auto\", cache_ignore_inputs=[\"api_key\"])\nasync def fetch_data(api_key: str, url: str) -> flyte.io.File:\n    \"\"\"Don't include api_key in cache key.\"\"\"\n    ...\n```\n\n### Cache policies\n\n```python\n@env.task(cache=\"auto\", cache_policy=flyte.CachePolicy(min_cached_age=\"1h\"))\nasync def cached_task(data: str) -> flyte.io.File:\n    \"\"\"Only use cache if result is at least 1 hour old.\"\"\"\n    ...\n```\n\n## Resource Tuning\n\n### Setting task resources\n\n```python\n@env.task(\n    requests=flyte.Resources(cpu=\"500m\", memory=\"1Gi\"),\n    limits=flyte.Resources(cpu=\"2\", memory=\"4Gi\"),\n)\nasync def light_task(data: str) -> str:\n    \"\"\"Lightweight task — small resources.\"\"\"\n    ...\n\n@env.task(\n    requests=flyte.Resources(cpu=\"4\", memory=\"16Gi\"),\n    limits=flyte.Resources(cpu=\"8\", memory=\"32Gi\"),\n)\nasync def heavy_task(data: flyte.io.DataFrame) -> flyte.io.DataFrame:\n    \"\"\"Heavy data processing — large resources.\"\"\"\n    ...\n\n@env.task(\n    requests=flyte.Resources(cpu=\"1\", memory=\"4Gi\", gpu=\"1\", gpu_model=\"nvidia-a10g\"),\n    limits=flyte.Resources(cpu=\"2\", memory=\"8Gi\", gpu=\"1\", gpu_model=\"nvidia-a10g\"),\n)\nasync def train_model(data: flyte.io.File) -> flyte.io.File:\n    \"\"\"GPU training task.\"\"\"\n    ...\n```\n\n### GPU resource configuration\n\n```python\n@env.task(\n    requests=flyte.Resources(\n        cpu=\"2\",\n        memory=\"8Gi\",\n        gpu=\"1\",\n        gpu_model=\"nvidia-a10g\",  # or \"nvidia-a100\", \"nvidia-h100\"\n    ),\n)\nasync def inference(batch: flyte.io.DataFrame) -> flyte.io.DataFrame:\n    ...\n```\n\n### Resource recommendations by workload\n\n| Workload | CPU | Memory | GPU |\n|---|---|---|---|\n| Light ETL | 500m-1 | 1-2 Gi | none |\n| Data processing | 2-4 | 8-16 Gi | none |\n| Embedding | 2-4 | 8-16 Gi | none |\n| Model training | 4-8 | 16-32 Gi | 1-8 |\n| Batch inference | 2-4 | 8-16 Gi | 1-4 |\n| LLM serving | 8-16 | 32-64 Gi | 1-8 |\n| Data quality | 1-2 | 4-8 Gi | none |\n\n## Parallelization Patterns\n\n### flyte.map for parallel execution\n\n```python\n@env.task\nasync def process_item(item: dict) -> dict:\n    \"\"\"Process a single item.\"\"\"\n    ...\n\n@env.task\nasync def main(items: list[dict]) -> list:\n    \"\"\"Fan out processing in parallel.\"\"\"\n    results = await flyte.map(process_item, items)\n    return results\n```\n\n### flyte.trace for lightweight parallelism\n\n```python\n@env.task\nasync def fetch_url(url: str) -> str:\n    \"\"\"Lightweight HTTP fetch — use trace (no container overhead).\"\"\"\n    ...\n\n@env.task\nasync def main(urls: list[str]) -> list:\n    \"\"\"Use trace for light ops (no container spin-up cost).\"\"\"\n    results = await flyte.trace(fetch_url, urls)\n    return results\n```\n\n### asyncio.gather for sequential fan-out\n\n```python\n@env.task\nasync def main(data: list[str]) -> dict:\n    \"\"\"Chain tasks with parallel fan-out at each step.\"\"\"\n    # Step 1: parallel preprocessing\n    preprocessed = await asyncio.gather(*(preprocess(d) for d in data))\n\n    # Step 2: sequential aggregation\n    aggregated = aggregate(preprocessed)\n\n    # Step 3: parallel evaluation\n    metrics = await asyncio.gather(*(evaluate(p) for p in preprocessed))\n\n    return {\"aggregated\": aggregated, \"metrics\": metrics}\n```\n\n### Controlling concurrency\n\n```python\n@env.task\nasync def main(urls: list[str]) -> list:\n    \"\"\"Limit concurrency with asyncio.Semaphore.\"\"\"\n    import asyncio\n    sem = asyncio.Semaphore(10)  # max 10 concurrent\n\n    async def bounded(item):\n        async with sem:\n            return await process_item(item)\n\n    return await asyncio.gather(*(bounded(u) for u in urls))\n```\n\n## Data Format Optimization\n\n### Choosing efficient formats\n\n| Use case | Recommended format | Why |\n|---|---|---|\n| Tabular data | Parquet | Columnar, compressed, fast |\n| JSON data | JSONL | Line-delimited, streaming |\n| Images | PNG/WebP | Lossless/lossy compression |\n| Audio | WAV/FLAC | Lossless |\n| Model checkpoints | .pt/.safetensors | Native framework format |\n| Embeddings | .npy/.npz | NumPy binary format |\n\n### Reducing inline I/O\n\n```python\n# Bad: large dict passed inline (JSON serialization overhead)\n@env.task\nasync def process(large_data: dict) -> dict:\n    ...\n\n# Good: pass by reference\n@env.task\nasync def process(data_file: flyte.io.File) -> flyte.io.File:\n    ...\n\n# Good: set inline output limit\n@env.task(inline_output_limit=\"5MB\")\nasync def process(data: dict) -> dict:\n    ...\n```\n\n## Run Metadata Inspection\n\n### Using MCP to inspect runs\n\nIf Flyte MCP tools are available, use them to list a task's recent runs for performance\nanalysis, fetch a run's metadata (status, duration), and read its inputs and outputs.\n\n\n### Performance analysis checklist\n\n1. **Check run duration** — fetch the run's metadata and read `durationMs`\n2. **Check cache status** — `CACHE_HIT` vs `CACHE_MISS` in run metadata\n3. **Check resource utilization** — compare requested vs actual usage\n4. **Check data transfer** — large inline I/O indicates format issues\n5. **Check retry count** — frequent retries indicate instability\n\n## Retry and Timeout Configuration\n\n### Retries for resilience\n\n```python\n@env.task(retries=3)  # retry up to 3 times on failure\nasync def flaky_task(data: str) -> str:\n    \"\"\"Task that may fail transiently.\"\"\"\n    ...\n\n@env.task(retries=flyte.RetryStrategy(count=3))\nasync def critical_task(data: str) -> str:\n    \"\"\"Always retry, never fail fast.\"\"\"\n    ...\n```\n\n### Timeouts for bounding execution\n\n```python\n@env.task(max_runtime=\"1h\")  # bound single attempt\nasync def long_task(data: str) -> str:\n    ...\n\n@env.task(max_queued_time=\"30m\")  # fail fast if no capacity\nasync def urgent_task(data: str) -> str:\n    ...\n\n@env.task(deadline=\"2h\")  # bound total wall-clock (all attempts)\nasync def deadline_task(data: str) -> str:\n    ...\n```\n\n## Interruptible Tasks (Spot Instances)\n\n```python\n@env.task(interruptible=True)  # can be preempted, falls back to on-demand\nasync def spot_task(data: str) -> str:\n    \"\"\"Cost-effective for fault-tolerant workloads.\"\"\"\n    ...\n```\n\n## Optimization Anti-Patterns\n\n1. **Don't over-cache** — avoid `cache=\"auto\"` on tasks with side effects or non-deterministic outputs\n2. **Don't set resources too high** — over-provisioning wastes money; too low causes OOM\n3. **Don't use `asyncio.gather` for heavy workloads** — use `flyte.map` for parallel container execution\n4. **Don't skip caching on ETL** — idempotent data transforms should always cache\n5. **Don't use Union-only features** — avoid `ReusePolicy` and other Union-specific APIs\n6. **Don't ignore run metadata** — always check `durationMs` and cache status when debugging performance\n"
}

SHA-256 of public snapshot: 0550d3ae11ac24dcb8e9179372357e608636aed1159c48c2571d764c03e05adf