OpenLineage and Marquez: A Hands-On Tutorial
Chat2DB TeamMost lineage tools want to be your catalog, your quality platform and your governance layer. OpenLineage wants to be none of those. It is an open specification for how a job reports what it read and wrote, and nothing else. That narrow scope is why it has ended up integrated into Airflow, dbt, Spark, Flink, Dagster and Great Expectations, and why it is the sensible foundation to build on rather than a vendor format you will regret.
Marquez is the reference implementation that receives those events, stores them in PostgreSQL, and serves a graph API and a UI. Together they get you working lineage in an afternoon.
This tutorial runs both locally, emits events by hand so you understand the model, then wires up Airflow and dbt.
The OpenLineage model
Four concepts, and they compose in a way that stays sane at scale.
- Job — a recurring process.
etl.load_orders. It exists independently of any execution. - Run — one execution of a job, identified by a UUID. It has a lifecycle:
START, thenRUNNING, thenCOMPLETE,FAILorABORT. - Dataset — a table, file or topic, identified by a namespace and a name.
postgres://db:5432plusanalytics.fct_orders. - Facet — a piece of typed metadata attached to any of the above. Schema, column lineage, row counts, SQL text, data quality assertions, ownership.
Facets are the extension mechanism. The core event is deliberately tiny; everything interesting arrives as a facet, and you can define your own without breaking consumers that do not understand it.
An event is a JSON document posted to /api/v1/lineage:
{
"eventType": "START",
"eventTime": "2026-09-08T09:00:00.000Z",
"producer": "https://github.com/example/etl/blob/main/load_orders.py",
"run": { "runId": "d46e465b-d358-4d32-83d4-df660ff614dd" },
"job": { "namespace": "etl", "name": "load_orders" },
"inputs": [
{ "namespace": "postgres://db:5432", "name": "public.raw_orders" }
],
"outputs": [
{ "namespace": "postgres://db:5432", "name": "analytics.fct_orders" }
]
}That is the whole idea. A job says "I am starting, here is what I read and what I will write." Later it says COMPLETE, optionally attaching row counts and schemas.
Run Marquez locally
Marquez ships a Docker Compose setup. The only reliable requirement is Docker itself.
git clone https://github.com/MarquezProject/marquez.git
cd marquez
./docker/up.sh --seedThat starts four things: the API on port 5000, an admin port on 5001, the web UI on 3000, and a PostgreSQL instance holding the metadata. The --seed flag loads example data so the UI is not empty on first open.
If you prefer an explicit compose file:
services:
db:
image: postgres:16
environment:
POSTGRES_USER: marquez
POSTGRES_PASSWORD: marquez
POSTGRES_DB: marquez
volumes:
- marquez-db:/var/lib/postgresql/data
api:
image: marquezproject/marquez:latest
ports: ["5000:5000", "5001:5001"]
depends_on: [db]
environment:
MARQUEZ_CONFIG: /opt/marquez/marquez.yml
POSTGRES_HOST: db
POSTGRES_PORT: "5432"
POSTGRES_DB: marquez
POSTGRES_USER: marquez
POSTGRES_PASSWORD: marquez
web:
image: marquezproject/marquez-web:latest
ports: ["3000:3000"]
depends_on: [api]
environment:
MARQUEZ_HOST: api
MARQUEZ_PORT: "5000"
volumes:
marquez-db:Check it is alive:
curl -s http://localhost:5001/healthcheck | python3 -m json.toolThen open http://localhost:3000 (opens in a new tab).
Emit your first event by hand
Understanding the raw event is worth ten minutes, because when an integration misbehaves you will be reading these.
RUN_ID=$(python3 -c "import uuid; print(uuid.uuid4())")
curl -s -X POST http://localhost:5000/api/v1/lineage \
-H 'Content-Type: application/json' \
-d "{
\"eventType\": \"START\",
\"eventTime\": \"$(date -u +%Y-%m-%dT%H:%M:%S.000Z)\",
\"producer\": \"https://github.com/example/etl\",
\"run\": { \"runId\": \"$RUN_ID\" },
\"job\": { \"namespace\": \"tutorial\", \"name\": \"load_orders\" },
\"inputs\": [
{ \"namespace\": \"postgres://localhost:5432\", \"name\": \"public.raw_orders\" }
],
\"outputs\": [
{ \"namespace\": \"postgres://localhost:5432\", \"name\": \"analytics.fct_orders\" }
]
}"A 201 with an empty body means accepted. Now complete the run, and attach facets this time:
curl -s -X POST http://localhost:5000/api/v1/lineage \
-H 'Content-Type: application/json' \
-d "{
\"eventType\": \"COMPLETE\",
\"eventTime\": \"$(date -u +%Y-%m-%dT%H:%M:%S.000Z)\",
\"producer\": \"https://github.com/example/etl\",
\"run\": { \"runId\": \"$RUN_ID\" },
\"job\": { \"namespace\": \"tutorial\", \"name\": \"load_orders\" },
\"outputs\": [{
\"namespace\": \"postgres://localhost:5432\",
\"name\": \"analytics.fct_orders\",
\"facets\": {
\"schema\": {
\"_producer\": \"https://github.com/example/etl\",
\"_schemaURL\": \"https://openlineage.io/spec/facets/1-0-0/SchemaDatasetFacet.json\",
\"fields\": [
{ \"name\": \"order_id\", \"type\": \"BIGINT\" },
{ \"name\": \"customer_id\", \"type\": \"BIGINT\" },
{ \"name\": \"total_cents\", \"type\": \"BIGINT\" }
]
},
\"outputStatistics\": {
\"_producer\": \"https://github.com/example/etl\",
\"_schemaURL\": \"https://openlineage.io/spec/facets/1-0-0/OutputStatisticsOutputDatasetFacet.json\",
\"rowCount\": 128394,
\"size\": 9128374
}
}
}]
}"Read it back:
curl -s "http://localhost:5000/api/v1/namespaces/tutorial/jobs/load_orders" | python3 -m json.tool
curl -s "http://localhost:5000/api/v1/lineage?nodeId=dataset:postgres://localhost:5432:analytics.fct_orders&depth=3" | python3 -m json.toolThat second call is the one that matters — it returns the graph around a node, which is what the UI renders and what you would query from your own tooling.
Emit events from Python
Hand-rolling JSON is fine for learning and terrible for production. The client library handles the boilerplate:
pip install openlineage-pythonimport uuid
from datetime import datetime, timezone
from openlineage.client import OpenLineageClient
from openlineage.client.event_v2 import Job, Run, RunEvent, RunState
from openlineage.client.facet_v2 import schema_dataset, sql_job
from openlineage.client.event_v2 import InputDataset, OutputDataset
client = OpenLineageClient(url="http://localhost:5000")
NAMESPACE = "postgres://localhost:5432"
PRODUCER = "https://github.com/example/etl/blob/main/load_orders.py"
QUERY = """
INSERT INTO analytics.fct_orders (order_id, customer_id, total_cents)
SELECT o.id, o.customer_id, o.total_cents
FROM public.raw_orders o
WHERE o.created_at >= current_date - 1
"""
run_id = str(uuid.uuid4())
job = Job(namespace="etl", name="load_orders",
facets={"sql": sql_job.SQLJobFacet(query=QUERY)})
run = Run(runId=run_id)
def now():
return datetime.now(timezone.utc).isoformat()
# START
client.emit(RunEvent(
eventType=RunState.START,
eventTime=now(),
run=run, job=job, producer=PRODUCER,
inputs=[InputDataset(namespace=NAMESPACE, name="public.raw_orders")],
outputs=[OutputDataset(namespace=NAMESPACE, name="analytics.fct_orders")],
))
try:
rows = run_the_actual_etl() # your real work here
client.emit(RunEvent(
eventType=RunState.COMPLETE,
eventTime=now(),
run=run, job=job, producer=PRODUCER,
inputs=[],
outputs=[OutputDataset(
namespace=NAMESPACE,
name="analytics.fct_orders",
facets={"schema": schema_dataset.SchemaDatasetFacet(fields=[
schema_dataset.SchemaDatasetFacetFields(name="order_id", type="BIGINT"),
schema_dataset.SchemaDatasetFacetFields(name="customer_id", type="BIGINT"),
schema_dataset.SchemaDatasetFacetFields(name="total_cents", type="BIGINT"),
])},
)],
))
except Exception:
client.emit(RunEvent(
eventType=RunState.FAIL,
eventTime=now(),
run=run, job=job, producer=PRODUCER,
inputs=[], outputs=[],
))
raiseTwo things worth copying from that structure. First, always emit FAIL in an exception handler — a run that starts and never finishes shows as perpetually running, and after a week of those the UI is unreadable. Second, put the SQL in a SQLJobFacet; Marquez displays it, and it turns the lineage graph into something you can debug from directly.
Column-level lineage
Column lineage arrives as a facet on the output dataset, mapping each output field to its input fields:
from openlineage.client.facet_v2 import column_lineage_dataset as cl
column_lineage = cl.ColumnLineageDatasetFacet(fields={
"lifetime_value_cents": cl.Fields(
inputFields=[
cl.InputField(
namespace=NAMESPACE,
name="public.fct_orders",
field="total_cents",
transformations=[cl.Transformation(
type="AGGREGATION", subtype="SUM", masking=False
)],
)
],
transformationDescription="SUM of order totals per customer",
transformationType="AGGREGATION",
),
"email": cl.Fields(
inputFields=[cl.InputField(
namespace=NAMESPACE, name="public.raw_users", field="email",
)],
transformationType="IDENTITY",
),
})The masking flag on a transformation is genuinely useful and underused: it marks that a column was hashed or redacted on the way through, which is what lets you answer "does any downstream consumer see raw PII?" from the lineage graph rather than by reading every model.
You do not have to build these by hand. openlineage-sql parses SQL and derives column lineage, and the Airflow and dbt integrations use it automatically.
Wire up Airflow
Airflow 2.7 and later ship an OpenLineage provider that instruments operators without changing DAG code.
pip install apache-airflow-providers-openlineage# airflow.cfg
[openlineage]
transport = {"type": "http", "url": "http://localhost:5000"}
namespace = airflow-prodOr by environment variable, which is easier in containers:
export AIRFLOW__OPENLINEAGE__TRANSPORT='{"type":"http","url":"http://localhost:5000"}'
export AIRFLOW__OPENLINEAGE__NAMESPACE='airflow-prod'Operators with built-in extractors — PostgresOperator, SnowflakeOperator, BigQueryInsertJobOperator, SparkSubmitOperator, the dbt operators — now emit lineage automatically, parsing their SQL to determine inputs and outputs. A DAG like this needs no lineage-specific code at all:
from airflow import DAG
from airflow.providers.postgres.operators.postgres import PostgresOperator
import pendulum
with DAG(
"orders_pipeline",
start_date=pendulum.datetime(2026, 9, 1, tz="UTC"),
schedule="@daily",
catchup=False,
) as dag:
stage = PostgresOperator(
task_id="stage_orders",
postgres_conn_id="warehouse",
sql="""
INSERT INTO staging.stg_orders
SELECT id, customer_id, total_cents, created_at
FROM public.raw_orders
WHERE created_at >= current_date - 1
""",
)
enrich = PostgresOperator(
task_id="enrich_orders",
postgres_conn_id="warehouse",
sql="""
INSERT INTO analytics.fct_orders
SELECT s.id, s.customer_id,
s.total_cents - coalesce(r.amount_cents, 0)
FROM staging.stg_orders s
LEFT JOIN public.refunds r ON r.order_id = s.id
""",
)
stage >> enrichRun it once and the graph appears in Marquez, including the raw_orders → stg_orders → fct_orders chain derived from the SQL rather than from the task dependencies. That distinction matters: task order tells you what ran when, SQL parsing tells you what actually depends on what, and they are not always the same.
For a PythonOperator doing something the extractor cannot see, emit manually inside the callable using the client code from earlier, reusing Airflow's run id so the events attach to the right run.
Wire up dbt
dbt does not emit OpenLineage natively; a wrapper reads the run artifacts and converts them.
pip install openlineage-dbt
export OPENLINEAGE_URL=http://localhost:5000
export OPENLINEAGE_NAMESPACE=dbt-prod
dbt-ol run --project-dir . --profiles-dir .
dbt-ol test --project-dir . --profiles-dir .dbt-ol runs dbt normally, then parses target/manifest.json and target/run_results.json and posts equivalent OpenLineage events. You get model-level lineage from ref() dependencies, execution status per model, and — with recent dbt versions and a supported adapter — column-level lineage from the compiled SQL.
Running dbt-ol test as well is worth it: test results arrive as data quality facets attached to the datasets, so a failing not_null test shows up on the node in the lineage graph rather than only in CI logs.
Query the graph from your own tools
The UI is useful for humans; the API is what you build automation on.
What does this dataset depend on, three hops out?
curl -s "http://localhost:5000/api/v1/lineage?nodeId=dataset:postgres://localhost:5432:analytics.fct_orders&depth=3"List every dataset in a namespace:
curl -s "http://localhost:5000/api/v1/namespaces/postgres%3A%2F%2Flocalhost%3A5432/datasets?limit=100"Recent runs of a job, to see failures:
curl -s "http://localhost:5000/api/v1/namespaces/etl/jobs/load_orders/runs?limit=20"The obvious automation is a CI check: when a pull request modifies a model, call the lineage endpoint for that dataset, collect downstream nodes, and post them as a review comment. That single integration prevents most "we dropped a column and broke finance" incidents, and it takes about fifty lines.
Inspecting Marquez's own database
Marquez stores everything in PostgreSQL, and querying it directly is often faster than the API when you are exploring:
docker exec -it marquez-db-1 psql -U marquez -d marquez-- Datasets and their current schema version
SELECT d.name, d.namespace_name, dv.created_at
FROM datasets d
JOIN dataset_versions dv ON dv.dataset_uuid = d.uuid
ORDER BY dv.created_at DESC
LIMIT 20;
-- Job runs and outcomes over the last day
SELECT j.name, r.current_run_state, r.started_at, r.ended_at,
r.ended_at - r.started_at AS duration
FROM runs r
JOIN jobs j ON j.uuid = r.job_uuid
WHERE r.started_at > now() - interval '1 day'
ORDER BY r.started_at DESC;
-- Datasets nothing has read in 30 days: deletion candidates
SELECT d.namespace_name, d.name, max(le.event_time) AS last_seen
FROM datasets d
LEFT JOIN lineage_events le ON le.event::text LIKE '%' || d.name || '%'
GROUP BY 1, 2
HAVING max(le.event_time) < now() - interval '30 days'
OR max(le.event_time) IS NULL;That last query is rough — string matching against the raw event JSON — but it is the kind of thing worth running once a quarter, and it usually finds a surprising amount of dead weight.
Poking at an unfamiliar schema like Marquez's goes faster with a client that renders the foreign keys as a diagram. Chat2DB (opens in a new tab) is a free AI-powered database client that generates ER diagrams from a live connection and translates plain-English questions into SQL, which helps when you are reverse-engineering a metadata store you did not design. It also runs in the browser at app.chat2db.ai (opens in a new tab).
Things to get right in production
- Never let lineage emission fail the pipeline. Configure a short timeout and swallow transport errors. Losing a lineage event is an inconvenience; failing a production job because a metadata service was restarting is an outage.
- Buffer through Kafka at volume. The HTTP transport is fine to start; past a few thousand events a minute, use the Kafka transport so Marquez restarts do not drop events.
- Use stable dataset naming. OpenLineage has naming conventions per source type — follow them, because
postgres://host:5432andpostgres://hostproduce two disconnected nodes for the same table, and you will spend an afternoon confused. - Always emit a terminal event.
COMPLETE,FAILorABORT. OrphanedSTARTevents accumulate fast. - Prune old events.
lineage_eventsgrows without bound. Marquez has a retention job; schedule it before the table reaches a hundred million rows. - Treat Marquez as replaceable. Because OpenLineage is a specification, moving to DataHub, OpenMetadata, Atlan or a managed service later means changing a transport URL, not re-instrumenting every job. That optionality is the main argument for adopting the standard rather than a proprietary agent.
Summary
OpenLineage is a small, well-scoped spec: jobs, runs, datasets, facets. Marquez is a working reference implementation you can run with one Docker command. Between them you can go from no lineage at all to an automatically-maintained graph covering Airflow, dbt and Spark in an afternoon — and because the format is open, the instrumentation you write today survives whatever catalog you standardise on next year.
Start by running Marquez locally and pointing one real Airflow DAG at it. Seeing your own pipeline appear as a graph is more persuasive than any amount of planning.
