MaksYermak commented on code in PR #72809:
URL: https://github.com/apache/airflow/pull/72809#discussion_r4025292823


##########
providers/google/src/airflow/providers/google/cloud/hooks/dataflow.py:
##########
@@ -339,6 +339,15 @@ def _get_current_jobs(self) -> list[dict]:
             jobs = self._fetch_jobs_by_prefix_name(self._job_name.lower())
             if len(jobs) == 1:
                 self._job_id = jobs[0]["id"]
+            elif len(jobs) > 1 and not self._multiple_jobs:
+                active_jobs = [
+                    job for job in jobs if job.get("currentState") not in 
DataflowJobStatus.TERMINAL_STATES
+                ]
+                if len(active_jobs) == 1:
+                    self._job_id = active_jobs[0]["id"]
+                else:
+                    jobs.sort(key=lambda j: j.get("createTime", ""), 
reverse=True)
+                    self._job_id = jobs[0]["id"]

Review Comment:
   @aaron-y-chen hmm, I have tested this scenario locally during development. 
It was one of the problem which I tried to solve and I did not see issue which 
you describe. Because I was working on this issue in the June, I will retest 
one more time and let you know.
   Here is part of logs for similar run from June:
   ```
   {"timestamp":"2026-06-30T07:52:24.437128Z","level":"info","event":"Beam 
version: 
2.67.0","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.task.hooks.airflow.providers.apache.beam.hooks.beam.BeamHook","filename":"beam.py","lineno":319}
   {"timestamp":"2026-06-30T07:52:24.437437Z","level":"info","event":"Running 
command: /tmp/apache-beam-venvgg6s0hmf/bin/python 
/files/dags/resources/wordcounr_debug.py --runner=DataflowRunner 
--job_name=start-python-deferrable --project=TEST --region=europe-west3 
--labels=airflow-version=v3-3-0 
--output=gs://bucket_dataflow_native_python_bug_yermaklocal/output","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.task.hooks.airflow.providers.apache.beam.hooks.beam.BeamHook","filename":"beam.py","lineno":169}
   {"timestamp":"2026-06-30T07:52:24.439508Z","level":"info","event":"Start 
waiting for Apache Beam process to 
complete.","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.task.hooks.airflow.providers.apache.beam.hooks.beam.BeamHook","filename":"beam.py","lineno":185}
   
{"timestamp":"2026-06-30T07:52:27.893801Z","level":"warning","event":"/tmp/apache-beam-venvgg6s0hmf/lib/python3.10/site-packages/google/api_core/_python_version_support.py:255:
 FutureWarning: You are using a Python version (3.10.20) which Google will stop 
supporting in new releases of google.api_core once it reaches its end of life 
(2026-10-04). Please upgrade to the latest Python version, or at least Python 
3.11, to continue receiving updates for google.api_core past that 
date.","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.task.hooks.airflow.providers.apache.beam.hooks.beam.BeamHook","filename":"beam.py","lineno":148}
   {"timestamp":"2026-06-30T07:52:27.894111Z","level":"warning","event":"  
warnings.warn(message, 
FutureWarning)","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.task.hooks.airflow.providers.apache.beam.hooks.beam.BeamHook","filename":"beam.py","lineno":148}
   
{"timestamp":"2026-06-30T07:52:28.920120Z","level":"warning","event":"/tmp/apache-beam-venvgg6s0hmf/lib/python3.10/site-packages/google/api_core/_python_version_support.py:255:
 FutureWarning: You are using a Python version (3.10.20) which Google will stop 
supporting in new releases of google.cloud.bigquery_storage_v1 once it reaches 
its end of life (2026-10-04). Please upgrade to the latest Python version, or 
at least Python 3.11, to continue receiving updates for 
google.cloud.bigquery_storage_v1 past that 
date.","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.task.hooks.airflow.providers.apache.beam.hooks.beam.BeamHook","filename":"beam.py","lineno":148}
   {"timestamp":"2026-06-30T07:52:28.920430Z","level":"warning","event":"  
warnings.warn(message, 
FutureWarning)","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.task.hooks.airflow.providers.apache.beam.hooks.beam.BeamHook","filename":"beam.py","lineno":148}
   
{"timestamp":"2026-06-30T07:52:29.250451Z","level":"warning","event":"/tmp/apache-beam-venvgg6s0hmf/lib/python3.10/site-packages/google/api_core/_python_version_support.py:255:
 FutureWarning: You are using a Python version (3.10.20) which Google will stop 
supporting in new releases of google.pubsub_v1 once it reaches its end of life 
(2026-10-04). Please upgrade to the latest Python version, or at least Python 
3.11, to continue receiving updates for google.pubsub_v1 past that 
date.","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.task.hooks.airflow.providers.apache.beam.hooks.beam.BeamHook","filename":"beam.py","lineno":148}
   {"timestamp":"2026-06-30T07:52:29.250692Z","level":"warning","event":"  
warnings.warn(message, 
FutureWarning)","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.task.hooks.airflow.providers.apache.beam.hooks.beam.BeamHook","filename":"beam.py","lineno":148}
   
{"timestamp":"2026-06-30T07:52:29.600427Z","level":"warning","event":"WARNING:apache_beam.io.gcp.gcsio:Unexpected
 error occurred when checking soft delete policy for 
gs://dataflow-staging-europe-west3-50d011efd56dbb1d8aba79a3b2551fa3","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.task.hooks.airflow.providers.apache.beam.hooks.beam.BeamHook","filename":"beam.py","lineno":148}
   
{"timestamp":"2026-06-30T07:52:29.604767Z","level":"warning","event":"WARNING:apache_beam.io.gcp.gcsio:Unexpected
 error occurred when checking soft delete policy for 
gs://dataflow-staging-europe-west3-50d011efd56dbb1d8aba79a3b2551fa3","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.task.hooks.airflow.providers.apache.beam.hooks.beam.BeamHook","filename":"beam.py","lineno":148}
   
{"timestamp":"2026-06-30T07:52:30.521179Z","level":"warning","event":"WARNING:apache_beam.io.gcp.gcsio:Unexpected
 error occurred when checking soft delete policy for 
gs://dataflow-staging-europe-west3-50d011efd56dbb1d8aba79a3b2551fa3","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.task.hooks.airflow.providers.apache.beam.hooks.beam.BeamHook","filename":"beam.py","lineno":148}
   
{"timestamp":"2026-06-30T07:52:30.525575Z","level":"warning","event":"WARNING:apache_beam.io.gcp.gcsio:Unexpected
 error occurred when checking soft delete policy for 
gs://dataflow-staging-europe-west3-50d011efd56dbb1d8aba79a3b2551fa3","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.task.hooks.airflow.providers.apache.beam.hooks.beam.BeamHook","filename":"beam.py","lineno":148}
   {"timestamp":"2026-06-30T07:52:35.876123Z","level":"info","event":"Google 
Cloud DataFlow job start-python-deferrable is state: 
JOB_STATE_PENDING","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.providers.google.cloud.hooks.dataflow._DataflowJobsController","filename":"dataflow.py","lineno":488}
   {"timestamp":"2026-06-30T07:52:35.876370Z","level":"info","event":"Google 
Cloud DataFlow job start-python-deferrable-c84651d8 is state: 
JOB_STATE_DONE","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.providers.google.cloud.hooks.dataflow._DataflowJobsController","filename":"dataflow.py","lineno":488}
   {"timestamp":"2026-06-30T07:52:35.876463Z","level":"info","event":"Google 
Cloud DataFlow job start-python-deferrable-ce1175dd is state: 
JOB_STATE_DONE","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.providers.google.cloud.hooks.dataflow._DataflowJobsController","filename":"dataflow.py","lineno":488}
   {"timestamp":"2026-06-30T07:52:35.876523Z","level":"info","event":"Google 
Cloud DataFlow job start-python-deferrable is state: 
JOB_STATE_DONE","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.providers.google.cloud.hooks.dataflow._DataflowJobsController","filename":"dataflow.py","lineno":488}
   {"timestamp":"2026-06-30T07:52:35.876574Z","level":"info","event":"Google 
Cloud DataFlow job start-python-deferrable is state: 
JOB_STATE_DONE","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.providers.google.cloud.hooks.dataflow._DataflowJobsController","filename":"dataflow.py","lineno":488}
   {"timestamp":"2026-06-30T07:52:35.876624Z","level":"info","event":"Google 
Cloud DataFlow job start-python-deferrable is state: 
JOB_STATE_DONE","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.providers.google.cloud.hooks.dataflow._DataflowJobsController","filename":"dataflow.py","lineno":488}
   {"timestamp":"2026-06-30T07:52:35.876671Z","level":"info","event":"Google 
Cloud DataFlow job start-python-deferrable is state: 
JOB_STATE_DONE","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.providers.google.cloud.hooks.dataflow._DataflowJobsController","filename":"dataflow.py","lineno":488}
   {"timestamp":"2026-06-30T07:52:35.876718Z","level":"info","event":"Google 
Cloud DataFlow job start-python-deferrable is state: 
JOB_STATE_DONE","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.providers.google.cloud.hooks.dataflow._DataflowJobsController","filename":"dataflow.py","lineno":488}
   {"timestamp":"2026-06-30T07:52:35.876765Z","level":"info","event":"Google 
Cloud DataFlow job start-python-deferrable-3aeb2c18 is state: 
JOB_STATE_DONE","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.providers.google.cloud.hooks.dataflow._DataflowJobsController","filename":"dataflow.py","lineno":488}
   {"timestamp":"2026-06-30T07:52:35.876814Z","level":"info","event":"Google 
Cloud DataFlow job start-python-deferrable-32d0d195 is state: 
JOB_STATE_DONE","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.providers.google.cloud.hooks.dataflow._DataflowJobsController","filename":"dataflow.py","lineno":488}
   {"timestamp":"2026-06-30T07:52:35.876861Z","level":"info","event":"Google 
Cloud DataFlow job start-python-deferrable is state: 
JOB_STATE_DONE","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.providers.google.cloud.hooks.dataflow._DataflowJobsController","filename":"dataflow.py","lineno":488}
   {"timestamp":"2026-06-30T07:52:35.876924Z","level":"info","event":"Google 
Cloud DataFlow job start-python-deferrable is state: 
JOB_STATE_FAILED","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.providers.google.cloud.hooks.dataflow._DataflowJobsController","filename":"dataflow.py","lineno":488}
   {"timestamp":"2026-06-30T07:52:35.876972Z","level":"info","event":"Google 
Cloud DataFlow job start-python-deferrable is state: 
JOB_STATE_DONE","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.providers.google.cloud.hooks.dataflow._DataflowJobsController","filename":"dataflow.py","lineno":488}
   {"timestamp":"2026-06-30T07:52:35.877017Z","level":"info","event":"Google 
Cloud DataFlow job start-python-deferrable-ed85f1e0 is state: 
JOB_STATE_DONE","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.providers.google.cloud.hooks.dataflow._DataflowJobsController","filename":"dataflow.py","lineno":488}
   {"timestamp":"2026-06-30T07:52:35.877063Z","level":"info","event":"Google 
Cloud DataFlow job start-python-deferrable-964332a5 is state: 
JOB_STATE_DONE","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.providers.google.cloud.hooks.dataflow._DataflowJobsController","filename":"dataflow.py","lineno":488}
   {"timestamp":"2026-06-30T07:52:35.877106Z","level":"info","event":"Google 
Cloud DataFlow job start-python-deferrable-babaa18f is state: 
JOB_STATE_DONE","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"airflow.providers.google.cloud.hooks.dataflow._DataflowJobsController","filename":"dataflow.py","lineno":488}
   
{"timestamp":"2026-06-30T07:52:37.376846Z","level":"info","event":"::group::Post
 
Execute","dag_id":"dataflow_native_python_bug","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","run_id":"manual__2026-06-30T07:51:47.583474+00:00","try_number":1,"task_id":"start_python_job_dataflow_deferrable","map_index":-1,"logger":"task","filename":"task_runner.py","lineno":1609}
   {"timestamp":"2026-06-30T07:52:37.377095Z","level":"info","event":"Pausing 
task as DEFERRED. 
","dag_id":"dataflow_native_python_bug","task_id":"start_python_job_dataflow_deferrable","run_id":"manual__2026-06-30T07:51:47.583474+00:00","ti_id":"019f1783-48a9-71f7-8896-10c4ce89e240","try_number":1,"map_index":-1,"logger":"task","filename":"task_runner.py","lineno":1453}
   ```
   As you can see I have a new started Job in the `Pending` state. For me is 
interesting why in your list you do not have a new started Job in the `Pending` 
state?
   In my run I used this code:
   ```python
   start_python_job_dataflow_deferrable = BeamRunPythonPipelineOperator(
           runner=BeamRunnerType.DataflowRunner,
           task_id="start_python_job_dataflow_deferrable",
           py_file=LOCAL_PYTHON_SCRIPT,
           py_options=[],
           pipeline_options={
               "output": GCS_OUTPUT,
           },
           py_requirements=["apache-beam[gcp]==2.67.0"],
           py_interpreter="python3",
           py_system_site_packages=False,
           dataflow_config={"location": LOCATION, "job_name": 
"start_python_deferrable", "max_num_workers": 1, "append_job_name": False},
           deferrable=True,
       )
   ```
   which is similar to yours.
   As I mentioned I will recheck with this code one more time and let you know. 



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to