← 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": "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