← astronomer-dataCONTENT HISTORYWHAT CHANGED · RULE-BASED ANALYSIS
Update to astronomer-data
Snapshot Sep 30, 2026 · 23:17 UTC · version 0.1.0
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
{
"name": "annotating-task-lineage",
"description": "Annotate Airflow tasks with data lineage using inlets and outlets. Use when the user wants to add lineage metadata to tasks, specify input/output datasets, or enable lineage tracking for operators without built-in OpenLineage extraction.",
"included_files": [],
"skill_md_contents": "---\nname: annotating-task-lineage\ndescription: Annotate Airflow tasks with data lineage using inlets and outlets. Use when the user wants to add lineage metadata to tasks, specify input/output datasets, or enable lineage tracking for operators without built-in OpenLineage extraction.\n---\n\n# Annotating Task Lineage with Inlets & Outlets\n\nThis skill guides you through adding manual lineage annotations to Airflow tasks using `inlets` and `outlets`.\n\n> **Reference:** See the [OpenLineage provider developer guide](https://airflow.apache.org/docs/apache-airflow-providers-openlineage/stable/guides/developer.html) for the latest supported operators and patterns.\n\n### On Astro\n\nLineage annotations defined with inlets and outlets are visualized in Astro's enhanced **Lineage tab**, which provides cross-DAG and cross-deployment lineage views. This means your annotations are immediately visible in the Astro UI, giving you a unified view of data flow across your entire Astro organization.\n\n## When to Use This Approach\n\n| Scenario | Use Inlets/Outlets? |\n|----------|---------------------|\n| Operator has OpenLineage methods (`get_openlineage_facets_on_*`) | ❌ Modify the OL method directly |\n| Operator has no built-in OpenLineage extractor | ✅ Yes |\n| Simple table-level lineage is sufficient | ✅ Yes |\n| Quick lineage setup without custom code | ✅ Yes |\n| Need column-level lineage | ❌ Use OpenLineage methods or custom extractor |\n| Complex extraction logic needed | ❌ Use OpenLineage methods or custom extractor |\n\n> **Note:** Inlets/outlets are the lowest-priority fallback. If an OpenLineage extractor or method exists for the operator, it takes precedence. Use this approach for operators without extractors.\n\n---\n\n## Supported Types for Inlets/Outlets\n\nYou can use **OpenLineage Dataset** objects or **Airflow Assets** for inlets and outlets:\n\n### OpenLineage Datasets (Recommended)\n\n```python\nfrom openlineage.client.event_v2 import Dataset\n\n# Database tables\nsource_table = Dataset(\n namespace=\"postgres://mydb:5432\",\n name=\"public.orders\",\n)\ntarget_table = Dataset(\n namespace=\"snowflake://account.snowflakecomputing.com\",\n name=\"staging.orders_clean\",\n)\n\n# Files\ninput_file = Dataset(\n namespace=\"s3://my-bucket\",\n name=\"raw/events/2024-01-01.json\",\n)\n```\n\n### Airflow Assets (Airflow 3+)\n\n```python\nfrom airflow.sdk import Asset\n\n# Using Airflow's native Asset type\norders_asset = Asset(uri=\"s3://my-bucket/data/orders\")\n```\n\n### Airflow Datasets (Airflow 2.4+)\n\n```python\nfrom airflow.datasets import Dataset\n\n# Using Airflow's Dataset type (Airflow 2.4-2.x)\norders_dataset = Dataset(uri=\"s3://my-bucket/data/orders\")\n```\n\n---\n\n## Basic Usage\n\n### Setting Inlets and Outlets on Operators\n\n```python\nfrom airflow import DAG\nfrom airflow.operators.bash import BashOperator\nfrom openlineage.client.event_v2 import Dataset\nimport pendulum\n\n# Define your lineage datasets\nsource_table = Dataset(\n namespace=\"snowflake://account.snowflakecomputing.com\",\n name=\"raw.orders\",\n)\ntarget_table = Dataset(\n namespace=\"snowflake://account.snowflakecomputing.com\",\n name=\"staging.orders_clean\",\n)\noutput_file = Dataset(\n namespace=\"s3://my-bucket\",\n name=\"exports/orders.parquet\",\n)\n\nwith DAG(\n dag_id=\"etl_with_lineage\",\n start_date=pendulum.datetime(2024, 1, 1, tz=\"UTC\"),\n schedule=\"@daily\",\n) as dag:\n\n transform = BashOperator(\n task_id=\"transform_orders\",\n bash_command=\"echo 'transforming...'\",\n inlets=[source_table], # What this task reads\n outlets=[target_table], # What this task writes\n )\n\n export = BashOperator(\n task_id=\"export_to_s3\",\n bash_command=\"echo 'exporting...'\",\n inlets=[target_table], # Reads from previous output\n outlets=[output_file], # Writes to S3\n )\n\n transform >> export\n```\n\n### Multiple Inputs and Outputs\n\nTasks often read from multiple sources and write to multiple destinations:\n\n```python\nfrom openlineage.client.event_v2 import Dataset\n\n# Multiple source tables\ncustomers = Dataset(namespace=\"postgres://crm:5432\", name=\"public.customers\")\norders = Dataset(namespace=\"postgres://sales:5432\", name=\"public.orders\")\nproducts = Dataset(namespace=\"postgres://inventory:5432\", name=\"public.products\")\n\n# Multiple output tables\ndaily_summary = Dataset(namespace=\"snowflake://account\", name=\"analytics.daily_summary\")\ncustomer_metrics = Dataset(namespace=\"snowflake://account\", name=\"analytics.customer_metrics\")\n\naggregate_task = PythonOperator(\n task_id=\"build_daily_aggregates\",\n python_callable=build_aggregates,\n inlets=[customers, orders, products], # All inputs\n outlets=[daily_summary, customer_metrics], # All outputs\n)\n```\n\n---\n\n## Setting Lineage in Custom Operators\n\nWhen building custom operators, you have two options:\n\n### Option 1: Implement OpenLineage Methods (Recommended)\n\nThis is the preferred approach as it gives you full control over lineage extraction:\n\n```python\nfrom airflow.models import BaseOperator\n\n\nclass MyCustomOperator(BaseOperator):\n def __init__(self, source_table: str, target_table: str, **kwargs):\n super().__init__(**kwargs)\n self.source_table = source_table\n self.target_table = target_table\n\n def execute(self, context):\n # ... perform the actual work ...\n self.log.info(f\"Processing {self.source_table} -> {self.target_table}\")\n\n def get_openlineage_facets_on_complete(self, task_instance):\n \"\"\"Return lineage after successful execution.\"\"\"\n from openlineage.client.event_v2 import Dataset\n from airflow.providers.openlineage.extractors import OperatorLineage\n\n return OperatorLineage(\n inputs=[Dataset(namespace=\"warehouse://db\", name=self.source_table)],\n outputs=[Dataset(namespace=\"warehouse://db\", name=self.target_table)],\n )\n```\n\n### Option 2: Set Inlets/Outlets Dynamically\n\nFor simpler cases, set lineage within the `execute` method (non-deferrable operators only):\n\n```python\nfrom airflow.models import BaseOperator\nfrom openlineage.client.event_v2 import Dataset\n\n\nclass MyCustomOperator(BaseOperator):\n def __init__(self, source_table: str, target_table: str, **kwargs):\n super().__init__(**kwargs)\n self.source_table = source_table\n self.target_table = target_table\n\n def execute(self, context):\n # Set lineage dynamically based on operator parameters\n self.inlets = [\n Dataset(namespace=\"warehouse://db\", name=self.source_table)\n ]\n self.outlets = [\n Dataset(namespace=\"warehouse://db\", name=self.target_table)\n ]\n\n # ... perform the actual work ...\n self.log.info(f\"Processing {self.source_table} -> {self.target_table}\")\n```\n\n---\n\n## Dataset Naming Helpers\n\nUse the [OpenLineage dataset naming helpers](https://openlineage.io/docs/client/python/best-practices#dataset-naming-helpers) to ensure consistent naming across platforms:\n\n```python\nfrom openlineage.client.event_v2 import Dataset\n\n# Snowflake\nfrom openlineage.client.naming.snowflake import SnowflakeDatasetNaming\n\nnaming = SnowflakeDatasetNaming(\n account_identifier=\"myorg-myaccount\",\n database=\"mydb\",\n schema=\"myschema\",\n table=\"mytable\",\n)\ndataset = Dataset(namespace=naming.get_namespace(), name=naming.get_name())\n# -> namespace: \"snowflake://myorg-myaccount\", name: \"mydb.myschema.mytable\"\n\n# BigQuery\nfrom openlineage.client.naming.bigquery import BigQueryDatasetNaming\n\nnaming = BigQueryDatasetNaming(\n project=\"my-project\",\n dataset=\"my_dataset\",\n table=\"my_table\",\n)\ndataset = Dataset(namespace=naming.get_namespace(), name=naming.get_name())\n# -> namespace: \"bigquery\", name: \"my-project.my_dataset.my_table\"\n\n# S3\nfrom openlineage.client.naming.s3 import S3DatasetNaming\n\nnaming = S3DatasetNaming(bucket=\"my-bucket\", key=\"path/to/file.parquet\")\ndataset = Dataset(namespace=naming.get_namespace(), name=naming.get_name())\n# -> namespace: \"s3://my-bucket\", name: \"path/to/file.parquet\"\n\n# PostgreSQL\nfrom openlineage.client.naming.postgres import PostgresDatasetNaming\n\nnaming = PostgresDatasetNaming(\n host=\"localhost\",\n port=5432,\n database=\"mydb\",\n schema=\"public\",\n table=\"users\",\n)\ndataset = Dataset(namespace=naming.get_namespace(), name=naming.get_name())\n# -> namespace: \"postgres://localhost:5432\", name: \"mydb.public.users\"\n```\n\n> **Note:** Always use the naming helpers instead of constructing namespaces manually. If a helper is missing for your platform, check the [OpenLineage repo](https://github.com/OpenLineage/OpenLineage) or request it.\n\n---\n\n## Precedence Rules\n\nOpenLineage uses this precedence for lineage extraction:\n\n1. **Custom Extractors** (highest) - User-registered extractors\n2. **OpenLineage Methods** - `get_openlineage_facets_on_*` in operator\n3. **Hook-Level Lineage** - Lineage collected from hooks via `HookLineageCollector`\n4. **Inlets/Outlets** (lowest) - Falls back to these if nothing else extracts lineage\n\n> **Note:** If an extractor or method exists but returns no datasets, OpenLineage will check hook-level lineage, then fall back to inlets/outlets.\n\n---\n\n## Best Practices\n\n### Use the Naming Helpers\n\nAlways use OpenLineage naming helpers for consistent dataset creation:\n\n```python\nfrom openlineage.client.event_v2 import Dataset\nfrom openlineage.client.naming.snowflake import SnowflakeDatasetNaming\n\n\ndef snowflake_dataset(schema: str, table: str) -> Dataset:\n \"\"\"Create a Snowflake Dataset using the naming helper.\"\"\"\n naming = SnowflakeDatasetNaming(\n account_identifier=\"mycompany\",\n database=\"analytics\",\n schema=schema,\n table=table,\n )\n return Dataset(namespace=naming.get_namespace(), name=naming.get_name())\n\n\n# Usage\nsource = snowflake_dataset(\"raw\", \"orders\")\ntarget = snowflake_dataset(\"staging\", \"orders_clean\")\n```\n\n### Document Your Lineage\n\nAdd comments explaining the data flow:\n\n```python\ntransform = SqlOperator(\n task_id=\"transform_orders\",\n sql=\"...\",\n # Lineage: Reads raw orders, joins with customers, writes to staging\n inlets=[\n snowflake_dataset(\"raw\", \"orders\"),\n snowflake_dataset(\"raw\", \"customers\"),\n ],\n outlets=[\n snowflake_dataset(\"staging\", \"order_details\"),\n ],\n)\n```\n\n### Keep Lineage Accurate\n\n- Update inlets/outlets when SQL queries change\n- Include all tables referenced in JOINs as inlets\n- Include all tables written to (including temp tables if relevant)\n- **Outlet-only and inlet-only annotations are valid.** One-sided annotations are encouraged for lineage visibility even without a corresponding inlet or outlet in another DAG.\n\n---\n\n## Limitations\n\n| Limitation | Workaround |\n|------------|------------|\n| Table-level only (no column lineage) | Use OpenLineage methods or custom extractor |\n| Overridden by extractors/methods | Only use for operators without extractors |\n| Static at DAG parse time | Set dynamically in `execute()` or use OL methods |\n| Deferrable operators lose dynamic lineage | Use OL methods instead; attributes set in `execute()` are lost when deferring |\n\n---\n\n## Related Skills\n\n- **creating-openlineage-extractors**: For column-level lineage or complex extraction\n- **tracing-upstream-lineage**: Investigate where data comes from\n- **tracing-downstream-lineage**: Investigate what depends on data\n"
}SHA-256: 2a892c468389606bc80b8554ff85870250890b28f04be053c552df3e90fd9df2