agentsclimarketplace

Openlineage

Skill ivanshamaev/de-agent-skills/skills/openlineage

Профессиональные Data Engineering Agent Skills для разработки AI Agentic Data Platform

Install
npx -y skills add ivanshamaev/de-agent-skills --skill openlineage

Assembled from the repository path, not quoted from the project. Check it against their README if it does not work.

2 things to look at

  • no licenseNo license file was found in the repository. Code published without one is not open source by default, so using it at work is a question for whoever answers licensing questions where you are.
  • 13 stars13 stars. Stars are a popularity signal and not a quality one, but at this level it is likely that nobody has read this closely except its author, and you would be relying on your own review.

What its author says it does

Copied from the file, not written here

OpenLineage data lineage tracking — RunEvent/Job/Dataset/facet spec, Marquez backend setup, Airflow/Spark/dbt integrations, column-level lineage, custom emitters, lineage-based impact analysis

SKILL.md

22.3 KB, as published. Nobody here has run it

OpenLineage Data Lineage

When to Use

Activate this skill when the task involves:

  • Implementing data lineage tracking across Airflow, Spark, dbt, or Trino pipelines
  • Setting up the OpenLineage spec emitters and a Marquez or OpenMetadata backend
  • Designing column-level lineage for impact analysis
  • Writing custom OpenLineage clients or enriching events with facets
  • Debugging missing or broken lineage edges in the lineage graph
  • Integrating lineage events with data catalogs (DataHub, Atlan, OpenMetadata)

Core Model

┌─────────────────────────────────────────────────────────────┐
│  OpenLineage Model                                          │
│                                                             │
│   Dataset ──────── Job ──────── Dataset                    │
│  (input)          (Run)         (output)                    │
│                                                             │
│   Each entity carries Facets — atomic metadata blocks:     │
│   • Schema facet: column names + types                      │
│   • Column-level lineage: input col → output col mapping   │
│   • Data quality assertions                                 │
│   • SQL query text                                          │
│   • Source code location                                    │
└─────────────────────────────────────────────────────────────┘

Events flow:  Pipeline Tool → OpenLineage Client → HTTP Transport → Backend
              (Airflow, Spark, dbt)                                 (Marquez / DataHub)

Lineage is collected passively — tools emit events without modifying pipeline logic.


RunEvent Specification

The canonical event JSON for an OpenLineage RunEvent:

{
  "eventType": "COMPLETE",
  "eventTime": "2024-03-15T10:30:00.000Z",
  "producer": "https://github.com/OpenLineage/OpenLineage/tree/1.0.0/integration/spark",
  "schemaURL": "https://openlineage.io/spec/1-0-5/OpenLineage.json#/definitions/RunEvent",

  "run": {
    "runId": "d46e465b-d358-4d32-83d4-df660ff614dd",
    "facets": {
      "nominalTime": {
        "_producer": "...",
        "_schemaURL": "...",
        "nominalStartTime": "2024-03-15T10:00:00Z",
        "nominalEndTime":   "2024-03-15T10:30:00Z"
      },
      "parent": {
        "_producer": "...",
        "_schemaURL": "...",
        "run":  {"runId": "e9c85741-93ab-4b27-9c8a-f3f4a3c0e001"},
        "job":  {"namespace": "airflow", "name": "etl_pipeline.run_spark_job"}
      }
    }
  },

  "job": {
    "namespace": "spark://spark-master:7077",
    "name": "silver.transform_orders",
    "facets": {
      "sql": {
        "_producer": "...",
        "_schemaURL": "...",
        "query": "INSERT INTO silver.orders SELECT id, customer_id, total FROM bronze.orders WHERE status = 'valid'"
      },
      "sourceCodeLocation": {
        "_producer": "...",
        "_schemaURL": "...",
        "type": "git",
        "url": "https://github.com/org/repo",
        "repoUrl": "https://github.com/org/repo",
        "path": "jobs/transform_orders.py",
        "version": "abc123"
      }
    }
  },

  "inputs": [
    {
      "namespace": "spark://spark-master:7077",
      "name": "bronze.orders",
      "facets": {
        "schema": {
          "_producer": "...",
          "_schemaURL": "...",
          "fields": [
            {"name": "id",          "type": "BIGINT"},
            {"name": "customer_id", "type": "BIGINT"},
            {"name": "total",       "type": "DECIMAL(10,2)"},
            {"name": "status",      "type": "VARCHAR"}
          ]
        },
        "dataSource": {
          "_producer": "...",
          "_schemaURL": "...",
          "name": "spark_iceberg",
          "uri":  "iceberg://lakehouse/bronze"
        }
      }
    }
  ],

  "outputs": [
    {
      "namespace": "spark://spark-master:7077",
      "name": "silver.orders",
      "facets": {
        "schema": {
          "_producer": "...",
          "_schemaURL": "...",
          "fields": [
            {"name": "id",          "type": "BIGINT"},
            {"name": "customer_id", "type": "BIGINT"},
            {"name": "total",       "type": "DECIMAL(10,2)"}
          ]
        },
        "columnLineage": {
          "_producer": "...",
          "_schemaURL": "...",
          "fields": {
            "id": {
              "inputFields": [
                {"namespace": "spark://spark-master:7077", "name": "bronze.orders", "field": "id"}
              ]
            },
            "customer_id": {
              "inputFields": [
                {"namespace": "spark://spark-master:7077", "name": "bronze.orders", "field": "customer_id"}
              ]
            },
            "total": {
              "inputFields": [
                {"namespace": "spark://spark-master:7077", "name": "bronze.orders", "field": "total"}
              ]
            }
          }
        },
        "outputStatistics": {
          "_producer": "...",
          "_schemaURL": "...",
          "rowCount": 150000,
          "size": 45000000
        }
      }
    }
  ]
}

Run States

eventTypeMeaningWhen to Emit
STARTJob execution beganBefore first data read
RUNNINGPeriodic progress updateLong-running jobs (checkpoints)
COMPLETEFinished successfullyAfter last write
FAILExecution failedOn exception/error
ABORTKilled externallyOn timeout/cancel
OTHERCustom stateCustom tooling

Every job must emit at least START + (COMPLETE | FAIL | ABORT).


Facet Reference

Job Facets

FacetKey FieldsPurpose
sqlqueryFull SQL text of the transformation
sourceCodeLocationtype, url, path, versionGit repo + file
jobTypejobType, processingType, integrationBATCH vs STREAMING
ownershipowners[].name, owners[].typeTeam/service owner

Run Facets

FacetKey FieldsPurpose
nominalTimenominalStartTime, nominalEndTimeLogical execution window
parentrun.runId, job.namespace, job.nameParent Airflow task → child Spark job
errorMessagemessage, programmingLanguage, stackTraceStructured error on FAIL
externalQueryexternalQueryId, sourceMaps to DW query ID

Dataset Facets

FacetKey FieldsPurpose
schemafields[].name, fields[].typeColumn definitions
columnLineagefields.{col}.inputFieldsColumn-level lineage
dataSourcename, uriConnection identifier
symlinksidentifiers[].name, .typeAlternate dataset names
dataQualityMetricscolumnMetrics.{col}.*, rowCountQuality stats
dataQualityAssertionsassertions[].success, .assertionGE/Soda results
lifecycleStateChangelifecycleStateChangeCREATE/DROP/OVERWRITE/RENAME
storagestorageLayer, fileFormatIceberg/Delta/Parquet
outputStatisticsrowCount, sizeWrite volume

Marquez Backend

Marquez is the OpenLineage reference implementation — stores and visualizes lineage events.

Docker Compose

version: "3.8"
services:
  marquez-db:
    image: postgres:14
    environment:
      POSTGRES_DB: marquez
      POSTGRES_USER: marquez
      POSTGRES_PASSWORD: marquez
    volumes:
      - marquez-db:/var/lib/postgresql/data

  marquez:
    image: marquezproject/marquez:0.47.0
    environment:
      MARQUEZ_PORT: 5000
      MARQUEZ_ADMIN_PORT: 5001
      MARQUEZ_DB_HOST: marquez-db
      MARQUEZ_DB_PORT: 5432
      MARQUEZ_DB_NAME: marquez
      MARQUEZ_DB_USER: marquez
      MARQUEZ_DB_PASSWORD: marquez
    ports:
      - "5000:5000"   # API
      - "5001:5001"   # Admin
    depends_on:
      - marquez-db

  marquez-web:
    image: marquezproject/marquez-web:0.47.0
    environment:
      MARQUEZ_HOST: marquez
      MARQUEZ_PORT: 5000
    ports:
      - "3000:3000"   # UI
    depends_on:
      - marquez

volumes:
  marquez-db:

Marquez API

# List namespaces
curl http://localhost:5000/api/v1/namespaces | jq

# List jobs in namespace
curl "http://localhost:5000/api/v1/namespaces/airflow/jobs" | jq

# Get dataset with lineage
curl "http://localhost:5000/api/v1/namespaces/spark://master:7077/datasets/silver.orders" | jq

# Get lineage graph (upstream + downstream, depth=2)
curl "http://localhost:5000/api/v1/lineage?nodeId=dataset:spark://master:7077:silver.orders&depth=2" | jq

# Search datasets
curl "http://localhost:5000/api/v1/search?q=orders&type=DATASET" | jq

Airflow Integration

Installation

pip install apache-airflow-providers-openlineage

Configuration (airflow.cfg or environment variables)

[openlineage]
transport = {"type": "http", "url": "http://marquez:5000", "endpoint": "api/v1/lineage"}
namespace = airflow
disabled = false
disabled_for_operators = airflow.operators.empty.EmptyOperator

Or via environment variables:

export OPENLINEAGE_URL=http://marquez:5000
export OPENLINEAGE_NAMESPACE=airflow
export AIRFLOW__OPENLINEAGE__TRANSPORT='{"type": "http", "url": "http://marquez:5000", "endpoint": "api/v1/lineage"}'

What Airflow Captures Automatically

SourceLineage Captured
SQLExecuteQueryOperatorInput/output tables (SQL parsing via openlineage-sql)
PythonOperatorJob START/COMPLETE events (no dataset lineage without manual instrumentation)
SparkSubmitOperatorParent facet linking to child Spark job run
BigQueryInsertJobOperatorInput/output tables + SQL
S3CopyObjectOperatorInput/output S3 dataset
ExternalTaskSensorUpstream job dependency
TriggerDagRunOperatorParent DAG → child DAG lineage

Custom Dataset Lineage in PythonOperator

from openlineage.client.run import Dataset
from openlineage.client.facet import SchemaDatasetFacet, SchemaField
from airflow.providers.openlineage.extractors.base import OperatorLineage

class MyPandasOperator(BaseOperator):
    def get_openlineage_facets_on_complete(self, ti):
        return OperatorLineage(
            inputs=[
                Dataset(
                    namespace="postgres://prod-db:5432",
                    name="public.customers",
                    facets={
                        "schema": SchemaDatasetFacet(
                            fields=[
                                SchemaField("id", "BIGINT"),
                                SchemaField("email", "VARCHAR"),
                            ]
                        )
                    },
                )
            ],
            outputs=[
                Dataset(
                    namespace="s3://data-lake",
                    name="silver/customers",
                )
            ],
        )

Spark Integration

Maven / pip

<!-- pom.xml or build.sbt for JVM Spark -->
<dependency>
  <groupId>io.openlineage</groupId>
  <artifactId>openlineage-spark_2.12</artifactId>
  <version>1.15.0</version>
</dependency>
# Download JAR for spark-submit
wget https://repo1.maven.org/maven2/io/openlineage/openlineage-spark_2.12/1.15.0/openlineage-spark_2.12-1.15.0.jar

spark-submit Configuration

spark-submit \
  --jars /opt/spark/jars/openlineage-spark_2.12-1.15.0.jar \
  --conf spark.extraListeners=io.openlineage.spark.agent.OpenLineageSparkListener \
  --conf spark.openlineage.transport.type=http \
  --conf spark.openlineage.transport.url=http://marquez:5000 \
  --conf spark.openlineage.transport.endpoint=/api/v1/lineage \
  --conf spark.openlineage.namespace=spark://spark-master:7077 \
  --conf spark.openlineage.parentJobNamespace=airflow \
  --conf spark.openlineage.parentJobName=etl_dag.run_spark_transform \
  --conf spark.openlineage.parentRunId=d46e465b-d358-4d32-83d4-df660ff614dd \
  my_etl_job.py

PySpark Programmatic Configuration

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("silver_transform") \
    .config("spark.extraListeners",
            "io.openlineage.spark.agent.OpenLineageSparkListener") \
    .config("spark.openlineage.transport.type", "http") \
    .config("spark.openlineage.transport.url", "http://marquez:5000") \
    .config("spark.openlineage.transport.endpoint", "/api/v1/lineage") \
    .config("spark.openlineage.namespace", "spark://spark-master:7077") \
    .getOrCreate()

# All subsequent spark.read / spark.write / df.createOrReplaceTempView
# and SQL queries are automatically tracked
df = spark.read.table("bronze.orders")
result = df.filter("status = 'valid'").select("id", "customer_id", "total")
result.writeTo("silver.orders").append()
# → Emits RunEvent with inputs=[bronze.orders], outputs=[silver.orders]

What Spark Captures

OperationLineage
spark.read.table(name)Input dataset
df.write.saveAsTable(name)Output dataset
spark.sql("INSERT INTO ...")Parsed SQL → input/output
df.join(other, ...)Both DFs as inputs
Column transformationsColumn-level lineage (when SQL-based)

dbt Integration

pip install openlineage-dbt

profiles.yml — no changes needed; OL is configured separately

# Run dbt with OpenLineage emission
export OPENLINEAGE_URL=http://marquez:5000
export OPENLINEAGE_NAMESPACE=dbt_prod
dbt run --target prod

# Or via dbt-openlineage package flags
dbt-ol run --target prod

dbt_project.yml metadata captured

# These model-level configs appear in lineage facets
models:
  my_project:
    staging:
      +meta:
        owner: "data-team"
        tags: ["daily", "customers"]

What dbt Captures

dbt ArtifactOpenLineage Event
dbt run (model)RunEvent per model: inputs = ref() / source(), output = model target
dbt testRunEvent per test: dataset = tested model
dbt snapshotRunEvent: input = source table, output = snapshot table
Column-level lineageParsed from model SQL via openlineage-sql

Custom Python Emitter

from openlineage.client import OpenLineageClient
from openlineage.client.run import (
    RunEvent, RunState, Run, Job,
    Dataset, InputDataset, OutputDataset,
)
from openlineage.client.facet import (
    SchemaDatasetFacet, SchemaField,
    SqlJobFacet, NominalTimeRunFacet,
    ColumnLineageDatasetFacet, ColumnLineageDatasetFacetFieldsAdditional,
    ColumnLineageDatasetFacetFieldsAdditionalInputFields,
    OutputStatisticsOutputDatasetFacet,
)
from openlineage.client.transport.http import HttpTransport, HttpConfig
import uuid
from datetime import datetime, timezone

client = OpenLineageClient(
    transport=HttpTransport(
        HttpConfig(url="http://marquez:5000", endpoint="api/v1/lineage")
    )
)

run_id = str(uuid.uuid4())
now = datetime.now(timezone.utc).isoformat()

# Emit START
client.emit(RunEvent(
    eventType=RunState.START,
    eventTime=now,
    run=Run(runId=run_id),
    job=Job(namespace="my-pipeline", name="daily_revenue"),
    producer="https://github.com/org/pipeline",
    inputs=[InputDataset(namespace="postgres://db:5432", name="public.orders")],
    outputs=[OutputDataset(namespace="s3://warehouse", name="gold/revenue")],
))

# ... do actual work ...

# Emit COMPLETE with column lineage
client.emit(RunEvent(
    eventType=RunState.COMPLETE,
    eventTime=datetime.now(timezone.utc).isoformat(),
    run=Run(runId=run_id),
    job=Job(
        namespace="my-pipeline",
        name="daily_revenue",
        facets={"sql": SqlJobFacet(query="SELECT order_date, SUM(total) AS revenue FROM orders GROUP BY 1")},
    ),
    producer="https://github.com/org/pipeline",
    inputs=[
        InputDataset(
            namespace="postgres://db:5432",
            name="public.orders",
            facets={
                "schema": SchemaDatasetFacet(fields=[
                    SchemaField("order_date", "DATE"),
                    SchemaField("total", "NUMERIC"),
                ])
            },
        )
    ],
    outputs=[
        OutputDataset(
            namespace="s3://warehouse",
            name="gold/revenue",
            facets={
                "schema": SchemaDatasetFacet(fields=[
                    SchemaField("order_date", "DATE"),
                    SchemaField("revenue", "NUMERIC"),
                ]),
                "columnLineage": ColumnLineageDatasetFacet(
                    fields={
                        "order_date": ColumnLineageDatasetFacetFieldsAdditional(
                            inputFields=[ColumnLineageDatasetFacetFieldsAdditionalInputFields(
                                namespace="postgres://db:5432",
                                name="public.orders",
                                field="order_date",
                            )]
                        ),
                        "revenue": ColumnLineageDatasetFacetFieldsAdditional(
                            inputFields=[ColumnLineageDatasetFacetFieldsAdditionalInputFields(
                                namespace="postgres://db:5432",
                                name="public.orders",
                                field="total",
                                transformationType="AGGREGATE",
                                transformationDescription="SUM",
                            )]
                        ),
                    }
                ),
                "outputStatistics": OutputStatisticsOutputDatasetFacet(
                    rowCount=365,
                    size=4096,
                ),
            },
        )
    ],
))

Namespace Conventions

Consistent namespaces are critical — if Airflow and Spark use different namespace strings for the same dataset, lineage will show two disconnected nodes.

SystemNamespace PatternExample
PostgreSQLpostgres://<host>:<port>postgres://prod-db:5432
MySQLmysql://<host>:<port>mysql://rds-mysql:3306
S3s3://<bucket>s3://data-lake
HDFShdfs://<host>:<port>hdfs://namenode:8020
Kafkakafka://<bootstrap>kafka://kafka:9092
Spark tablesspark://<master>spark://spark-master:7077
Trinotrino://<host>:<port>trino://trino:8080
dbtdbt://<project>dbt://my_project

Dataset names follow <schema>.<table> or <database>.<schema>.<table> conventions — match exactly what the SQL engine uses.


Transport Types

TransportConfigUse Case
httpurl, endpoint, authMarquez, DataHub, OpenMetadata
filelog_file_path, appendLocal debugging, batch import
consoleDevelopment (prints to stdout)
kafkatopic, bootstrap.serversHigh-throughput / decoupled
compositetransports: [t1, t2]Fan-out to multiple backends
# Composite transport example: HTTP + file
from openlineage.client.transport.composite import CompositeTransport, CompositeConfig

transport = CompositeTransport(CompositeConfig(transports=[
    {"type": "http", "url": "http://marquez:5000", "endpoint": "api/v1/lineage"},
    {"type": "file", "log_file_path": "/var/log/openlineage.jsonl"},
]))

Impact Analysis Pattern

Once lineage is in Marquez, use the API to answer "what breaks if I drop column X?":

import requests

def get_downstream_jobs(namespace: str, dataset_name: str, depth: int = 5) -> list[dict]:
    """Return all jobs that read from the given dataset (transitively)."""
    resp = requests.get(
        "http://marquez:5000/api/v1/lineage",
        params={
            "nodeId": f"dataset:{namespace}:{dataset_name}",
            "depth": depth,
        },
    )
    graph = resp.json()
    return [
        node for node in graph["graph"]
        if node["type"] == "JOB" and node["id"] != f"dataset:{namespace}:{dataset_name}"
    ]

# Example: find all downstream jobs of silver.orders
downstream = get_downstream_jobs(
    namespace="spark://spark-master:7077",
    dataset_name="silver.orders",
)
for job in downstream:
    print(job["id"], job["data"]["latestRun"]["state"])

Anti-Patterns

  1. Different namespace strings for the same database — Airflow uses postgres://host:5432, Spark uses postgresql — graph shows two unconnected clusters. Standardize namespaces across all tools.

  2. Emitting only COMPLETE events without START — many backends require START to create the run record; COMPLETE alone is silently dropped or creates orphan records.

  3. Not linking Spark/dbt child runs to parent Airflow run — lineage is fragmented per system. Always pass parentJobNamespace, parentJobName, parentRunId to Spark and dbt.

  4. Omitting nominalTime facet on scheduled jobs — makes time-partitioned lineage unusable. Use the Airflow logical date as nominalStartTime.

  5. Using file transport in production — JSON files grow unboundedly and are not queryable. Use HTTP transport to Marquez.

  6. Relying solely on job-level lineage — table-level lineage is insufficient for column impact analysis. Instrument SQL-based jobs with openlineage-sql to get column lineage.

  7. Custom dataset names that diverge from SQL names — e.g., calling a dataset "orders_v2" while SQL uses silver.orders creates phantom nodes. Use the actual schema.table name.

  8. Not handling FAIL events — a job that silently dies without emitting FAIL leaves runs in START state forever. Always wrap job logic in try/except and emit FAIL on error.


References to Consult When Needed

  • OpenLineage spec: openlineage.io/spec
  • Marquez project: marquezproject.ai
  • Airflow provider: airflow.apache.org/docs/apache-airflow-providers-openlineage/
  • Spark integration: openlineage.io/docs/integrations/spark
  • dbt integration: openlineage.io/docs/integrations/dbt
  • openlineage-sql (SQL parser): github.com/OpenLineage/OpenLineage/tree/main/integration/sql

Keep looking

Skills are one crate of 328,083. Ordering is by how many stacks a row turns up in, so the top of any crate is what has actually been picked rather than what has the most stars.