← 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": "Handles data engineering patterns: ETL pipelines, data processing, data quality checks, fanout/map tasks, conditions, dynamic workflows, and batch data transformations. Use when the user wants to build ETL pipelines, process large datasets, run data quality checks, fan out data processing tasks, or handle batch data transformations. Trigger words: \"ETL\", \"data pipeline\", \"data processing\", \"fanout\", \"map\", \"transform\", \"data quality\", \"parquet\", \"CSV\", \"batch\", \"extract\", \"load\", \"validate\", \"schema\".",
  "included_files": [],
  "name": "flyte-sdk-data",
  "skill_md_contents": "---\nname: flyte-sdk-data\ndescription: 'Handles data engineering patterns: ETL pipelines, data processing, data quality checks, fanout/map tasks, conditions, dynamic workflows, and batch data transformations. Use when the user wants to build ETL pipelines, process large datasets, run data quality checks, fan out data processing tasks, or handle batch data transformations. Trigger words: \"ETL\", \"data pipeline\", \"data processing\", \"fanout\", \"map\", \"transform\", \"data quality\", \"parquet\", \"CSV\", \"batch\", \"extract\", \"load\", \"validate\", \"schema\".'\n---\n\n# Flyte 2 SDK Data Engineering Skill\n\nBuild ETL pipelines, data processing workflows, and data quality checks with Flyte 2.\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## ETL Pipeline Patterns\n\n### Basic Extract-Transform-Load\n\n```python\nimport flyte\nimport flyte.io\n\nenv = flyte.TaskEnvironment(\n    name=\"etl-pipeline\",\n    image=flyte.Image.from_debian_base(python_version=(3, 12)).with_pip_packages(\n        \"pandas\", \"polars\", \"pyarrow\", \"boto3\", \"sqlalchemy\",\n    ),\n)\n\n@env.task(retries=3, cache=\"auto\")\nasync def extract(source_uri: str) -> flyte.io.DataFrame:\n    \"\"\"Extract data from various sources.\"\"\"\n    import polars as pl\n    if source_uri.endswith(\".csv\"):\n        df = pl.read_csv(source_uri)\n    elif source_uri.endswith(\".parquet\"):\n        df = pl.read_parquet(source_uri)\n    else:\n        raise ValueError(f\"Unsupported format: {source_uri}\")\n    return flyte.io.DataFrame(df)\n\n@env.task(retries=2, cache=\"auto\")\nasync def transform(df: flyte.io.DataFrame) -> flyte.io.DataFrame:\n    \"\"\"Clean and transform data.\"\"\"\n    inner = df.to_polars()\n    cleaned = (\n        inner\n        .drop_nulls()\n        .unique()\n        .with_columns([\n            pl.col(\"date\").str.strptime(pl.Date, \"%Y-%m-%d\").alias(\"date_parsed\"),\n        ])\n    )\n    return flyte.io.DataFrame(cleaned)\n\n@env.task(retries=1, cache=\"auto\")\nasync def load(df: flyte.io.DataFrame, destination: str) -> str:\n    \"\"\"Load transformed data to destination.\"\"\"\n    inner = df.to_polars()\n    if destination.endswith(\".parquet\"):\n        inner.write_parquet(destination)\n    elif destination.endswith(\".csv\"):\n        inner.write_csv(destination)\n    return destination\n\n@env.task\nasync def etl_pipeline(source_uri: str, destination: str) -> dict:\n    \"\"\"Orchestrate the ETL pipeline.\"\"\"\n    raw = await extract(source_uri)\n    cleaned = await transform(raw)\n    loaded_path = await load(cleaned, destination)\n    return {\"source\": source_uri, \"destination\": loaded_path}\n```\n\n### Multi-step ETL with intermediate storage\n\n```python\n@env.task(cache=\"auto\")\nasync def extract_raw(source_uri: str) -> flyte.io.File:\n    \"\"\"Extract and save raw data to remote storage.\"\"\"\n    import polars as pl\n    df = pl.read_csv(source_uri)\n    path = \"/tmp/raw.parquet\"\n    df.write_parquet(path)\n    return flyte.io.File(path=path)\n\n@env.task(cache=\"auto\")\nasync def validate_raw(raw: flyte.io.File) -> dict:\n    \"\"\"Validate raw data quality.\"\"\"\n    df = pl.read_parquet(raw.path)\n    return {\n        \"row_count\": len(df),\n        \"column_count\": len(df.columns),\n        \"null_counts\": df.null_count().to_dict(),\n    }\n\n@env.task(cache=\"auto\")\nasync def clean(raw: flyte.io.File) -> flyte.io.File:\n    \"\"\"Clean and normalize data.\"\"\"\n    df = pl.read_parquet(raw.path)\n    cleaned = df.drop_nulls().unique()\n    path = \"/tmp/cleaned.parquet\"\n    cleaned.write_parquet(path)\n    return flyte.io.File(path=path)\n\n@env.task(cache=\"auto\")\nasync def enrich(cleaned: flyte.io.File) -> flyte.io.File:\n    \"\"\"Enrich data with external features.\"\"\"\n    df = pl.read_parquet(cleaned.path)\n    # Join with external feature store\n    ...\n    path = \"/tmp/enriched.parquet\"\n    df.write_parquet(path)\n    return flyte.io.File(path=path)\n\n@env.task\nasync def load_enriched(enriched: flyte.io.File, destination: str) -> str:\n    \"\"\"Load enriched data to final destination.\"\"\"\n    ...\n    return destination\n\n@env.task\nasync def etl_with_validation(source_uri: str, destination: str) -> dict:\n    \"\"\"ETL pipeline with validation gates.\"\"\"\n    raw = await extract_raw(source_uri)\n    quality = await validate_raw(raw)\n\n    # Quality gate: fail if too many nulls\n    if quality[\"null_counts\"].get(\"critical_field\", 0) / quality[\"row_count\"] > 0.5:\n        raise ValueError(\"Too many nulls in critical field\")\n\n    cleaned = await clean(raw)\n    enriched = await enrich(cleaned)\n    loaded = await load_enriched(enriched, destination)\n    return {\"quality\": quality, \"destination\": loaded}\n```\n\n## Fan-out Data Processing\n\n### Map over large datasets\n\n```python\n@env.task(cache=\"auto\")\nasync def process_file(file_uri: str) -> flyte.io.DataFrame:\n    \"\"\"Process a single data file.\"\"\"\n    import polars as pl\n    df = pl.read_parquet(file_uri)\n    cleaned = df.drop_nulls().unique()\n    return flyte.io.DataFrame(cleaned)\n\n@env.task\nasync def process_dataset(file_uris: list[str]) -> list:\n    \"\"\"Fan out processing across all files in parallel.\"\"\"\n    results = await flyte.map(process_file, file_uris)\n    return results\n\n@env.task\nasync def merge_results(results: list) -> flyte.io.DataFrame:\n    \"\"\"Merge processed results into a single DataFrame.\"\"\"\n    import polars as pl\n    combined = pl.concat([r.to_polars() for r in results])\n    return flyte.io.DataFrame(combined)\n\n@env.task\nasync def main(file_uris: list[str]) -> flyte.io.DataFrame:\n    processed = await process_dataset(file_uris)\n    return await merge_results(processed)\n```\n\n### Fan-out with error handling\n\n```python\n@env.task\nasync def process_file_safe(file_uri: str) -> dict:\n    \"\"\"Process a file with error handling.\"\"\"\n    try:\n        df = await process_file(file_uri)\n        return {\"status\": \"success\", \"file\": file_uri, \"rows\": len(df.to_polars())}\n    except Exception as e:\n        return {\"status\": \"error\", \"file\": file_uri, \"error\": str(e)}\n\n@env.task\nasync def process_with_errors(file_uris: list[str]) -> dict:\n    \"\"\"Process files, collecting both successes and errors.\"\"\"\n    results = await flyte.map(process_file_safe, file_uris)\n    successes = [r for r in results if r[\"status\"] == \"success\"]\n    errors = [r for r in results if r[\"status\"] == \"error\"]\n    return {\"successes\": successes, \"errors\": errors, \"total\": len(results)}\n```\n\n### Limited concurrency fan-out\n\n```python\n@env.task\nasync def main(file_uris: list[str]) -> list:\n    \"\"\"Fan out with limited concurrency.\"\"\"\n    import asyncio\n    sem = asyncio.Semaphore(20)  # max 20 concurrent\n\n    async def bounded(uri):\n        async with sem:\n            return await process_file(uri)\n\n    return await asyncio.gather(*(bounded(u) for u in file_uris))\n```\n\n## Data Quality Checks\n\n### Comprehensive data quality\n\n```python\n@env.task(cache=\"auto\")\nasync def validate_schema(df: flyte.io.DataFrame, expected_schema: dict) -> dict:\n    \"\"\"Validate DataFrame schema matches expected schema.\"\"\"\n    inner = df.to_polars()\n    checks = {}\n\n    # Column names\n    expected_cols = set(expected_schema.keys())\n    actual_cols = set(inner.columns)\n    checks[\"columns_match\"] = expected_cols == actual_cols\n    checks[\"missing_columns\"] = list(expected_cols - actual_cols)\n    checks[\"extra_columns\"] = list(actual_cols - expected_cols)\n\n    # Column types\n    for col, expected_type in expected_schema.items():\n        if col in inner.columns:\n            actual_type = str(inner[col].dtype)\n            checks[f\"type_{col}\"] = {\n                \"expected\": expected_type,\n                \"actual\": actual_type,\n                \"match\": expected_type in actual_type,\n            }\n\n    return checks\n\n@env.task(cache=\"auto\")\nasync def validate_nulls(df: flyte.io.DataFrame, max_null_pct: float = 0.1) -> dict:\n    \"\"\"Validate null percentages per column.\"\"\"\n    inner = df.to_polars()\n    row_count = len(inner)\n    checks = {}\n\n    for col in inner.columns:\n        null_count = inner[col].null_count()\n        null_pct = null_count / row_count if row_count > 0 else 0\n        checks[col] = {\n            \"null_count\": null_count,\n            \"null_pct\": null_pct,\n            \"passed\": null_pct <= max_null_pct,\n        }\n\n    return checks\n\n@env.task(cache=\"auto\")\nasync def validate_values(df: flyte.io.DataFrame, constraints: dict) -> dict:\n    \"\"\"Validate value constraints (ranges, enums, patterns).\"\"\"\n    inner = df.to_polars()\n    checks = {}\n\n    for col, constraint in constraints.items():\n        if col not in inner.columns:\n            continue\n\n        if \"min\" in constraint:\n            checks[f\"{col}_min\"] = inner[col].min() >= constraint[\"min\"]\n        if \"max\" in constraint:\n            checks[f\"{col}_max\"] = inner[col].max() <= constraint[\"max\"]\n        if \"allowed_values\" in constraint:\n            unique = set(inner[col].unique())\n            checks[f\"{col}_values\"] = unique.issubset(set(constraint[\"allowed_values\"]))\n\n    return checks\n\n@env.task\nasync def data_quality_gate(\n    df: flyte.io.DataFrame,\n    schema: dict,\n    max_null_pct: float = 0.1,\n    constraints: dict = None,\n) -> dict:\n    \"\"\"Run all data quality checks and pass/fail.\"\"\"\n    schema_check = await validate_schema(df, schema)\n    null_check = await validate_nulls(df, max_null_pct)\n    value_checks = await validate_values(df, constraints or {})\n\n    all_passed = (\n        schema_check[\"columns_match\"]\n        and all(c[\"passed\"] for c in null_check.values())\n        and all(value_checks.values())\n    )\n\n    return {\n        \"passed\": all_passed,\n        \"schema\": schema_check,\n        \"nulls\": null_check,\n        \"values\": value_checks,\n    }\n```\n\n### Data quality with custom checks\n\n```python\n@env.task(cache=\"auto\")\nasync def check_distribution(df: flyte.io.DataFrame, column: str, expected_stats: dict) -> dict:\n    \"\"\"Check if data distribution matches expected statistics.\"\"\"\n    inner = df.to_polars()\n    col_data = inner[column].drop_nulls()\n\n    actual_mean = col_data.mean()\n    actual_std = col_data.std()\n    actual_min = col_data.min()\n    actual_max = col_data.max()\n\n    return {\n        \"column\": column,\n        \"mean\": {\"actual\": actual_mean, \"expected\": expected_stats.get(\"mean\"),\n                 \"within_tolerance\": abs(actual_mean - expected_stats.get(\"mean\", 0)) < expected_stats.get(\"tolerance\", 0.1)},\n        \"std\": {\"actual\": actual_std, \"expected\": expected_stats.get(\"std\"),\n                \"within_tolerance\": abs(actual_std - expected_stats.get(\"std\", 0)) < expected_stats.get(\"tolerance\", 0.1)},\n    }\n```\n\n## Dynamic Workflows for Data\n\n### Dynamic file processing\n\n```python\n@env.task\nasync def discover_files(prefix: str) -> list[str]:\n    \"\"\"Discover data files in a storage prefix.\"\"\"\n    import boto3\n    s3 = boto3.client(\"s3\")\n    files = []\n    paginator = s3.get_paginator(\"list_objects_v2\")\n    for page in paginator.paginate(Bucket=\"my-data-bucket\", Prefix=prefix):\n        for obj in page.get(\"Contents\", []):\n            if obj[\"Key\"].endswith((\".parquet\", \".csv\")):\n                files.append(f\"s3://{obj['Bucket']}/{obj['Key']}\")\n    return files\n\n@env.task\nasync def process_file(file_uri: str) -> flyte.io.File:\n    \"\"\"Process a single file.\"\"\"\n    import polars as pl\n    df = pl.read_parquet(file_uri)\n    cleaned = df.drop_nulls()\n    path = f\"/tmp/cleaned_{file_uri.split('/')[-1]}\"\n    cleaned.write_parquet(path)\n    return flyte.io.File(path=path)\n\n@env.task\nasync def merge_files(files: list[flyte.io.File]) -> flyte.io.File:\n    \"\"\"Merge processed files.\"\"\"\n    import polars as pl\n    dfs = [pl.read_parquet(f.path) for f in files]\n    combined = pl.concat(dfs)\n    path = \"/tmp/merged.parquet\"\n    combined.write_parquet(path)\n    return flyte.io.File(path=path)\n\n@env.task\nasync def dynamic_etl(prefix: str, destination: str) -> dict:\n    \"\"\"Dynamic ETL: discover files, process, merge.\"\"\"\n    files = await discover_files(prefix)\n    processed = await flyte.map(process_file, files)\n    merged = await merge_files(processed)\n    # Copy to destination\n    return {\"source_prefix\": prefix, \"destination\": destination, \"file_count\": len(files)}\n```\n\n### Conditional data routing\n\n```python\n@env.task\nasync def route_data(df: flyte.io.DataFrame, threshold: float) -> dict:\n    \"\"\"Route data based on quality score.\"\"\"\n    score = compute_quality_score(df)\n    if score >= threshold:\n        return {\"route\": \"production\", \"score\": score}\n    else:\n        return {\"route\": \"review\", \"score\": score}\n\n@env.task\nasync def process_production(df: flyte.io.DataFrame) -> flyte.io.File:\n    \"\"\"Process data for production.\"\"\"\n    ...\n\n@env.task\nasync def process_review(df: flyte.io.DataFrame) -> flyte.io.File:\n    \"\"\"Flag data for manual review.\"\"\"\n    ...\n\n@env.task\nasync def conditional_pipeline(df: flyte.io.DataFrame, threshold: float) -> dict:\n    \"\"\"Route data based on quality.\"\"\"\n    routed = await route_data(df, threshold)\n    if routed[\"route\"] == \"production\":\n        result = await process_production(df)\n    else:\n        result = await process_review(df)\n    return {**routed, \"result\": result}\n```\n\n## JsonlFile and JsonlDir for Large Datasets\n\n### JsonlFile — streaming JSONL\n\n```python\n@env.task\nasync def process_jsonl(path: str) -> int:\n    \"\"\"Process a JSONL file with streaming.\"\"\"\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\n### JsonlDir — batched JSONL directories\n\n```python\n@env.task\nasync def process_jsonl_dir(dir_path: str) -> dict:\n    \"\"\"Process JSONL directory with batched streaming.\"\"\"\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## Data Format Reference\n\n| Format | Flyte Type | Best For |\n|---|---|---|\n| Parquet | `flyte.io.DataFrame` | Tabular data, ETL |\n| CSV | `flyte.io.File` | Small datasets, interchange |\n| JSONL | `JsonlFile` / `JsonlDir` | Streaming records |\n| JSON | inline (dict) | Small structured data |\n| Pickle | `flyte.io.File` | Python objects |\n| NumPy (.npy/.npz) | `flyte.io.File` | Arrays, embeddings |\n| PNG/JPEG | `flyte.io.File` | Images |\n| Model (.pt/.safetensors) | `flyte.io.File` | Model checkpoints |\n\n## Performance Tips for Data Pipelines\n\n1. **Use Parquet over CSV** — columnar format, compressed, faster I/O\n2. **Cache idempotent transforms** — `cache=\"auto\"` on ETL steps\n3. **Fan out with `flyte.map`** — parallel processing for independent files\n4. **Use `flyte.trace` for lightweight ops** — no container spin-up cost\n5. **Set `raw_data_path`** — control where intermediate data is stored\n6. **Use `inline_output_limit`** — control when data goes by reference vs inline\n7. **Use `interruptible=True`** — spot instances for fault-tolerant data processing\n\n## Anti-Patterns\n\n1. **Don't load entire datasets into memory** — use streaming (`JsonlFile`, Polars lazy) for large data.\n2. **Don't pass DataFrames inline** — they go by reference automatically, but small dicts do go inline.\n3. **Don't skip data quality gates** — always validate before and after transforms.\n4. **Don't use Union-only features** — avoid `ReusePolicy` and other Union-specific APIs.\n"
}

SHA-256 of public snapshot: c49b3808e476eff1de889eabc3850d8c1899ed252fd49c69e434455552e97e31