Skip to content

Fix an issue with getting dataflow_job_id for Dataflow Runner when logging is disabled in Pipeline - #72809

Open
MaksYermak wants to merge 1 commit into
apache:mainfrom
VladaZakharova:fix-getting-dataflow-job-if-when-logs-is-off
Open

MaksYermak wants to merge 1 commit into
apache:mainfrom
VladaZakharova:fix-getting-dataflow-job-if-when-logs-is-off

Conversation

@MaksYermak

Copy link
Copy Markdown
Contributor

In this PR I have fixed an issue in BeamRun*PipelineOperators with getting dataflow_job_id for Dataflow Runner when logging is disabled in user's Pipeline code.


Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

@MaksYermak
MaksYermak requested a review from shahar1 as a code owner September 9, 2026 13:55
@boring-cyborg boring-cyborg Bot added area:providers provider:apache-beam provider:google Google (including GCP) related issues labels Sep 9, 2026
@VladaZakharova

Copy link
Copy Markdown
Contributor

@potiuk Can you please check these changes here ?:)

@aaron-y-chen aaron-y-chen left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

  1. With append_job_name=False, could the first five-second timeout pick a previous completed job while the new job is still being submitted?
  2. If the launcher exits without logging a job ID before the timeout, when would the API lookup run before deferring?

@MaksYermak

Copy link
Copy Markdown
Contributor Author
  1. With append_job_name=False, could the first five-second timeout pick a previous completed job while the new job is still being submitted?
  2. If the launcher exits without logging a job ID before the timeout, when would the API lookup run before deferring?
  1. Sorry I do not clearly understand the question, about what "first five-second timeout" are you talking? Could you please share a line in the code?
    About creating Dataflow jobs when append_job_name=False.
    For Apache Beam SDK only, Dataflow API raise DataflowJobAlreadyExistsError in case when user tries to create Job with the same job_name as the currently active Job has. Because for Beam operators we use only SDK for creating Job it means for as that we can use job_name as a parameter for monitoring Job's state. As we can be sure that in Dataflow we have only one Job with this name and in the running state.
  2. Sorry I do not understand this question either. What launcher is? What timeout are you talking about? What API did you mean? Could you please clarify it with sharing lines in the code?

@aaron-y-chen aaron-y-chen left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sorry for the unclear comment. I hope this one is clearer and more helpful.

self._job_id = active_jobs[0]["id"]
else:
jobs.sort(key=lambda j: j.get("createTime", ""), reverse=True)
self._job_id = jobs[0]["id"]

@aaron-y-chen aaron-y-chen Sep 16, 2026 •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

With append_job_name=False, the job name stays the same for every run.

For example, this Dag could hit the issue.

BeamRunPythonPipelineOperator(
    task_id="run_beam",
    runner="DataflowRunner",
    py_file="gs://my-bucket/etl.py",
    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(5 seconds), before this run's job is created, since staging takes much longer than that.

So _fetch_jobs_by_prefix_name() returns every job in the region, including finished ones, and keeps the ones whose name starts with daily-etl, like the jobs from the past two days. jobs would look something like this:

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:

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

Neither of these jobs is active, so the DataflowJobAlreadyExistsError rule doesn't prevent this. 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.

Would it make sense to only accept jobs created after this task started, and keep waiting when every match is terminal?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@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:

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@aaron-y-chen I have tested with the same code and I did not see any problem which you had described.
I ran two Jobs with the same name sequentially:
Screenshot 2026-09-22 at 16 32 24

Here is the screenshot from the first run:
Screenshot 2026-09-22 at 16 33 14
we can see that start-python-deferrable Job in Pending state and the next line is starting the deferreble mode.

Here is the screenshot from the second run:
Screenshot 2026-09-22 at 16 33 33
on this screenshot we can see on line 188 that start-python-deferrable Job in Pending state and on line 189 that start-python-deferrable Job in Done state. On the line 189 information about previous Jobs runs with the same name and it is not a problem for code at all. Because, we have start-python-deferrable Job in Pending state the code goes to deferreble mode on line 190.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Got it. Thank you for the investigation. I think I have no issues here now.

@VladaZakharova

Copy link
Copy Markdown
Contributor

@potiuk Can we please merge this change? Looks like we got several +2 even :)

@potiuk potiuk left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for tracking this down. Falling back to an API lookup when the launcher is silent is the right idea, but a few things need fixing before this can go in:

  1. The "running job" lookup can return a finished job. get_job_id_of_running_job → _get_current_jobs matches by name prefix and then takes the single match regardless of state, or, among several, the newest by createTime even if it's JOB_STATE_DONE. With append_job_name=False and a silent launcher (the exact case this PR targets), the first lookup fires 5 s after start, usually before staging has created the new job. On the second run of such a Dag the previous DONE job is returned, the operator defers on it and the task succeeds while the real job runs unmonitored. Your manual test had stderr warnings delaying the first timeout until the job was already PENDING, so it doesn't cover this. Please restrict this path to non-terminal jobs with an exact name match (jobs.list with filter=ACTIVE also makes it much cheaper), and keep the "newest terminal job" fallback out of the stop-reading decision.

  2. Cross-provider compatibility. The beam operator calls DataflowHook.get_job_id_of_running_job, which this PR adds to the google provider, so it will only exist from the next google release. Beam's google extra doesn't depend on the google provider at all (it only lists apache-beam[gcp]), so new beam plus an already-installed google provider fails with AttributeError after 5 s of launcher silence.

    Please declare the dependency in providers/apache/beam/pyproject.toml, in the google extra, with the special # use next version marker:

    "google" = [
        "apache-beam[gcp]>=2.76.0",
        "apache-airflow-providers-google>=22.6.0",  # use next version
    ]

    Keep the version at the current one (22.6.0) — don't bump it yourself. The exact # use next version comment tells the release tooling to raise it to the google version being released alongside this beam release, i.e. the first one that has get_job_id_of_running_job, and the Release Manager releases both together. See the # use next version explanation in contributing-docs/13_airflow_dependencies_and_extras.rst. Alternatively, guard the call (e.g. hasattr(self.dataflow_hook, "get_job_id_of_running_job")) and fall back to the current behaviour, but the dependency marker is the cleaner fix.

  3. Tests for the beam side. Nothing tests the timeout → API-lookup path in run_beam_command, the os.set_blocking change, the operator callback, or get_job_id_of_running_job / fetch_job_id. Each changed behaviour needs a test that fails without the change. (Please also mock os.set_blocking there; with a MagicMock Popen it currently ends up calling set_blocking(1, False) on the test process's stdout.)

  4. Non-blocking reads for all runners. process_fd decodes whatever readline() returns. In non-blocking mode that can be a partial line, so a split UTF-8 character raises UnicodeDecodeError, and the job-id regex can miss an ID split across reads. Please decode defensively (or buffer up to newline) and consider limiting the non-blocking switch to the Dataflow path.

  5. Lookup robustness. The lookup runs on every 5 s of silence and lists all jobs in the region each time. If it raises (e.g. missing dataflow.jobs.list permission) the task fails while the launcher keeps running. Please catch and log the error, then continue reading.

Smaller points: the new .. warning:: says the subprocess is terminated and only in deferrable mode, but the code just stops reading it and also returns early in non-deferrable mode. Also, fetch_job_id / get_job_id_of_running_job could use -> str | None annotations.


Drafted-by: Claude Code (Opus 5); reviewed by @potiuk before posting

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

4 participants