This is an automated email from the ASF dual-hosted git repository.

potiuk pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new b2c215cc037 Fix Beam async hook launching pipelines through a shell 
(#72510)
b2c215cc037 is described below

commit b2c215cc03744640f381e4e77b47d2c24d74baf7
Author: Zhaoqi Xu <[email protected]>
AuthorDate: Sun Oct 4 08:24:36 2026 +0800

    Fix Beam async hook launching pipelines through a shell (#72510)
    
    * Fix Beam async hook launching pipelines through a shell
    
    * 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
    
    * 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
    
    ---------
    
    Co-authored-by: Jarek Potiuk <[email protected]>
---
 .../airflow/providers/apache/beam/hooks/beam.py    | 45 ++++++++++++----------
 .../beam/tests/unit/apache/beam/hooks/test_beam.py | 42 ++++++++++++++++++++
 2 files changed, 67 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 ae44d7f42bc..86ec5595503 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,28 @@ class BeamAsyncHook(BeamHook):
 
     @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
+        returncode: int | 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 +636,12 @@ class BeamAsyncHook(BeamHook):
         :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 e9750170280..5b25e6c0929 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 @@ from unittest.mock import ANY, AsyncMock, MagicMock
 import pytest
 
 from airflow.providers.apache.beam.hooks.beam import (
+    _APACHE_BEAM_VERSION_SCRIPT,
     BeamAsyncHook,
     BeamHook,
     beam_options_to_args,
@@ -479,6 +480,23 @@ class TestBeamAsyncHook:
         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,27 @@ class TestBeamAsyncHook:
             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

Reply via email to