diff --git a/airflow/providers/google/cloud/operators/dataflow.py b/airflow/providers/google/cloud/operators/dataflow.py index 677bc94940dba..468aa9ba4a606 100644 --- a/airflow/providers/google/cloud/operators/dataflow.py +++ b/airflow/providers/google/cloud/operators/dataflow.py @@ -499,6 +499,7 @@ class DataflowTemplatedJobStartOperator(BaseOperator): `__ :param cancel_timeout: How long (in seconds) operator should wait for the pipeline to be successfully cancelled when task is being killed. + :param append_job_name: True if unique suffix has to be appended to job name. :param wait_until_finished: (Optional) If True, wait for the end of pipeline execution before exiting. If False, only submits job. @@ -612,6 +613,7 @@ def __init__( environment: Optional[Dict] = None, cancel_timeout: Optional[int] = 10 * 60, wait_until_finished: Optional[bool] = None, + append_job_name: bool = True, **kwargs, ) -> None: super().__init__(**kwargs) @@ -631,6 +633,7 @@ def __init__( self.environment = environment self.cancel_timeout = cancel_timeout self.wait_until_finished = wait_until_finished + self.append_job_name = append_job_name def execute(self, context: 'Context') -> dict: self.hook = DataflowHook( @@ -657,6 +660,7 @@ def set_current_job(current_job): project_id=self.project_id, location=self.location, environment=self.environment, + append_job_name=self.append_job_name, ) return job diff --git a/tests/providers/google/cloud/operators/test_dataflow.py b/tests/providers/google/cloud/operators/test_dataflow.py index 2cc37cb589629..790101cb4e915 100644 --- a/tests/providers/google/cloud/operators/test_dataflow.py +++ b/tests/providers/google/cloud/operators/test_dataflow.py @@ -482,6 +482,7 @@ def test_exec(self, dataflow_mock): project_id=None, location=TEST_LOCATION, environment={'maxWorkers': 2}, + append_job_name=True, )