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
17 changes: 17 additions & 0 deletions providers/amazon/docs/operators/emr/emr.rst
Original file line number Diff line number Diff line change
Expand Up @@ -118,6 +118,23 @@ available in your deployment.
:start-after: [START howto_operator_emr_add_steps]
:end-before: [END howto_operator_emr_add_steps]

OpenLineage parent job information
""""""""""""""""""""""""""""""""""

For Spark steps launched through ``command-runner.jar`` with ``spark-submit`` or ``run-example``,
:class:`~airflow.providers.amazon.aws.operators.emr.EmrAddStepsOperator` can inject OpenLineage parent job
information into the Spark arguments. This links OpenLineage events emitted by the Spark application to the
Airflow task that submitted the step.

Enable injection globally with the ``[openlineage] spark_inject_parent_job_info`` configuration option, or for
one operator by passing ``openlineage_inject_parent_job_info=True``. The operator preserves manually configured
``spark.openlineage.parent*`` properties and does not modify non-Spark steps. The Spark application must still
have the OpenLineage Spark integration installed and enabled.

See the `OpenLineage Spark automatic injection documentation
<https://airflow.apache.org/docs/apache-airflow-providers-openlineage/stable/spark.html#automatic-injection>`__
for configuration details and an example.

.. _howto/operator:EmrTerminateJobFlowOperator:

Terminate an EMR job flow
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
from __future__ import annotations

import ast
import copy
import warnings
from collections.abc import Sequence
from datetime import timedelta
Expand Down Expand Up @@ -60,6 +61,7 @@
from airflow.providers.amazon.version_compat import NOTSET, ArgNotSet
from airflow.providers.common.compat.openlineage.utils.spark import (
inject_parent_job_information_into_emr_serverless_properties,
inject_parent_job_information_into_spark_properties,
inject_transport_information_into_emr_serverless_properties,
)
from airflow.providers.common.compat.sdk import AirflowException, conf
Expand Down Expand Up @@ -100,6 +102,9 @@ class EmrAddStepsOperator(AwsBaseOperator[EmrHook]):
:param deferrable: If True, the operator will wait asynchronously for the job to complete.
This implies waiting for completion. This mode requires aiobotocore module to be installed.
(default: False)
:param openlineage_inject_parent_job_info: If True, injects OpenLineage parent job information
into Spark steps so the Spark job emits a ``parentRunFacet`` linking back to the Airflow task.
Defaults to the ``openlineage.spark_inject_parent_job_info`` config value.
"""

aws_hook_class = EmrHook
Expand Down Expand Up @@ -130,6 +135,9 @@ def __init__(
waiter_max_attempts: int = 60,
execution_role_arn: str | None = None,
deferrable: bool = conf.getboolean("operators", "default_deferrable", fallback=False),
openlineage_inject_parent_job_info: bool = conf.getboolean(
"openlineage", "spark_inject_parent_job_info", fallback=False
),
**kwargs,
):
if not exactly_one(job_flow_id is None, job_flow_name is None):
Expand All @@ -146,6 +154,36 @@ def __init__(
self.waiter_max_attempts = waiter_max_attempts
self.execution_role_arn = execution_role_arn
self.deferrable = deferrable
self.openlineage_inject_parent_job_info = openlineage_inject_parent_job_info

def _inject_openlineage_parent_job_information(self, steps: list[dict], context: Context) -> list[dict]:
parent_job_information = inject_parent_job_information_into_spark_properties({}, context)
if not parent_job_information:
return steps

parent_conf = [
argument
for key, value in parent_job_information.items()
for argument in ("--conf", f"{key}={value}")
]
result = copy.deepcopy(steps)
for step in result:
hadoop_jar_step = step.get("HadoopJarStep", {})
arguments = hadoop_jar_step.get("Args", [])
if hadoop_jar_step.get("Jar", "").rsplit("/", 1)[-1] != "command-runner.jar" or not arguments:
continue
command = arguments[0].rsplit("/", 1)[-1]
if command not in {"run-example", "spark-submit"}:
continue
if any("spark.openlineage.parent" in argument for argument in arguments):
self.log.info(
"Some OpenLineage properties with parent job information are already present "
"in EMR Spark step arguments. Skipping injection for step `%s`.",
step.get("Name", ""),
)
continue
hadoop_jar_step["Args"] = [arguments[0], *parent_conf, *arguments[1:]]
return result

def execute(self, context: Context) -> list[str]:
job_flow_id = self.job_flow_id or self.hook.get_cluster_id_by_name(
Expand Down Expand Up @@ -181,6 +219,9 @@ def execute(self, context: Context) -> list[str]:
steps = self.steps
if isinstance(steps, str):
steps = ast.literal_eval(steps)
if self.openlineage_inject_parent_job_info:
self.log.info("Injecting OpenLineage parent job information into EMR Spark steps.")
steps = self._inject_openlineage_parent_job_information(steps, context)
step_ids = self.hook.add_job_flow_steps(
job_flow_id=job_flow_id,
steps=steps,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -323,3 +323,140 @@ def test_template_fields(self):
steps=self._config,
)
validate_template_fields(op)

@patch(
"airflow.providers.amazon.aws.operators.emr.inject_parent_job_information_into_spark_properties",
autospec=True,
)
def test_inject_openlineage_parent_job_information(self, mock_inject, mocked_hook_client):
mock_inject.return_value = {
"spark.openlineage.parentRunId": "parent-run-id",
"spark.openlineage.parentJobName": "test_dag_id.test_task",
}
mocked_hook_client.add_job_flow_steps.return_value = ADD_STEPS_SUCCESS_RETURN
steps = [
{
"Name": "spark-submit",
"HadoopJarStep": {
"Jar": "command-runner.jar",
"Args": ["spark-submit", "--deploy-mode", "cluster", "s3://bucket/job.py"],
},
},
{
"Name": "run-example",
"HadoopJarStep": {
"Jar": "command-runner.jar",
"Args": ["/usr/lib/spark/bin/run-example", "SparkPi", "10"],
},
},
{
"Name": "existing-parent-information",
"HadoopJarStep": {
"Jar": "command-runner.jar",
"Args": [
"spark-submit",
"--conf",
"spark.openlineage.parentRunId=existing-run-id",
"s3://bucket/job.py",
],
},
},
{
"Name": "non-spark",
"HadoopJarStep": {
"Jar": "command-runner.jar",
"Args": ["bash", "-c", "echo done"],
},
},
{
"Name": "custom-jar",
"HadoopJarStep": {
"Jar": "s3://bucket/job.jar",
"Args": ["spark-submit", "application-argument"],
},
},
{
"Name": "no-arguments",
"HadoopJarStep": {"Jar": "command-runner.jar"},
},
]
context = MagicMock(spec=dict)
operator = EmrAddStepsOperator(
task_id="test_task",
job_flow_id="j-8989898989",
aws_conn_id="aws_default",
steps=steps,
openlineage_inject_parent_job_info=True,
dag=DAG("test_dag_id", schedule=None, default_args=self.args),
)

operator.execute(context)

submitted_steps = mocked_hook_client.add_job_flow_steps.call_args.kwargs["Steps"]
assert submitted_steps[0]["HadoopJarStep"]["Args"] == [
"spark-submit",
"--conf",
"spark.openlineage.parentRunId=parent-run-id",
"--conf",
"spark.openlineage.parentJobName=test_dag_id.test_task",
"--deploy-mode",
"cluster",
"s3://bucket/job.py",
]
assert submitted_steps[1]["HadoopJarStep"]["Args"] == [
"/usr/lib/spark/bin/run-example",
"--conf",
"spark.openlineage.parentRunId=parent-run-id",
"--conf",
"spark.openlineage.parentJobName=test_dag_id.test_task",
"SparkPi",
"10",
]
assert submitted_steps[2:] == steps[2:]
assert operator.steps == steps
mock_inject.assert_called_once_with({}, context)

@pytest.mark.parametrize(
("enabled", "parent_job_information", "expected_call_count"),
[
pytest.param(False, {"spark.openlineage.parentRunId": "parent-run-id"}, 0, id="disabled"),
pytest.param(True, {}, 1, id="no-parent-information"),
],
)
@patch(
"airflow.providers.amazon.aws.operators.emr.inject_parent_job_information_into_spark_properties",
autospec=True,
)
def test_does_not_inject_openlineage_parent_job_information(
self,
mock_inject,
enabled,
parent_job_information,
expected_call_count,
mocked_hook_client,
):
mock_inject.return_value = parent_job_information
mocked_hook_client.add_job_flow_steps.return_value = ADD_STEPS_SUCCESS_RETURN
steps = [
{
"Name": "spark-submit",
"HadoopJarStep": {
"Jar": "command-runner.jar",
"Args": ["spark-submit", "s3://bucket/job.py"],
},
}
]
operator = EmrAddStepsOperator(
task_id="test_task",
job_flow_id="j-8989898989",
aws_conn_id="aws_default",
steps=steps,
openlineage_inject_parent_job_info=enabled,
dag=DAG("test_dag_id", schedule=None, default_args=self.args),
)

operator.execute(MagicMock(spec=dict))

submitted_steps = mocked_hook_client.add_job_flow_steps.call_args.kwargs["Steps"]
assert submitted_steps == steps
assert mock_inject.call_count == expected_call_count
38 changes: 38 additions & 0 deletions providers/openlineage/docs/spark.rst
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,9 @@ to manually configure these properties in every Spark operator.

Automatic injection is supported for the following operators:

- :class:`~airflow.providers.amazon.aws.operators.emr.EmrAddStepsOperator`
- :class:`~airflow.providers.amazon.aws.operators.emr.EmrContainerOperator`
- :class:`~airflow.providers.amazon.aws.operators.emr.EmrServerlessStartJobOperator`
- :class:`~airflow.providers.amazon.aws.operators.glue.GlueJobOperator`
- :class:`~airflow.providers.apache.livy.operators.livy.LivyOperator`
- :class:`~airflow.providers.apache.spark.operators.spark_submit.SparkSubmitOperator`
Expand All @@ -110,6 +113,18 @@ Automatic injection is supported for the following operators:
job itself must still have the Spark OpenLineage integration (the ``OpenLineageSparkListener``) enabled for
lineage to be emitted.

.. note::

:class:`~airflow.providers.amazon.aws.operators.emr.EmrAddStepsOperator` supports parent job information
injection for classic EMR steps that use ``command-runner.jar`` to launch ``spark-submit`` or ``run-example``.
It adds each property as a separate ``--conf key=value`` argument before the existing Spark arguments.
Steps with manually configured ``spark.openlineage.parent*`` properties and non-Spark steps are left unchanged.
Automatic transport injection is not supported for this operator.

:class:`~airflow.providers.amazon.aws.operators.emr.EmrContainerOperator` and
:class:`~airflow.providers.amazon.aws.operators.emr.EmrServerlessStartJobOperator` inject parent and transport
properties through their ``spark-defaults`` configuration.


.. _options:spark_inject_parent_job_info:

Expand Down Expand Up @@ -217,6 +232,29 @@ This allows you to customize the injection behavior for specific operators while
},
)

For a classic EMR Spark step, enable parent job information injection on
:class:`~airflow.providers.amazon.aws.operators.emr.EmrAddStepsOperator`:

.. code-block:: python

from airflow.providers.amazon.aws.operators.emr import EmrAddStepsOperator

EmrAddStepsOperator(
task_id="submit_spark_step",
job_flow_id="j-1234567890",
steps=[
{
"Name": "process-data",
"ActionOnFailure": "CONTINUE",
"HadoopJarStep": {
"Jar": "command-runner.jar",
"Args": ["spark-submit", "s3://example-bucket/jobs/process_data.py"],
},
}
],
openlineage_inject_parent_job_info=True,
)



Complete Example
Expand Down