aaron-y-chen commented on code in PR #72809:
URL: https://github.com/apache/airflow/pull/72809#discussion_r4023666984
##########
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:
With `append_job_name=False`, the [job name stays the
same](https://github.com/apache/airflow/blob/be95313ae7d67eb04868f381f61087df62df2ec5/providers/google/src/airflow/providers/google/cloud/hooks/dataflow.py#L1066-L1069)
for every run.
For example, this Dag could hit the issue.
```python
BeamRunPythonPipelineOperator(
task_id="run_beam",
runner="DataflowRunner",
deferrable=True,
dataflow_config=DataflowConfiguration(
job_name="daily-etl",
append_job_name=False,
location="us-central1",
),
)
```
Step 1.
The lookup runs at the first [`select`
timeout](https://github.com/apache/airflow/blob/be95313ae7d67eb04868f381f61087df62df2ec5/providers/apache/beam/src/airflow/providers/apache/beam/hooks/beam.py#L184)(5
seconds), before this run's job is created, staging takes much longer than
that.
So `_fetch_jobs_by_prefix_name()` returns every job in the region, like jobs
that finished in the past two days. After filtering, jobs would look something
like this:
```python
jobs = [
{"id": "2026-09-14_03_00_12-1182736549827",
"name": "daily-etl", "currentState": "JOB_STATE_DONE",
"createTime": "2026-09-14T03:00:12Z"},
{"id": "2026-09-15_03_00_09-9938174462051",
"name": "daily-etl", "currentState": "JOB_STATE_DONE",
"createTime": "2026-09-15T03:00:09Z"}
]
```
Step 2.
Since both are in `DataflowJobStatus.TERMINAL_STATES`, `active_jobs` ends up
empty, so we fall nto the `else` branch and `self._job_id` gets yesterday's job:
```python
else:
jobs.sort(key=lambda j: j.get("createTime", ""), reverse=True)
# "2026-09-15T03:00:09Z"
# "2026-09-14T03:00:12Z"
self._job_id = jobs[0]["id"]
# = "2026-09-15_03_00_09-9938174462051" <- yesterday's finished job
```
The operator then defers on a job that already reached `JOB_STATE_DONE`,
while the process that is still submitting this run's job is dropped.
--
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]