From 3f9c22faebffb433be9b3b32b3bf0bd511d577df Mon Sep 17 00:00:00 2001 From: Voldurk <40298266+Voldurk@users.noreply.github.com> Date: Tue, 16 Aug 2022 18:42:26 +0200 Subject: [PATCH 1/3] Update DataflowTemplatedJobStartOperator added possibility to pass append_job_name as DataflowHook.start_template_dataflow() argument when initializing DataflowTemplatedJobStartOperator --- airflow/providers/google/cloud/operators/dataflow.py | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/airflow/providers/google/cloud/operators/dataflow.py b/airflow/providers/google/cloud/operators/dataflow.py index 677bc94940dba..6f3dadde87f66 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: Optional[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 From 0d9bc3bfe8a549825f34af8d7f4cfe0bcc5407fe Mon Sep 17 00:00:00 2001 From: Voldurk <40298266+Voldurk@users.noreply.github.com> Date: Wed, 17 Aug 2022 01:26:42 +0200 Subject: [PATCH 2/3] Removed Optional Removed optional since None is not acceptable in DataflowHook --- airflow/providers/google/cloud/operators/dataflow.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/airflow/providers/google/cloud/operators/dataflow.py b/airflow/providers/google/cloud/operators/dataflow.py index 6f3dadde87f66..468aa9ba4a606 100644 --- a/airflow/providers/google/cloud/operators/dataflow.py +++ b/airflow/providers/google/cloud/operators/dataflow.py @@ -613,7 +613,7 @@ def __init__( environment: Optional[Dict] = None, cancel_timeout: Optional[int] = 10 * 60, wait_until_finished: Optional[bool] = None, - append_job_name: Optional[bool] = True, + append_job_name: bool = True, **kwargs, ) -> None: super().__init__(**kwargs) From 9747a7eba2d6a120d44829fe9ac6a41e28af9174 Mon Sep 17 00:00:00 2001 From: Voldurk <40298266+Voldurk@users.noreply.github.com> Date: Sat, 20 Aug 2022 15:57:39 +0200 Subject: [PATCH 3/3] updated unittest --- tests/providers/google/cloud/operators/test_dataflow.py | 1 + 1 file changed, 1 insertion(+) 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, )