From aff724a95fb05fcbc0ec6aae7806ef41abdf9bf7 Mon Sep 17 00:00:00 2001 From: Zhaoqi Xu Date: Fri, 4 Sep 2026 17:20:22 +0800 Subject: [PATCH 1/3] Fix Beam async hook launching pipelines through a shell --- .../providers/apache/beam/hooks/beam.py | 44 +++++++------ .../tests/unit/apache/beam/hooks/test_beam.py | 61 +++++++++++++++++++ 2 files changed, 85 insertions(+), 20 deletions(-) diff --git a/providers/apache/beam/src/airflow/providers/apache/beam/hooks/beam.py b/providers/apache/beam/src/airflow/providers/apache/beam/hooks/beam.py index ae44d7f42bcaf..96320aeedbf29 100644 --- a/providers/apache/beam/src/airflow/providers/apache/beam/hooks/beam.py +++ b/providers/apache/beam/src/airflow/providers/apache/beam/hooks/beam.py @@ -469,19 +469,27 @@ async def _cleanup_tmp_dir(tmp_dir: str) -> None: @staticmethod async def _beam_version(py_interpreter: str) -> str: - version_script_cmd = shlex.join([py_interpreter, "-c", _APACHE_BEAM_VERSION_SCRIPT]) - proc = await asyncio.create_subprocess_shell( - version_script_cmd, - stdout=asyncio.subprocess.PIPE, - stderr=asyncio.subprocess.PIPE, - ) - stdout, stderr = await proc.communicate() - if proc.returncode != 0: + start_error: OSError | None = None + try: + proc = await asyncio.create_subprocess_exec( + py_interpreter, + "-c", + _APACHE_BEAM_VERSION_SCRIPT, + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + ) + except OSError as e: + start_error = e + stdout, stderr, returncode = b"", str(e).encode(), 1 + else: + stdout, stderr = await proc.communicate() + returncode = proc.returncode + if returncode != 0: msg = ( - f"Unable to retrieve Apache Beam version, return code {proc.returncode}." + f"Unable to retrieve Apache Beam version, return code {returncode}." f"\nstdout: {stdout.decode()}\nstderr: {stderr.decode()}" ) - raise AirflowException(msg) + raise AirflowException(msg) from start_error return stdout.decode().strip() async def start_python_pipeline_async( @@ -627,16 +635,12 @@ async def run_beam_command_async( :param process_line_callback: Optional callback which can be used to process stdout and stderr to detect job id """ - cmd_str_representation = " ".join(shlex.quote(c) for c in cmd) - log.info("Running command: %s", cmd_str_representation) - - # Creating a separate asynchronous process - process = await asyncio.create_subprocess_shell( - cmd_str_representation, - shell=True, - stdout=subprocess.PIPE, - stderr=subprocess.PIPE, - close_fds=True, + log.info("Running command: %s", " ".join(shlex.quote(c) for c in cmd)) + + process = await asyncio.create_subprocess_exec( + *cmd, + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, cwd=working_directory, ) # Waits for Apache Beam pipeline to complete. diff --git a/providers/apache/beam/tests/unit/apache/beam/hooks/test_beam.py b/providers/apache/beam/tests/unit/apache/beam/hooks/test_beam.py index e9750170280c0..0a5b1b5672ff3 100644 --- a/providers/apache/beam/tests/unit/apache/beam/hooks/test_beam.py +++ b/providers/apache/beam/tests/unit/apache/beam/hooks/test_beam.py @@ -29,6 +29,7 @@ import pytest from airflow.providers.apache.beam.hooks.beam import ( + _APACHE_BEAM_VERSION_SCRIPT, BeamAsyncHook, BeamHook, beam_options_to_args, @@ -479,6 +480,23 @@ async def test_beam_version_error(self): with pytest.raises(AirflowException, match="Unable to retrieve Apache Beam version"): await BeamAsyncHook._beam_version("python1") + @pytest.mark.asyncio + async def test_beam_version_invokes_interpreter_without_shell(self): + fake_proc = AsyncMock() + fake_proc.communicate = AsyncMock(return_value=(b"2.39.0\n", b"")) + fake_proc.returncode = 0 + interpreter = r"C:\Program Files\Python\python.exe" + with ( + mock.patch("asyncio.create_subprocess_exec", new=AsyncMock(return_value=fake_proc)) as mock_exec, + mock.patch("asyncio.create_subprocess_shell", new=AsyncMock()) as mock_shell, + ): + version = await BeamAsyncHook._beam_version(interpreter) + + mock_shell.assert_not_called() + mock_exec.assert_awaited_once() + assert mock_exec.await_args.args == (interpreter, "-c", _APACHE_BEAM_VERSION_SCRIPT) + assert version == "2.39.0" + @pytest.mark.asyncio @mock.patch("airflow.providers.apache.beam.hooks.beam.BeamAsyncHook.run_beam_command_async") async def test_start_pipline_async(self, mock_runner): @@ -690,3 +708,46 @@ async def test_start_java_pipeline_async(self, mock_start_pipeline, job_class, c command_prefix=command_prefix, process_line_callback=None, ) + + @pytest.mark.asyncio + async def test_run_beam_command_async_uses_exec_with_argv(self): + hook = BeamAsyncHook(runner=DEFAULT_RUNNER) + fake_proc = AsyncMock() + fake_proc.stdout.readline = AsyncMock(return_value=b"") + fake_proc.stderr.readline = AsyncMock(return_value=b"") + fake_proc.wait = AsyncMock(return_value=0) + cmd = [ + r"C:\Program Files\Python\python.exe", + r"C:\Program Files\pipelines\word count.py", + "--output=gs://test/output", + ] + with ( + mock.patch("asyncio.create_subprocess_exec", new=AsyncMock(return_value=fake_proc)) as mock_exec, + mock.patch("asyncio.create_subprocess_shell", new=AsyncMock()) as mock_shell, + ): + return_code = await hook.run_beam_command_async(cmd=cmd, log=logging.getLogger("beam-test")) + + mock_shell.assert_not_called() + mock_exec.assert_awaited_once() + assert mock_exec.await_args.args == tuple(cmd) + assert mock_exec.await_args.kwargs.get("shell") is None + assert return_code == 0 + + @pytest.mark.asyncio + async def test_run_beam_command_async_preserves_arguments_with_spaces(self, tmp_path): + hook = BeamAsyncHook(runner=DEFAULT_RUNNER) + marker = tmp_path / "got-arg.txt" + cmd = [ + sys.executable, + "-c", + "import pathlib, sys; pathlib.Path(sys.argv[1]).write_text(sys.argv[2])", + str(marker), + "hello world", + ] + return_code = await hook.run_beam_command_async( + cmd=cmd, + log=logging.getLogger("beam-test"), + working_directory=str(tmp_path), + ) + assert return_code == 0 + assert marker.read_text() == "hello world" From bfba8a0aada0ad655e26fee7875bb35e1ebc4d36 Mon Sep 17 00:00:00 2001 From: Jarek Potiuk Date: Sun, 4 Oct 2026 01:24:19 +0200 Subject: [PATCH 2/3] Drop live-process Beam test that does not exercise the fix The test passes on main too, because shlex.quote plus sh already keeps an argument with spaces together, and it spawns a real interpreter. test_run_beam_command_async_uses_exec_with_argv already covers the switch to create_subprocess_exec. Generated-by: Claude Opus 5 --- .../tests/unit/apache/beam/hooks/test_beam.py | 19 ------------------- 1 file changed, 19 deletions(-) diff --git a/providers/apache/beam/tests/unit/apache/beam/hooks/test_beam.py b/providers/apache/beam/tests/unit/apache/beam/hooks/test_beam.py index 0a5b1b5672ff3..5b25e6c0929ec 100644 --- a/providers/apache/beam/tests/unit/apache/beam/hooks/test_beam.py +++ b/providers/apache/beam/tests/unit/apache/beam/hooks/test_beam.py @@ -732,22 +732,3 @@ async def test_run_beam_command_async_uses_exec_with_argv(self): assert mock_exec.await_args.args == tuple(cmd) assert mock_exec.await_args.kwargs.get("shell") is None assert return_code == 0 - - @pytest.mark.asyncio - async def test_run_beam_command_async_preserves_arguments_with_spaces(self, tmp_path): - hook = BeamAsyncHook(runner=DEFAULT_RUNNER) - marker = tmp_path / "got-arg.txt" - cmd = [ - sys.executable, - "-c", - "import pathlib, sys; pathlib.Path(sys.argv[1]).write_text(sys.argv[2])", - str(marker), - "hello world", - ] - return_code = await hook.run_beam_command_async( - cmd=cmd, - log=logging.getLogger("beam-test"), - working_directory=str(tmp_path), - ) - assert return_code == 0 - assert marker.read_text() == "hello world" From dcda56fba26da33ac3b4e39549cee572494016a5 Mon Sep 17 00:00:00 2001 From: Jarek Potiuk Date: Sun, 4 Oct 2026 01:45:49 +0200 Subject: [PATCH 3/3] Fix mypy error in Beam async version check asyncio's Process.returncode is typed int | None, so assigning it to the int inferred from the OSError branch failed mypy-providers. Generated-by: Claude Opus 5 --- .../apache/beam/src/airflow/providers/apache/beam/hooks/beam.py | 1 + 1 file changed, 1 insertion(+) diff --git a/providers/apache/beam/src/airflow/providers/apache/beam/hooks/beam.py b/providers/apache/beam/src/airflow/providers/apache/beam/hooks/beam.py index 96320aeedbf29..86ec559550350 100644 --- a/providers/apache/beam/src/airflow/providers/apache/beam/hooks/beam.py +++ b/providers/apache/beam/src/airflow/providers/apache/beam/hooks/beam.py @@ -470,6 +470,7 @@ async def _cleanup_tmp_dir(tmp_dir: str) -> None: @staticmethod async def _beam_version(py_interpreter: str) -> str: start_error: OSError | None = None + returncode: int | None try: proc = await asyncio.create_subprocess_exec( py_interpreter,