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]