Skip to main content

DataFlow & DataJob

Why Would You Use DataFlow and DataJob?​

The DataFlow and DataJob entities are used to represent data processing pipelines and jobs within a data ecosystem. They allow users to define, manage, and monitor the flow of data through various stages of processing, from ingestion to transformation and storage.

Goal Of This Guide​

This guide will show you how to

  • Create a DataFlow.
  • Create a Datajob with a DataFlow.

Prerequisites​

For this tutorial, you need to deploy DataHub Quickstart and ingest sample data. For detailed steps, please refer to DataHub Quickstart Guide.

Create DataFlow​

# Inlined from /metadata-ingestion/examples/library/dataflow_create.py
from datahub.metadata.urns import TagUrn
from datahub.sdk import DataFlow, DataHubClient

client = DataHubClient.from_env()

dataflow = DataFlow(
name="example_dataflow",
platform="airflow",
description="airflow pipeline for production",
tags=[TagUrn(name="production"), TagUrn(name="data_engineering")],
)

client.entities.upsert(dataflow)

Create DataJob​

DataJob must be associated with a DataFlow. You can create a DataJob by providing the DataFlow object or the DataFlow URN and its platform instance.

# Inlined from /metadata-ingestion/examples/library/datajob_create_full.py
from datahub.metadata.urns import DatasetUrn, TagUrn
from datahub.sdk import DataFlow, DataHubClient, DataJob

client = DataHubClient.from_env()

# datajob will inherit the platform and platform instance from the flow

dataflow = DataFlow(
platform="airflow",
name="example_dag",
platform_instance="PROD",
description="example dataflow",
tags=[TagUrn(name="tag1"), TagUrn(name="tag2")],
)

datajob = DataJob(
name="example_datajob",
flow=dataflow,
inlets=[
DatasetUrn(platform="hdfs", name="dataset1", env="PROD"),
],
outlets=[
DatasetUrn(platform="hdfs", name="dataset2", env="PROD"),
],
)

client.entities.upsert(dataflow)
client.entities.upsert(datajob)

Read DataFlow​

# Inlined from /metadata-ingestion/examples/library/dataflow_read.py
from datahub.sdk import DataFlowUrn, DataHubClient

client = DataHubClient.from_env()

# Or get this from the UI (share -> copy urn) and use DataFlowUrn.from_string(...)
dataflow_urn = DataFlowUrn(
orchestrator="airflow", flow_id="example_dataflow", cluster="PROD"
)

dataflow_entity = client.entities.get(dataflow_urn)
print("DataFlow name:", dataflow_entity.name)
print("DataFlow platform:", dataflow_entity.platform)
print("DataFlow description:", dataflow_entity.description)

Example Output​

>> DataFlow name: example_dataflow
>> DataFlow platform: urn:li:dataPlatform:airflow
>> DataFlow description: airflow pipeline for production

Read DataJob​

# Inlined from /metadata-ingestion/examples/library/datajob_read.py
from datahub.sdk import DataFlowUrn, DataHubClient, DataJobUrn

client = DataHubClient.from_env()

# Or get this from the UI (share -> copy urn) and use DataJobUrn.from_string(...)
# The flow_id carries the flow's platform instance as a prefix ("PROD."), which is
# how DataFlow(platform_instance=...) builds its urn.
datajob_urn = DataJobUrn(
flow=DataFlowUrn(
orchestrator="airflow", flow_id="PROD.example_dag", cluster="PROD"
),
job_id="example_datajob",
)

datajob_entity = client.entities.get(datajob_urn)
print("DataJob name:", datajob_entity.name)
print("DataJob Flow URN:", datajob_entity.flow_urn)
print("DataJob description:", datajob_entity.description)

Example Output​

>> DataJob name: example_datajob
>> DataJob Flow URN: urn:li:dataFlow:(airflow,PROD.example_dag,PROD)
>> DataJob description: example datajob