diff --git a/providers/amazon/docs/operators/emr/emr.rst b/providers/amazon/docs/operators/emr/emr.rst index 2d63fafcc2612..426ddf604ba9f 100644 --- a/providers/amazon/docs/operators/emr/emr.rst +++ b/providers/amazon/docs/operators/emr/emr.rst @@ -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 +`__ +for configuration details and an example. + .. _howto/operator:EmrTerminateJobFlowOperator: Terminate an EMR job flow diff --git a/providers/amazon/src/airflow/providers/amazon/aws/operators/emr.py b/providers/amazon/src/airflow/providers/amazon/aws/operators/emr.py index c1e5ee91556b6..aab675201497c 100644 --- a/providers/amazon/src/airflow/providers/amazon/aws/operators/emr.py +++ b/providers/amazon/src/airflow/providers/amazon/aws/operators/emr.py @@ -18,6 +18,7 @@ from __future__ import annotations import ast +import copy import warnings from collections.abc import Sequence from datetime import timedelta @@ -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 @@ -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 @@ -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): @@ -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( @@ -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, diff --git a/providers/amazon/tests/unit/amazon/aws/operators/test_emr_add_steps.py b/providers/amazon/tests/unit/amazon/aws/operators/test_emr_add_steps.py index e52eb36e7bbd6..7a8d6b3892ed9 100644 --- a/providers/amazon/tests/unit/amazon/aws/operators/test_emr_add_steps.py +++ b/providers/amazon/tests/unit/amazon/aws/operators/test_emr_add_steps.py @@ -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 diff --git a/providers/openlineage/docs/spark.rst b/providers/openlineage/docs/spark.rst index 6e8647bddc969..a324508b55caf 100644 --- a/providers/openlineage/docs/spark.rst +++ b/providers/openlineage/docs/spark.rst @@ -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` @@ -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: @@ -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