Data lineage is crucial for understanding your data's journey, from its origin to its current state, including all transformations and dependencies. OpenLineage is an open standard that provides a common language and specification for collecting and exchanging this lineage metadata. Think of it as a universal protocol that allows different data systems (like Spark, Airflow, dbt) to describe their operations and the datasets they interact with in a consistent, interoperable way. This standardization is key for building a comprehensive view of your data pipelines without being locked into a proprietary format.
Marquez is an open-source metadata service that acts as a central repository for your OpenLineage events. When your data jobs (e.g., an Apache Spark job or an Airflow DAG) run, they can emit OpenLineage events describing their inputs, outputs, schemas, and processing logic. Marquez ingests these events, stores them, and uses them to construct a detailed, graph-based view of your data lineage. It provides APIs for programmatic access and a user interface to visualize your data pipelines, helping you trace data provenance, understand impact analysis, and improve data governance by seeing exactly how data flows and transforms across your entire ecosystem.
While OpenLineage and Marquez focus on building and storing the lineage graph, Monte Carlo is a commercial end-to-end data observability platform that leverages this understanding to ensure data quality and reliability. Monte Carlo monitors your data warehouses, lakes, and other data sources for freshness, volume anomalies, schema changes, and other quality issues. By understanding data dependencies (often through its own lineage capabilities or integrations with systems like Marquez), it can pinpoint the root cause of data problems, understand their blast radius, and proactively alert data engineers. It moves beyond just seeing lineage to actively monitoring, detecting, and helping resolve data incidents across your stack.
Key Takeaways
- OpenLineage is an open standard for consistent data lineage metadata collection.
- Marquez is an open-source metadata service that stores and visualizes OpenLineage events, acting as your central lineage hub.
- Monte Carlo is a commercial data observability platform focused on data quality and incident management, often leveraging lineage to detect and alert on data issues.
- Together, these tools help you build a clear lineage graph, monitor data health, and react quickly to data quality problems.
Code Example
from openlineage.client import OpenLineageClient, set_producer
from openlineage.client.utils import get_hostname
from openlineage.client.run import Dataset, Job, Run, RunEvent, RunState
# Configure OpenLineage client to send events to Marquez
set_producer("my_simple_etl", get_hostname())
client = OpenLineageClient(url="http://marquez:5000/api/v1", timeout=5.0)
# Define job, run, and datasets
job = Job("my_namespace", "users_transform_job")
run = Run(runId="unique-run-id-456")
input_dataset = Dataset("my_source_db", "public.raw_users")
output_dataset = Dataset("my_warehouse", "public.clean_users")
# Emit START event before processing
client.emit(RunEvent(RunState.START, run, job, [], [input_dataset]))
# Simulate ETL processing here
print("Processing data from raw_users to clean_users...")
# Emit COMPLETE event after successful processing
client.emit(RunEvent(RunState.COMPLETE, run, job, [output_dataset], [input_dataset]))
print("OpenLineage events emitted to Marquez.")How this code works
This code demonstrates how to use OpenLineage to track a simple data transformation job, specifically an ETL (Extract, Transform, Load) process that moves data from raw_users to clean_users. It configures an OpenLineageClient to send events to a Marquez server, identified by its URL. The set_producer function helps identify the source system sending these events, in this case, a 'my_simple_etl' process running on the current get_hostname.
The code defines the job name, a unique run identifier for this specific execution, and describes the input_dataset (public.raw_users) and output_dataset (public.clean_users). It then uses client.emit to send a RunEvent with RunState.START before the simulated data processing begins, signaling the job's initiation and identifying its inputs. After the processing, another client.emit sends a RunEvent with RunState.COMPLETE. A subtle but important detail is including the input_dataset again in the COMPLETE event, along with the output_dataset; this ensures Marquez fully understands the entire lineage, linking both sources and destinations to the completed job. This way, Marquez builds a comprehensive view of how clean_users was derived from raw_users.