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