← astronomer-dataCONTENT HISTORY

Update to astronomer-data

Snapshot Sep 30, 2026 · 23:17 UTC · version 0.1.0

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
{
  "name": "creating-openlineage-extractors",
  "description": "Create custom OpenLineage extractors for Airflow operators. Use when the user needs lineage from unsupported or third-party operators, wants column-level lineage, or needs complex extraction logic beyond what inlets/outlets provide.",
  "included_files": [],
  "skill_md_contents": "---\nname: creating-openlineage-extractors\ndescription: Create custom OpenLineage extractors for Airflow operators. Use when the user needs lineage from unsupported or third-party operators, wants column-level lineage, or needs complex extraction logic beyond what inlets/outlets provide.\n---\n\n# Creating OpenLineage Extractors\n\nThis skill guides you through creating custom OpenLineage extractors to capture lineage from Airflow operators that don't have built-in support.\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 patterns and list of supported operators/hooks.\n\n## When to Use Each Approach\n\n| Scenario | Approach |\n|----------|----------|\n| Operator you own/maintain | **OpenLineage Methods** (recommended, simplest) |\n| Third-party operator you can't modify | Custom Extractor |\n| Need column-level lineage | OpenLineage Methods or Custom Extractor |\n| Complex extraction logic | OpenLineage Methods or Custom Extractor |\n| Simple table-level lineage | Inlets/Outlets (simplest, but lowest priority) |\n\n> **Important:** Always prefer OpenLineage methods over custom extractors when possible. Extractors are harder to write, easier to diverge from operator behavior after changes, and harder to debug.\n\n### On Astro\n\nAstro includes built-in OpenLineage integration — no additional transport configuration is needed. Lineage events are automatically collected and displayed in the Astro UI's **Lineage tab**. Custom extractors deployed to an Astro project are automatically picked up, so you only need to register them in `airflow.cfg` or via environment variable and deploy.\n\n---\n\n## Two Approaches\n\n### 1. OpenLineage Methods (Recommended)\n\nUse when you can add methods directly to your custom operator. This is the **go-to solution** for operators you own.\n\n### 2. Custom Extractors\n\nUse when you need lineage from third-party or provider operators that you **cannot modify**.\n\n---\n\n## Approach 1: OpenLineage Methods (Recommended)\n\nWhen you own the operator, add OpenLineage methods directly:\n\n```python\nfrom airflow.models import BaseOperator\n\n\nclass MyCustomOperator(BaseOperator):\n    \"\"\"Custom operator with built-in OpenLineage support.\"\"\"\n\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        self._rows_processed = 0  # Set during execution\n\n    def execute(self, context):\n        # Do the actual work\n        self._rows_processed = self._process_data()\n        return self._rows_processed\n\n    def get_openlineage_facets_on_start(self):\n        \"\"\"Called when task starts. Return known inputs/outputs.\"\"\"\n        # Import locally to avoid circular imports\n        from openlineage.client.event_v2 import Dataset\n        from airflow.providers.openlineage.extractors import OperatorLineage\n\n        return OperatorLineage(\n            inputs=[Dataset(namespace=\"postgres://db\", name=self.source_table)],\n            outputs=[Dataset(namespace=\"postgres://db\", name=self.target_table)],\n        )\n\n    def get_openlineage_facets_on_complete(self, task_instance):\n        \"\"\"Called after success. Add runtime metadata.\"\"\"\n        from openlineage.client.event_v2 import Dataset\n        from openlineage.client.facet_v2 import output_statistics_output_dataset\n        from airflow.providers.openlineage.extractors import OperatorLineage\n\n        return OperatorLineage(\n            inputs=[Dataset(namespace=\"postgres://db\", name=self.source_table)],\n            outputs=[\n                Dataset(\n                    namespace=\"postgres://db\",\n                    name=self.target_table,\n                    facets={\n                        \"outputStatistics\": output_statistics_output_dataset.OutputStatisticsOutputDatasetFacet(\n                            rowCount=self._rows_processed\n                        )\n                    },\n                )\n            ],\n        )\n\n    def get_openlineage_facets_on_failure(self, task_instance):\n        \"\"\"Called after failure. Optional - for partial lineage.\"\"\"\n        return None\n```\n\n### OpenLineage Methods Reference\n\n| Method | When Called | Required |\n|--------|-------------|----------|\n| `get_openlineage_facets_on_start()` | Task enters RUNNING | No |\n| `get_openlineage_facets_on_complete(ti)` | Task succeeds | No |\n| `get_openlineage_facets_on_failure(ti)` | Task fails | No |\n\n> Implement only the methods you need. Unimplemented methods fall through to Hook-Level Lineage or inlets/outlets.\n\n---\n\n## Approach 2: Custom Extractors\n\nUse this approach only when you **cannot modify** the operator (e.g., third-party or provider operators).\n\n### Basic Structure\n\n```python\nfrom airflow.providers.openlineage.extractors.base import BaseExtractor, OperatorLineage\nfrom openlineage.client.event_v2 import Dataset\n\n\nclass MyOperatorExtractor(BaseExtractor):\n    \"\"\"Extract lineage from MyCustomOperator.\"\"\"\n\n    @classmethod\n    def get_operator_classnames(cls) -> list[str]:\n        \"\"\"Return operator class names this extractor handles.\"\"\"\n        return [\"MyCustomOperator\"]\n\n    def _execute_extraction(self) -> OperatorLineage | None:\n        \"\"\"Called BEFORE operator executes. Use for known inputs/outputs.\"\"\"\n        # Access operator properties via self.operator\n        source_table = self.operator.source_table\n        target_table = self.operator.target_table\n\n        return OperatorLineage(\n            inputs=[\n                Dataset(\n                    namespace=\"postgres://mydb:5432\",\n                    name=f\"public.{source_table}\",\n                )\n            ],\n            outputs=[\n                Dataset(\n                    namespace=\"postgres://mydb:5432\",\n                    name=f\"public.{target_table}\",\n                )\n            ],\n        )\n\n    def extract_on_complete(self, task_instance) -> OperatorLineage | None:\n        \"\"\"Called AFTER operator executes. Use for runtime-determined lineage.\"\"\"\n        # Access properties set during execution\n        # Useful for operators that determine outputs at runtime\n        return None\n```\n\n### OperatorLineage Structure\n\n```python\nfrom airflow.providers.openlineage.extractors.base import OperatorLineage\nfrom openlineage.client.event_v2 import Dataset\nfrom openlineage.client.facet_v2 import sql_job\n\nlineage = OperatorLineage(\n    inputs=[Dataset(namespace=\"...\", name=\"...\")],      # Input datasets\n    outputs=[Dataset(namespace=\"...\", name=\"...\")],     # Output datasets\n    run_facets={\"sql\": sql_job.SQLJobFacet(query=\"SELECT...\")},  # Run metadata\n    job_facets={},                                      # Job metadata\n)\n```\n\n### Extraction Methods\n\n| Method | When Called | Use For |\n|--------|-------------|---------|\n| `_execute_extraction()` | Before operator runs | Static/known lineage |\n| `extract_on_complete(task_instance)` | After success | Runtime-determined lineage |\n| `extract_on_failure(task_instance)` | After failure | Partial lineage on errors |\n\n### Registering Extractors\n\n**Option 1: Configuration file (`airflow.cfg`)**\n\n```ini\n[openlineage]\nextractors = mypackage.extractors.MyOperatorExtractor;mypackage.extractors.AnotherExtractor\n```\n\n**Option 2: Environment variable**\n\n```bash\nAIRFLOW__OPENLINEAGE__EXTRACTORS='mypackage.extractors.MyOperatorExtractor;mypackage.extractors.AnotherExtractor'\n```\n\n> **Important:** The path must be importable from the Airflow worker. Place extractors in your DAGs folder or installed package.\n\n---\n\n## Common Patterns\n\n### SQL Operator Extractor\n\n```python\nfrom airflow.providers.openlineage.extractors.base import BaseExtractor, OperatorLineage\nfrom openlineage.client.event_v2 import Dataset\nfrom openlineage.client.facet_v2 import sql_job\n\n\nclass MySqlOperatorExtractor(BaseExtractor):\n    @classmethod\n    def get_operator_classnames(cls) -> list[str]:\n        return [\"MySqlOperator\"]\n\n    def _execute_extraction(self) -> OperatorLineage | None:\n        sql = self.operator.sql\n        conn_id = self.operator.conn_id\n\n        # Parse SQL to find tables (simplified example)\n        # In practice, use a SQL parser like sqlglot\n        inputs, outputs = self._parse_sql(sql)\n\n        namespace = f\"postgres://{conn_id}\"\n\n        return OperatorLineage(\n            inputs=[Dataset(namespace=namespace, name=t) for t in inputs],\n            outputs=[Dataset(namespace=namespace, name=t) for t in outputs],\n            job_facets={\n                \"sql\": sql_job.SQLJobFacet(query=sql)\n            },\n        )\n\n    def _parse_sql(self, sql: str) -> tuple[list[str], list[str]]:\n        \"\"\"Parse SQL to extract table names. Use sqlglot for real parsing.\"\"\"\n        # Simplified example - use proper SQL parser in production\n        inputs = []\n        outputs = []\n        # ... parsing logic ...\n        return inputs, outputs\n```\n\n### File Transfer Extractor\n\n```python\nfrom airflow.providers.openlineage.extractors.base import BaseExtractor, OperatorLineage\nfrom openlineage.client.event_v2 import Dataset\n\n\nclass S3ToSnowflakeExtractor(BaseExtractor):\n    @classmethod\n    def get_operator_classnames(cls) -> list[str]:\n        return [\"S3ToSnowflakeOperator\"]\n\n    def _execute_extraction(self) -> OperatorLineage | None:\n        s3_bucket = self.operator.s3_bucket\n        s3_key = self.operator.s3_key\n        table = self.operator.table\n        schema = self.operator.schema\n\n        return OperatorLineage(\n            inputs=[\n                Dataset(\n                    namespace=f\"s3://{s3_bucket}\",\n                    name=s3_key,\n                )\n            ],\n            outputs=[\n                Dataset(\n                    namespace=\"snowflake://myaccount.snowflakecomputing.com\",\n                    name=f\"{schema}.{table}\",\n                )\n            ],\n        )\n```\n\n### Dynamic Lineage from Execution\n\n```python\nfrom openlineage.client.event_v2 import Dataset\n\n\nclass DynamicOutputExtractor(BaseExtractor):\n    @classmethod\n    def get_operator_classnames(cls) -> list[str]:\n        return [\"DynamicOutputOperator\"]\n\n    def _execute_extraction(self) -> OperatorLineage | None:\n        # Only inputs known before execution\n        return OperatorLineage(\n            inputs=[Dataset(namespace=\"...\", name=self.operator.source)],\n        )\n\n    def extract_on_complete(self, task_instance) -> OperatorLineage | None:\n        # Outputs determined during execution\n        # Access via operator properties set in execute()\n        outputs = self.operator.created_tables  # Set during execute()\n\n        return OperatorLineage(\n            inputs=[Dataset(namespace=\"...\", name=self.operator.source)],\n            outputs=[Dataset(namespace=\"...\", name=t) for t in outputs],\n        )\n```\n\n---\n\n## Common Pitfalls\n\n### 1. Circular Imports\n\n**Problem:** Importing Airflow modules at the top level causes circular imports.\n\n```python\n# ❌ BAD - can cause circular import issues\nfrom airflow.models import TaskInstance\nfrom openlineage.client.event_v2 import Dataset\n\nclass MyExtractor(BaseExtractor):\n    ...\n```\n\n```python\n# ✅ GOOD - import inside methods\nclass MyExtractor(BaseExtractor):\n    def _execute_extraction(self):\n        from openlineage.client.event_v2 import Dataset\n        # ...\n```\n\n### 2. Wrong Import Path\n\n**Problem:** Extractor path doesn't match actual module location.\n\n```bash\n# ❌ Wrong - path doesn't exist\nAIRFLOW__OPENLINEAGE__EXTRACTORS='extractors.MyExtractor'\n\n# ✅ Correct - full importable path\nAIRFLOW__OPENLINEAGE__EXTRACTORS='dags.extractors.my_extractor.MyExtractor'\n```\n\n### 3. Not Handling None\n\n**Problem:** Extraction fails when operator properties are None.\n\n```python\n# ✅ Handle optional properties\ndef _execute_extraction(self) -> OperatorLineage | None:\n    if not self.operator.source_table:\n        return None  # Skip extraction\n\n    return OperatorLineage(...)\n```\n\n---\n\n## Testing Extractors\n\n### Unit Testing\n\n```python\nimport pytest\nfrom unittest.mock import MagicMock\nfrom mypackage.extractors import MyOperatorExtractor\n\n\ndef test_extractor():\n    # Mock the operator\n    operator = MagicMock()\n    operator.source_table = \"input_table\"\n    operator.target_table = \"output_table\"\n\n    # Create extractor\n    extractor = MyOperatorExtractor(operator)\n\n    # Test extraction\n    lineage = extractor._execute_extraction()\n\n    assert len(lineage.inputs) == 1\n    assert lineage.inputs[0].name == \"input_table\"\n    assert len(lineage.outputs) == 1\n    assert lineage.outputs[0].name == \"output_table\"\n```\n\n---\n\n## Precedence Rules\n\nOpenLineage checks for lineage in this order:\n\n1. **Custom Extractors** (highest priority)\n2. **OpenLineage Methods** on operator\n3. **Hook-Level Lineage** (from `HookLineageCollector`)\n4. **Inlets/Outlets** (lowest priority)\n\nIf a custom extractor exists, it overrides built-in extraction and inlets/outlets.\n\n---\n\n## Related Skills\n\n- **annotating-task-lineage**: For simple table-level lineage with inlets/outlets\n- **tracing-upstream-lineage**: Investigate data origins\n- **tracing-downstream-lineage**: Investigate data dependencies\n"
}

SHA-256: 920eca42c21f9dbbd361da396577a9818275e53ffd5d8cdb82bb3819547ebe67