Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
58 changes: 47 additions & 11 deletions airflow/providers/openlineage/plugins/adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -164,17 +164,19 @@ def start_task(

if not run_facets:
run_facets = {}
if task:
run_facets = {**task.run_facets, **run_facets}
run_facets["processing_engine"] = processing_engine_version_facet # type: ignore
event = RunEvent(
eventType=RunState.START,
eventTime=event_time,
run=self._build_run(
run_id,
job_name,
parent_job_name,
parent_run_id,
nominal_start_time,
nominal_end_time,
run_id=run_id,
job_name=job_name,
parent_job_name=parent_job_name,
parent_run_id=parent_run_id,
nominal_start_time=nominal_start_time,
nominal_end_time=nominal_end_time,
run_facets=run_facets,
),
job=self._build_job(
Expand All @@ -190,40 +192,74 @@ def start_task(
)
self.emit(event)

def complete_task(self, run_id: str, job_name: str, end_time: str, task: OperatorLineage):
def complete_task(
self,
run_id: str,
job_name: str,
parent_job_name: str | None,
parent_run_id: str | None,
end_time: str,
task: OperatorLineage,
):
"""
Emits openlineage event of type COMPLETE.

:param run_id: globally unique identifier of task in dag run
:param job_name: globally unique identifier of task between dags
:param parent_job_name: the name of the parent job (typically the DAG,
but possibly a task group)
:param parent_run_id: identifier of job spawning this task
:param end_time: time of task completion
:param task: metadata container with information extracted from operator
"""
event = RunEvent(
eventType=RunState.COMPLETE,
eventTime=end_time,
run=self._build_run(run_id, job_name=job_name, run_facets=task.run_facets),
run=self._build_run(
run_id=run_id,
job_name=job_name,
parent_job_name=parent_job_name,
parent_run_id=parent_run_id,
run_facets=task.run_facets,
),
job=self._build_job(job_name, job_facets=task.job_facets),
inputs=task.inputs,
outputs=task.outputs,
producer=_PRODUCER,
)
self.emit(event)

def fail_task(self, run_id: str, job_name: str, end_time: str, task: OperatorLineage):
def fail_task(
self,
run_id: str,
job_name: str,
parent_job_name: str | None,
parent_run_id: str | None,
end_time: str,
task: OperatorLineage,
):
"""
Emits openlineage event of type FAIL.

:param run_id: globally unique identifier of task in dag run
:param job_name: globally unique identifier of task between dags
:param parent_job_name: the name of the parent job (typically the DAG,
but possibly a task group)
:param parent_run_id: identifier of job spawning this task
:param end_time: time of task completion
:param task: metadata container with information extracted from operator
"""
event = RunEvent(
eventType=RunState.FAIL,
eventTime=end_time,
run=self._build_run(run_id, job_name=job_name, run_facets=task.run_facets),
job=self._build_job(job_name),
run=self._build_run(
run_id=run_id,
job_name=job_name,
parent_job_name=parent_job_name,
parent_run_id=parent_run_id,
run_facets=task.run_facets,
),
job=self._build_job(job_name, job_facets=task.job_facets),
inputs=task.inputs,
outputs=task.outputs,
producer=_PRODUCER,
Expand Down
11 changes: 10 additions & 1 deletion airflow/providers/openlineage/plugins/listener.py
Original file line number Diff line number Diff line change
Expand Up @@ -101,7 +101,6 @@ def on_running():
owners=dag.owner.split(", "),
task=task_metadata,
run_facets={
**task_metadata.run_facets,
**get_custom_facets(task_instance),
**get_airflow_run_facet(dagrun, dag, task_instance, task, task_uuid),
},
Expand All @@ -115,13 +114,16 @@ def on_task_instance_success(self, previous_state, task_instance: TaskInstance,

dagrun = task_instance.dag_run
task = task_instance.task
dag = task.dag

task_uuid = OpenLineageAdapter.build_task_instance_run_id(
task.task_id, task_instance.execution_date, task_instance.try_number - 1
)

@print_warning(self.log)
def on_success():
parent_run_id = OpenLineageAdapter.build_dag_run_id(dag.dag_id, dagrun.run_id)

task_metadata = self.extractor_manager.extract_metadata(
dagrun, task, complete=True, task_instance=task_instance
)
Expand All @@ -131,6 +133,8 @@ def on_success():
self.adapter.complete_task(
run_id=task_uuid,
job_name=get_job_name(task),
parent_job_name=dag.dag_id,
parent_run_id=parent_run_id,
end_time=end_date.isoformat(),
task=task_metadata,
)
Expand All @@ -143,13 +147,16 @@ def on_task_instance_failed(self, previous_state, task_instance: TaskInstance, s

dagrun = task_instance.dag_run
task = task_instance.task
dag = task.dag

task_uuid = OpenLineageAdapter.build_task_instance_run_id(
task.task_id, task_instance.execution_date, task_instance.try_number
)

@print_warning(self.log)
def on_failure():
parent_run_id = OpenLineageAdapter.build_dag_run_id(dag.dag_id, dagrun.run_id)

task_metadata = self.extractor_manager.extract_metadata(
dagrun, task, complete=True, task_instance=task_instance
)
Expand All @@ -159,6 +166,8 @@ def on_failure():
self.adapter.fail_task(
run_id=task_uuid,
job_name=get_job_name(task),
parent_job_name=dag.dag_id,
parent_run_id=parent_run_id,
end_time=end_date.isoformat(),
task=task_metadata,
)
Expand Down
9 changes: 8 additions & 1 deletion tests/providers/openlineage/plugins/test_listener.py
Original file line number Diff line number Diff line change
Expand Up @@ -223,7 +223,6 @@ def test_adapter_start_task_is_called_with_proper_arguments(
owners=["Test Owner"],
task=listener.extractor_manager.extract_metadata(),
run_facets={
"run_facet": 1,
"custom_facet": 2,
"airflow_run_facet": 3,
},
Expand All @@ -243,11 +242,14 @@ def test_adapter_fail_task_is_called_with_proper_arguments(mock_get_job_name, mo
listener, task_instance = _create_listener_and_task_instance()
mock_get_job_name.return_value = "job_name"
mocked_adapter.build_task_instance_run_id.side_effect = lambda x, y, z: f"{x}.{y}.{z}"
mocked_adapter.build_dag_run_id.side_effect = lambda x, y: f"{x}.{y}"

listener.on_task_instance_failed(None, task_instance, None)
listener.adapter.fail_task.assert_called_once_with(
end_time="2023-01-03T13:01:01",
job_name="job_name",
parent_job_name="dag_id",
parent_run_id="dag_id.dag_run_run_id",
run_id="task_id.execution_date.1",
task=listener.extractor_manager.extract_metadata(),
)
Expand All @@ -267,13 +269,16 @@ def test_adapter_complete_task_is_called_with_proper_arguments(mock_get_job_name
listener, task_instance = _create_listener_and_task_instance()
mock_get_job_name.return_value = "job_name"
mocked_adapter.build_task_instance_run_id.side_effect = lambda x, y, z: f"{x}.{y}.{z}"
mocked_adapter.build_dag_run_id.side_effect = lambda x, y: f"{x}.{y}"

listener.on_task_instance_success(None, task_instance, None)
# This run_id will be different as we did NOT simulate increase of the try_number attribute,
# which happens in Airflow.
listener.adapter.complete_task.assert_called_once_with(
end_time="2023-01-03T13:01:01",
job_name="job_name",
parent_job_name="dag_id",
parent_run_id="dag_id.dag_run_run_id",
run_id="task_id.execution_date.0",
task=listener.extractor_manager.extract_metadata(),
)
Expand All @@ -285,6 +290,8 @@ def test_adapter_complete_task_is_called_with_proper_arguments(mock_get_job_name
listener.adapter.complete_task.assert_called_once_with(
end_time="2023-01-03T13:01:01",
job_name="job_name",
parent_job_name="dag_id",
parent_run_id="dag_id.dag_run_run_id",
run_id="task_id.execution_date.1",
task=listener.extractor_manager.extract_metadata(),
)
Expand Down
Loading