From 4e92962b917fc7496341903692da4f5748476108 Mon Sep 17 00:00:00 2001 From: Vincent Hsiao <124506982+fat-catTW@users.noreply.github.com> Date: Thu, 23 Jul 2026 07:03:59 +0000 Subject: [PATCH 01/12] Allow PythonOperator to override task interpreter mode Some Python tasks need fresh interpreter startup semantics without forcing every task in the deployment onto the slower global mode. --- .../src/airflow/executors/base_executor.py | 1 + .../src/airflow/executors/workloads/task.py | 10 +++++++ .../tests/unit/executors/test_workloads.py | 19 +++++++++++++ providers/standard/docs/operators/python.rst | 17 ++++++++++++ .../providers/standard/operators/python.py | 16 ++++++++++- .../unit/standard/operators/test_python.py | 20 ++++++++++++++ .../airflow/sdk/execution_time/coordinator.py | 3 +++ .../airflow/sdk/execution_time/supervisor.py | 27 +++++++++++++++---- .../execution_time/test_supervisor.py | 17 ++++++++++++ 9 files changed, 124 insertions(+), 6 deletions(-) diff --git a/airflow-core/src/airflow/executors/base_executor.py b/airflow-core/src/airflow/executors/base_executor.py index 8030e9c14455e..84acfe04c73ae 100644 --- a/airflow-core/src/airflow/executors/base_executor.py +++ b/airflow-core/src/airflow/executors/base_executor.py @@ -724,6 +724,7 @@ def run_workload( log_path=workload.log_path, subprocess_logs_to_stdout=subprocess_logs_to_stdout, sentry_integration=getattr(workload, "sentry_integration", ""), + execute_tasks_new_python_interpreter=workload.execute_tasks_new_python_interpreter, ) if isinstance(workload, ExecuteCallback): from airflow.sdk.execution_time.callback_supervisor import supervise_callback diff --git a/airflow-core/src/airflow/executors/workloads/task.py b/airflow-core/src/airflow/executors/workloads/task.py index 3099fe1d77485..a3f48df1fd439 100644 --- a/airflow-core/src/airflow/executors/workloads/task.py +++ b/airflow-core/src/airflow/executors/workloads/task.py @@ -27,6 +27,14 @@ from airflow.executors.workloads.base import BaseDagBundleWorkload, BundleInfo from airflow.utils.state import TaskInstanceState + +def _get_execute_tasks_new_python_interpreter(ti: object) -> bool | None: + value = getattr(getattr(ti, "task", None), "execute_tasks_new_python_interpreter", None) + if value is None or isinstance(value, bool): + return value + return None + + if TYPE_CHECKING: from airflow.api_fastapi.auth.tokens import JWTGenerator from airflow.models.taskinstance import TaskInstance as TIModel @@ -67,6 +75,7 @@ class ExecuteTask(BaseDagBundleWorkload): ti: TaskInstanceDTO sentry_integration: str = "" + execute_tasks_new_python_interpreter: bool | None = None type: Literal["ExecuteTask"] = Field(init=False, default="ExecuteTask") @@ -122,4 +131,5 @@ def make( log_path=fname, bundle_info=bundle_info, sentry_integration=sentry_integration, + execute_tasks_new_python_interpreter=_get_execute_tasks_new_python_interpreter(ti), ) diff --git a/airflow-core/tests/unit/executors/test_workloads.py b/airflow-core/tests/unit/executors/test_workloads.py index 37fbcd96ce950..209d86e918003 100644 --- a/airflow-core/tests/unit/executors/test_workloads.py +++ b/airflow-core/tests/unit/executors/test_workloads.py @@ -234,6 +234,25 @@ def _make_mock_ti( return ti + @pytest.mark.parametrize("execute_tasks_new_python_interpreter", [True, False, None]) + def test_make_populates_execute_tasks_new_python_interpreter(self, execute_tasks_new_python_interpreter): + from unittest.mock import Mock + + ti = self._make_mock_ti(bundle_version="abc123", version_data={}) + ti.task = Mock(execute_tasks_new_python_interpreter=execute_tasks_new_python_interpreter) + + workload = ExecuteTask.make(ti) + + assert workload.execute_tasks_new_python_interpreter is execute_tasks_new_python_interpreter + + def test_make_defaults_execute_tasks_new_python_interpreter_to_none(self): + ti = self._make_mock_ti(bundle_version="abc123", version_data={}) + ti.task = None + + workload = ExecuteTask.make(ti) + + assert workload.execute_tasks_new_python_interpreter is None + def test_pinned_run_populates_version_data(self): """When the run is pinned, version_data from the run's created_dag_version flows to BundleInfo.""" version_data = {"schema_version": 1, "files": {"dags/my_dag.py": "ver123"}} diff --git a/providers/standard/docs/operators/python.rst b/providers/standard/docs/operators/python.rst index 30f7151e55bef..16c6d897e52f4 100644 --- a/providers/standard/docs/operators/python.rst +++ b/providers/standard/docs/operators/python.rst @@ -47,6 +47,23 @@ Use the :class:`~airflow.providers.standard.operators.python.PythonOperator` to :start-after: [START howto_operator_python] :end-before: [END howto_operator_python] +Running selected tasks in a new interpreter +^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ + +By default, task process creation follows the global +:ref:`execute_tasks_new_python_interpreter ` +configuration. Set ``execute_tasks_new_python_interpreter`` on a ``PythonOperator`` task to override +that global setting for a single task. This is useful when one task needs a fresh Python interpreter, +for example to pick up plugin changes immediately, without changing the execution mode for every task. + +.. code-block:: python + + PythonOperator( + task_id="refresh_plugin_state", + python_callable=refresh_plugin_state, + execute_tasks_new_python_interpreter=True, + ) + Passing in arguments ^^^^^^^^^^^^^^^^^^^^ diff --git a/providers/standard/src/airflow/providers/standard/operators/python.py b/providers/standard/src/airflow/providers/standard/operators/python.py index 8698ddc403b60..0f36819258380 100644 --- a/providers/standard/src/airflow/providers/standard/operators/python.py +++ b/providers/standard/src/airflow/providers/standard/operators/python.py @@ -35,7 +35,7 @@ from itertools import chain from pathlib import Path from tempfile import TemporaryDirectory -from typing import TYPE_CHECKING, Any, NamedTuple, cast +from typing import TYPE_CHECKING, Any, ClassVar, NamedTuple, cast import lazy_object_proxy from packaging.requirements import InvalidRequirement, Requirement @@ -169,6 +169,8 @@ def my_python_callable(**kwargs): logs. Defaults to True, which allows return value log output. It can be set to False to prevent log output of return value when you return huge data such as transmission a large amount of XCom to TaskAPI. + :param execute_tasks_new_python_interpreter: If set, overrides the global + ``[core] execute_tasks_new_python_interpreter`` setting for this task. """ template_fields: Sequence[str] = ("templates_dict", "op_args", "op_kwargs") @@ -176,6 +178,8 @@ def my_python_callable(**kwargs): BLUE = "#ffefeb" ui_color = BLUE + __serialized_fields: ClassVar[frozenset[str] | None] = None + # since we won't mutate the arguments, we should just do the shallow copy # there are some cases we can't deepcopy the objects(e.g protobuf). shallow_copy_attrs: Sequence[str] = ("python_callable", "op_kwargs") @@ -189,6 +193,7 @@ def __init__( templates_dict: dict[str, Any] | None = None, templates_exts: Sequence[str] | None = None, show_return_value_in_logs: bool = True, + execute_tasks_new_python_interpreter: bool | None = None, **kwargs, ) -> None: super().__init__(**kwargs) @@ -201,6 +206,15 @@ def __init__( if templates_exts: self.template_ext = templates_exts self.show_return_value_in_logs = show_return_value_in_logs + self.execute_tasks_new_python_interpreter = execute_tasks_new_python_interpreter + + @classmethod + def get_serialized_fields(cls): + if not cls.__serialized_fields: + cls.__serialized_fields = frozenset( + super().get_serialized_fields() | {"execute_tasks_new_python_interpreter"} + ) + return cls.__serialized_fields @property def is_async(self) -> bool: diff --git a/providers/standard/tests/unit/standard/operators/test_python.py b/providers/standard/tests/unit/standard/operators/test_python.py index d8228b1032f30..970b73cc34f12 100644 --- a/providers/standard/tests/unit/standard/operators/test_python.py +++ b/providers/standard/tests/unit/standard/operators/test_python.py @@ -291,6 +291,26 @@ def a_fn(): rendered_op_kwargs = task.op_kwargs assert rendered_op_kwargs["a_callable"] == a_fn + @pytest.mark.parametrize("execute_tasks_new_python_interpreter", [True, False, None]) + def test_execute_tasks_new_python_interpreter_is_serialized(self, execute_tasks_new_python_interpreter): + from airflow.serialization.serialized_objects import OperatorSerialization + + task = PythonOperator( + python_callable=lambda: None, + task_id=self.task_id, + execute_tasks_new_python_interpreter=execute_tasks_new_python_interpreter, + ) + + serialized = OperatorSerialization.serialize_operator(task) + deserialized = OperatorSerialization.deserialize_operator(serialized) + + assert "execute_tasks_new_python_interpreter" in PythonOperator.get_serialized_fields() + assert serialized.get("execute_tasks_new_python_interpreter") is execute_tasks_new_python_interpreter + assert ( + getattr(deserialized, "execute_tasks_new_python_interpreter", None) + is execute_tasks_new_python_interpreter + ) + def test_python_operator_shallow_copy_attr(self): def not_callable(x): raise RuntimeError("Should not be triggered") diff --git a/task-sdk/src/airflow/sdk/execution_time/coordinator.py b/task-sdk/src/airflow/sdk/execution_time/coordinator.py index 1469b27888c51..7175a1f929ad2 100644 --- a/task-sdk/src/airflow/sdk/execution_time/coordinator.py +++ b/task-sdk/src/airflow/sdk/execution_time/coordinator.py @@ -98,6 +98,7 @@ def execute_task( logger: FilteringBoundLogger | None = None, sentry_integration: str = "", subprocess_logs_to_stdout: bool, + execute_tasks_new_python_interpreter: bool | None = None, **kwargs, ) -> ExecutionResult: """ @@ -171,6 +172,7 @@ def execute_task( logger: FilteringBoundLogger | None = None, sentry_integration: str = "", subprocess_logs_to_stdout: bool, + execute_tasks_new_python_interpreter: bool | None = None, **kwargs, ) -> BaseCoordinator.ExecutionResult: # TODO: Importing this at the top causes circular imports. @@ -191,6 +193,7 @@ def execute_task( logger=logger, bundle_info=bundle_info, subprocess_logs_to_stdout=subprocess_logs_to_stdout, + execute_tasks_new_python_interpreter=execute_tasks_new_python_interpreter, sentry_integration=sentry_integration, ) exit_code = process.wait() diff --git a/task-sdk/src/airflow/sdk/execution_time/supervisor.py b/task-sdk/src/airflow/sdk/execution_time/supervisor.py index 87311f02da7a1..c2344c3bd5824 100644 --- a/task-sdk/src/airflow/sdk/execution_time/supervisor.py +++ b/task-sdk/src/airflow/sdk/execution_time/supervisor.py @@ -519,6 +519,17 @@ def _should_use_exec() -> bool: return sys.platform in _FORK_EXEC_PLATFORMS +def _should_use_exec_for_task(execute_tasks_new_python_interpreter: bool | None = None) -> bool: + """Whether task execution should ``exec`` a fresh Python interpreter.""" + if sys.platform in _FORK_EXEC_PLATFORMS: + return True + if execute_tasks_new_python_interpreter is None: + execute_tasks_new_python_interpreter = conf.getboolean( + "core", "execute_tasks_new_python_interpreter", fallback=False + ) + return execute_tasks_new_python_interpreter + + def _resolve_child_target(dotted: str) -> Callable[[], None]: """ Resolve a ``module:qualname`` string to the callable the exec'd child runs. @@ -688,8 +699,7 @@ def start( """ Fork and start a new subprocess with the specified target function. - :param use_exec: If True, on platforms that need it (currently macOS), - immediately ``os.execv`` a fresh Python interpreter after ``os.fork``. + :param use_exec: If True, immediately ``os.execv`` a fresh Python interpreter after ``os.fork``. This avoids macOS fork-safety issues with Objective-C frameworks. ``target`` is rehydrated in the exec'd child from its ``module:qualname``, so any importable entry point (task execution, DAG processor, triggerer) @@ -1352,13 +1362,16 @@ def start( # type: ignore[override] target: Callable[[], None] = _subprocess_main, logger: FilteringBoundLogger | None = None, sentry_integration: str = "", + execute_tasks_new_python_interpreter: bool | None = None, **kwargs, ) -> Self: """Fork and start a new subprocess to execute the given task.""" - # Opt in to fork+exec on platforms that need it (currently macOS). # Tests override `target` with a local stub to exercise the base - # infrastructure; keep bare fork for those. - use_exec = target is _subprocess_main and _should_use_exec() + # infrastructure; keep bare fork for those. Otherwise task execution + # can opt in to fork+exec through platform safety or configuration. + use_exec = target is _subprocess_main and _should_use_exec_for_task( + execute_tasks_new_python_interpreter + ) proc: Self = super().start( id=what.id, client=client, target=target, logger=logger, use_exec=use_exec, **kwargs ) @@ -2454,6 +2467,7 @@ def supervise_task( subprocess_logs_to_stdout: bool = False, client: Client | None = None, sentry_integration: str = "", + execute_tasks_new_python_interpreter: bool | None = None, ) -> int: """ Run a single task execution to completion. @@ -2469,6 +2483,8 @@ def supervise_task( :param client: Optional preconfigured client for communication with the server (Mostly for tests). :param sentry_integration: If the executor has a Sentry integration, import path to a callable to initialize it (empty means no integration). + :param execute_tasks_new_python_interpreter: Override ``[core] execute_tasks_new_python_interpreter`` + for this task. ``None`` keeps the global configuration. :return: Exit code of the process. :raises ValueError: If server URL is empty or invalid. :raises InvalidCoordinatorError: If the coordinator for the task is not @@ -2551,6 +2567,7 @@ def supervise_task( logger=logger, sentry_integration=sentry_integration, subprocess_logs_to_stdout=subprocess_logs_to_stdout, + execute_tasks_new_python_interpreter=execute_tasks_new_python_interpreter, ) end = time.monotonic() log.info( diff --git a/task-sdk/tests/task_sdk/execution_time/test_supervisor.py b/task-sdk/tests/task_sdk/execution_time/test_supervisor.py index 7ba463567a17a..1de7737f5765e 100644 --- a/task-sdk/tests/task_sdk/execution_time/test_supervisor.py +++ b/task-sdk/tests/task_sdk/execution_time/test_supervisor.py @@ -185,6 +185,23 @@ TI_ID = uuid7() +@pytest.mark.parametrize( + ("task_override", "global_config", "expected"), + [(True, False, True), (False, True, False), (None, True, True), (None, False, False)], +) +def test_should_use_exec_for_task_honors_task_override(monkeypatch, task_override, global_config, expected): + monkeypatch.setattr(supervisor.sys, "platform", "linux") + monkeypatch.setattr(supervisor.conf, "getboolean", lambda *args, **kwargs: global_config) + + assert supervisor._should_use_exec_for_task(task_override) is expected + + +def test_should_use_exec_for_task_keeps_platform_exec_requirement(monkeypatch): + monkeypatch.setattr(supervisor.sys, "platform", "darwin") + + assert supervisor._should_use_exec_for_task(False) is True + + def lineno(): """Returns the current line number in our program.""" return inspect.currentframe().f_back.f_lineno From 8ae6cdb6bf988103f35265692e5916f81c6cd83f Mon Sep 17 00:00:00 2001 From: Vincent Hsiao <124506982+fat-catTW@users.noreply.github.com> Date: Thu, 23 Jul 2026 12:16:31 +0000 Subject: [PATCH 02/12] Fix PythonOperator docs heading underline The provider docs build treats mismatched RST heading underline lengths as warnings, which fail CI for documentation builds. --- providers/standard/docs/operators/python.rst | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/providers/standard/docs/operators/python.rst b/providers/standard/docs/operators/python.rst index 16c6d897e52f4..2b1a04d3c4480 100644 --- a/providers/standard/docs/operators/python.rst +++ b/providers/standard/docs/operators/python.rst @@ -48,7 +48,7 @@ Use the :class:`~airflow.providers.standard.operators.python.PythonOperator` to :end-before: [END howto_operator_python] Running selected tasks in a new interpreter -^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ +^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ By default, task process creation follows the global :ref:`execute_tasks_new_python_interpreter ` From fd2383189cd168015070a86b13c1fd230cf5da11 Mon Sep 17 00:00:00 2001 From: Vincent Hsiao <124506982+fat-catTW@users.noreply.github.com> Date: Thu, 23 Jul 2026 13:01:07 +0000 Subject: [PATCH 03/12] Fix PythonOperator serialized task roundtrip Serialized PythonOperator tasks can omit default None values, so the deserialized placeholder needs to know the field exists for roundtrip comparisons. --- .../src/airflow/serialization/definitions/baseoperator.py | 1 + 1 file changed, 1 insertion(+) diff --git a/airflow-core/src/airflow/serialization/definitions/baseoperator.py b/airflow-core/src/airflow/serialization/definitions/baseoperator.py index 6bafc5891235a..0f72d2870b228 100644 --- a/airflow-core/src/airflow/serialization/definitions/baseoperator.py +++ b/airflow-core/src/airflow/serialization/definitions/baseoperator.py @@ -192,6 +192,7 @@ def get_serialized_fields(cls): "execution_timeout", "executor", "executor_config", + "execute_tasks_new_python_interpreter", "ignore_first_depends_on_past", "inlets", "is_setup", From 203a657962e6a1772d2b34a9542f5f95921cda35 Mon Sep 17 00:00:00 2001 From: Vincent Hsiao <124506982+fat-catTW@users.noreply.github.com> Date: Thu, 23 Jul 2026 15:12:43 +0000 Subject: [PATCH 04/12] Use compat serialization helper in PythonOperator test Provider compatibility jobs run standard provider tests against older Airflow releases, so serialization tests need the cross-version helper used by other provider tests. --- providers/standard/tests/unit/standard/operators/test_python.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/providers/standard/tests/unit/standard/operators/test_python.py b/providers/standard/tests/unit/standard/operators/test_python.py index 970b73cc34f12..85d98791430ef 100644 --- a/providers/standard/tests/unit/standard/operators/test_python.py +++ b/providers/standard/tests/unit/standard/operators/test_python.py @@ -293,7 +293,7 @@ def a_fn(): @pytest.mark.parametrize("execute_tasks_new_python_interpreter", [True, False, None]) def test_execute_tasks_new_python_interpreter_is_serialized(self, execute_tasks_new_python_interpreter): - from airflow.serialization.serialized_objects import OperatorSerialization + from tests_common.test_utils.compat import OperatorSerialization task = PythonOperator( python_callable=lambda: None, From 82514c176565ef8a26ed9577b970e60dd45d42f3 Mon Sep 17 00:00:00 2001 From: Vincent Hsiao <124506982+fat-catTW@users.noreply.github.com> Date: Fri, 24 Jul 2026 13:12:50 +0000 Subject: [PATCH 05/12] Update Edge worker API schema for task interpreter mode The ExecuteTask workload now carries the per-task interpreter override, so generated worker API schemas need to expose the field used by static validation. --- .../providers/edge3/worker_api/v2-edge-generated.yaml | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/providers/edge3/src/airflow/providers/edge3/worker_api/v2-edge-generated.yaml b/providers/edge3/src/airflow/providers/edge3/worker_api/v2-edge-generated.yaml index 13cacba97be36..708171f3f5c66 100644 --- a/providers/edge3/src/airflow/providers/edge3/worker_api/v2-edge-generated.yaml +++ b/providers/edge3/src/airflow/providers/edge3/worker_api/v2-edge-generated.yaml @@ -1136,6 +1136,11 @@ components: type: string title: Sentry Integration default: '' + execute_tasks_new_python_interpreter: + anyOf: + - type: boolean + - type: 'null' + title: Execute Tasks New Python Interpreter type: type: string const: ExecuteTask From 781dade4f0a688f813a96114f8d4b44f2e903ec3 Mon Sep 17 00:00:00 2001 From: Vincent Hsiao <124506982+fat-catTW@users.noreply.github.com> Date: Mon, 3 Aug 2026 16:12:57 +0000 Subject: [PATCH 06/12] Add newsfragment for task interpreter config fix --- airflow-core/newsfragments/70280.bugfix.rst | 1 + 1 file changed, 1 insertion(+) create mode 100644 airflow-core/newsfragments/70280.bugfix.rst diff --git a/airflow-core/newsfragments/70280.bugfix.rst b/airflow-core/newsfragments/70280.bugfix.rst new file mode 100644 index 0000000000000..f6f8d8382c8bc --- /dev/null +++ b/airflow-core/newsfragments/70280.bugfix.rst @@ -0,0 +1 @@ +Restored support for ``[core] execute_tasks_new_python_interpreter`` when running tasks with LocalExecutor and CeleryExecutor. From e92c54f413365d8da8212f168fb9d5f6704a1db7 Mon Sep 17 00:00:00 2001 From: Vincent Hsiao <124506982+fat-catTW@users.noreply.github.com> Date: Tue, 4 Aug 2026 16:34:54 +0000 Subject: [PATCH 07/12] Trigger CI for task interpreter PR From ef9b5f2b3cf1a64e28480677f5e29b31094ded19 Mon Sep 17 00:00:00 2001 From: Vincent Hsiao <124506982+fat-catTW@users.noreply.github.com> Date: Wed, 5 Aug 2026 06:01:20 +0000 Subject: [PATCH 08/12] Guard task interpreter serialization for compat --- .../providers/standard/operators/python.py | 17 ++++++++++++--- .../unit/standard/operators/test_python.py | 21 ++++++++++++++++++- 2 files changed, 34 insertions(+), 4 deletions(-) diff --git a/providers/standard/src/airflow/providers/standard/operators/python.py b/providers/standard/src/airflow/providers/standard/operators/python.py index 0f36819258380..c186f54efe765 100644 --- a/providers/standard/src/airflow/providers/standard/operators/python.py +++ b/providers/standard/src/airflow/providers/standard/operators/python.py @@ -76,6 +76,16 @@ log = logging.getLogger(__name__) + +def has_execute_tasks_new_python_interpreter_serialization_support() -> bool: + try: + from airflow.serialization.definitions.baseoperator import SerializedBaseOperator + except ImportError: + from airflow.serialization.serialized_objects import SerializedBaseOperator + + return "execute_tasks_new_python_interpreter" in SerializedBaseOperator.get_serialized_fields() + + if TYPE_CHECKING: from typing import Literal @@ -211,9 +221,10 @@ def __init__( @classmethod def get_serialized_fields(cls): if not cls.__serialized_fields: - cls.__serialized_fields = frozenset( - super().get_serialized_fields() | {"execute_tasks_new_python_interpreter"} - ) + serialized_fields = super().get_serialized_fields() + if has_execute_tasks_new_python_interpreter_serialization_support(): + serialized_fields = serialized_fields | {"execute_tasks_new_python_interpreter"} + cls.__serialized_fields = frozenset(serialized_fields) return cls.__serialized_fields @property diff --git a/providers/standard/tests/unit/standard/operators/test_python.py b/providers/standard/tests/unit/standard/operators/test_python.py index 85d98791430ef..90a205f052c47 100644 --- a/providers/standard/tests/unit/standard/operators/test_python.py +++ b/providers/standard/tests/unit/standard/operators/test_python.py @@ -64,7 +64,7 @@ from airflow.utils.state import DagRunState, State, TaskInstanceState from airflow.utils.types import DagRunType -from tests_common.test_utils.compat import TriggerRule, timezone +from tests_common.test_utils.compat import SerializedBaseOperator, TriggerRule, timezone from tests_common.test_utils.db import clear_db_runs from tests_common.test_utils.in_process_taskrun import pushed_xcom, run_task_no_db from tests_common.test_utils.taskinstance import get_template_context, run_task_instance @@ -89,6 +89,10 @@ from airflow.sdk import Context +def has_execute_tasks_new_python_interpreter_serialization_support() -> bool: + return "execute_tasks_new_python_interpreter" in SerializedBaseOperator.get_serialized_fields() + + AIRFLOW_ROOT_PATH = Path(__file__).parents[6] TI = TaskInstance @@ -291,6 +295,10 @@ def a_fn(): rendered_op_kwargs = task.op_kwargs assert rendered_op_kwargs["a_callable"] == a_fn + @pytest.mark.skipif( + not has_execute_tasks_new_python_interpreter_serialization_support(), + reason="Current Airflow core does not serialize execute_tasks_new_python_interpreter", + ) @pytest.mark.parametrize("execute_tasks_new_python_interpreter", [True, False, None]) def test_execute_tasks_new_python_interpreter_is_serialized(self, execute_tasks_new_python_interpreter): from tests_common.test_utils.compat import OperatorSerialization @@ -311,6 +319,17 @@ def test_execute_tasks_new_python_interpreter_is_serialized(self, execute_tasks_ is execute_tasks_new_python_interpreter ) + @mock.patch( + "airflow.providers.standard.operators.python.has_execute_tasks_new_python_interpreter_serialization_support", + return_value=False, + ) + def test_execute_tasks_new_python_interpreter_serialized_field_requires_core_support(self, _): + PythonOperator._PythonOperator__serialized_fields = None + try: + assert "execute_tasks_new_python_interpreter" not in PythonOperator.get_serialized_fields() + finally: + PythonOperator._PythonOperator__serialized_fields = None + def test_python_operator_shallow_copy_attr(self): def not_callable(x): raise RuntimeError("Should not be triggered") From 59688f983a98149385d3c8b9465b96a5d4b41eef Mon Sep 17 00:00:00 2001 From: Vincent Hsiao <124506982+fat-catTW@users.noreply.github.com> Date: Wed, 5 Aug 2026 15:28:26 +0000 Subject: [PATCH 09/12] Handle older serialized operator compatibility --- .../src/airflow/providers/standard/operators/python.py | 5 ++++- .../standard/tests/unit/standard/operators/test_python.py | 5 ++++- 2 files changed, 8 insertions(+), 2 deletions(-) diff --git a/providers/standard/src/airflow/providers/standard/operators/python.py b/providers/standard/src/airflow/providers/standard/operators/python.py index c186f54efe765..6f6dad20bc6b1 100644 --- a/providers/standard/src/airflow/providers/standard/operators/python.py +++ b/providers/standard/src/airflow/providers/standard/operators/python.py @@ -83,7 +83,10 @@ def has_execute_tasks_new_python_interpreter_serialization_support() -> bool: except ImportError: from airflow.serialization.serialized_objects import SerializedBaseOperator - return "execute_tasks_new_python_interpreter" in SerializedBaseOperator.get_serialized_fields() + get_serialized_fields = getattr(SerializedBaseOperator, "get_serialized_fields", None) + if get_serialized_fields is None: + return False + return "execute_tasks_new_python_interpreter" in get_serialized_fields() if TYPE_CHECKING: diff --git a/providers/standard/tests/unit/standard/operators/test_python.py b/providers/standard/tests/unit/standard/operators/test_python.py index 90a205f052c47..52cc46ff6f0d4 100644 --- a/providers/standard/tests/unit/standard/operators/test_python.py +++ b/providers/standard/tests/unit/standard/operators/test_python.py @@ -90,7 +90,10 @@ def has_execute_tasks_new_python_interpreter_serialization_support() -> bool: - return "execute_tasks_new_python_interpreter" in SerializedBaseOperator.get_serialized_fields() + get_serialized_fields = getattr(SerializedBaseOperator, "get_serialized_fields", None) + if get_serialized_fields is None: + return False + return "execute_tasks_new_python_interpreter" in get_serialized_fields() AIRFLOW_ROOT_PATH = Path(__file__).parents[6] From 202a7f2dc256f1a22a590cda8e8e810bc7717a30 Mon Sep 17 00:00:00 2001 From: Vincent Hsiao <124506982+fat-catTW@users.noreply.github.com> Date: Fri, 7 Aug 2026 13:08:56 +0000 Subject: [PATCH 10/12] Read task interpreter override from serialized task --- .../src/airflow/executors/workloads/task.py | 11 ++----- .../src/airflow/jobs/scheduler_job_runner.py | 10 +++++++ .../tests/unit/executors/test_workloads.py | 12 ++++---- .../tests/unit/jobs/test_scheduler_job.py | 30 +++++++++++++++++++ 4 files changed, 48 insertions(+), 15 deletions(-) diff --git a/airflow-core/src/airflow/executors/workloads/task.py b/airflow-core/src/airflow/executors/workloads/task.py index a3f48df1fd439..c661e3d6cded5 100644 --- a/airflow-core/src/airflow/executors/workloads/task.py +++ b/airflow-core/src/airflow/executors/workloads/task.py @@ -27,14 +27,6 @@ from airflow.executors.workloads.base import BaseDagBundleWorkload, BundleInfo from airflow.utils.state import TaskInstanceState - -def _get_execute_tasks_new_python_interpreter(ti: object) -> bool | None: - value = getattr(getattr(ti, "task", None), "execute_tasks_new_python_interpreter", None) - if value is None or isinstance(value, bool): - return value - return None - - if TYPE_CHECKING: from airflow.api_fastapi.auth.tokens import JWTGenerator from airflow.models.taskinstance import TaskInstance as TIModel @@ -105,6 +97,7 @@ def make( generator: JWTGenerator | None = None, bundle_info: BundleInfo | None = None, sentry_integration: str = "", + execute_tasks_new_python_interpreter: bool | None = None, ) -> ExecuteTask: """Create an ExecuteTask workload from a TaskInstance ORM model.""" from airflow.utils.helpers import log_filename_template_renderer @@ -131,5 +124,5 @@ def make( log_path=fname, bundle_info=bundle_info, sentry_integration=sentry_integration, - execute_tasks_new_python_interpreter=_get_execute_tasks_new_python_interpreter(ti), + execute_tasks_new_python_interpreter=execute_tasks_new_python_interpreter, ) diff --git a/airflow-core/src/airflow/jobs/scheduler_job_runner.py b/airflow-core/src/airflow/jobs/scheduler_job_runner.py index a9efa1c03d62c..2b62aef1b6c70 100644 --- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py +++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py @@ -1112,6 +1112,15 @@ def _get_sentry_integration(executor: BaseExecutor) -> str: return "" return sentry_integration + def _get_execute_tasks_new_python_interpreter(ti: TI) -> bool | None: + serialized_dag = self.scheduler_dag_bag.get_dag_for_run(dag_run=ti.dag_run, session=session) + if not serialized_dag or not serialized_dag.has_task(ti.task_id): + return None + value = getattr(serialized_dag.get_task(ti.task_id), "execute_tasks_new_python_interpreter", None) + if value is None or isinstance(value, bool): + return value + return None + # actually enqueue them for ti in task_instances: if ti.dag_run.state in State.finished_dr_states: @@ -1137,6 +1146,7 @@ def _get_sentry_integration(executor: BaseExecutor) -> str: ti, generator=executor.jwt_generator, sentry_integration=_get_sentry_integration(executor), + execute_tasks_new_python_interpreter=_get_execute_tasks_new_python_interpreter(ti), ) executor.queue_workload(workload, session=session) diff --git a/airflow-core/tests/unit/executors/test_workloads.py b/airflow-core/tests/unit/executors/test_workloads.py index 209d86e918003..73acfc3e3534f 100644 --- a/airflow-core/tests/unit/executors/test_workloads.py +++ b/airflow-core/tests/unit/executors/test_workloads.py @@ -235,19 +235,19 @@ def _make_mock_ti( return ti @pytest.mark.parametrize("execute_tasks_new_python_interpreter", [True, False, None]) - def test_make_populates_execute_tasks_new_python_interpreter(self, execute_tasks_new_python_interpreter): - from unittest.mock import Mock - + def test_make_uses_explicit_execute_tasks_new_python_interpreter( + self, execute_tasks_new_python_interpreter + ): ti = self._make_mock_ti(bundle_version="abc123", version_data={}) - ti.task = Mock(execute_tasks_new_python_interpreter=execute_tasks_new_python_interpreter) - workload = ExecuteTask.make(ti) + workload = ExecuteTask.make( + ti, execute_tasks_new_python_interpreter=execute_tasks_new_python_interpreter + ) assert workload.execute_tasks_new_python_interpreter is execute_tasks_new_python_interpreter def test_make_defaults_execute_tasks_new_python_interpreter_to_none(self): ti = self._make_mock_ti(bundle_version="abc123", version_data={}) - ti.task = None workload = ExecuteTask.make(ti) diff --git a/airflow-core/tests/unit/jobs/test_scheduler_job.py b/airflow-core/tests/unit/jobs/test_scheduler_job.py index 4add965a6abda..3b15f18be6ad4 100644 --- a/airflow-core/tests/unit/jobs/test_scheduler_job.py +++ b/airflow-core/tests/unit/jobs/test_scheduler_job.py @@ -3130,6 +3130,36 @@ def test_enqueue_task_instances_with_queued_state(self, dag_maker, session): assert mock_queue_workload.called session.rollback() + def test_enqueue_task_instances_sets_task_interpreter_override_from_serialized_task( + self, dag_maker, session + ): + from airflow.providers.standard.operators.python import PythonOperator + + dag_id = "SchedulerJobTest.test_enqueue_sets_task_interpreter_override" + task_id = "dummy" + session = settings.Session() + with dag_maker(dag_id=dag_id, start_date=DEFAULT_DATE, session=session): + PythonOperator( + task_id=task_id, + python_callable=lambda: None, + execute_tasks_new_python_interpreter=True, + ) + + scheduler_job = Job() + self.job_runner = SchedulerJobRunner(job=scheduler_job, executors=[self.null_exec]) + + dr = dag_maker.create_dagrun() + ti = dr.get_task_instance(task_id, session=session) + + with patch.object(BaseExecutor, "queue_workload") as mock_queue_workload: + self.job_runner._enqueue_task_instances_with_queued_state( + [ti], executor=self.job_runner.executor, session=session + ) + + workload = mock_queue_workload.call_args.args[0] + assert workload.execute_tasks_new_python_interpreter is True + session.rollback() + def test_executable_task_instances_to_queued_sets_external_executor_id(self, dag_maker, session): """external_executor_id is written to the DB in the same UPDATE that sets state=QUEUED.""" dag_id = "SchedulerJobTest.test_executable_sets_external_executor_id" From a37ce04de87acacfe14d4f759510ab3df29b6501 Mon Sep 17 00:00:00 2001 From: Vincent Hsiao <124506982+fat-catTW@users.noreply.github.com> Date: Fri, 7 Aug 2026 15:31:26 +0000 Subject: [PATCH 11/12] Trigger CI for task interpreter override PR From e3cb594a90ab985a7698f5c05fcdd6b0ff737b05 Mon Sep 17 00:00:00 2001 From: Vincent Hsiao <124506982+fat-catTW@users.noreply.github.com> Date: Sat, 8 Aug 2026 09:44:39 +0000 Subject: [PATCH 12/12] Restore task interpreter config fallback --- .../src/airflow/executors/base_executor.py | 1 - .../src/airflow/executors/workloads/task.py | 3 -- .../src/airflow/jobs/scheduler_job_runner.py | 10 ----- .../serialization/definitions/baseoperator.py | 1 - .../tests/unit/executors/test_workloads.py | 19 -------- .../tests/unit/jobs/test_scheduler_job.py | 30 ------------- .../edge3/worker_api/v2-edge-generated.yaml | 5 --- providers/standard/docs/operators/python.rst | 17 ------- .../providers/standard/operators/python.py | 30 +------------ .../unit/standard/operators/test_python.py | 44 +------------------ .../airflow/sdk/execution_time/coordinator.py | 3 -- .../airflow/sdk/execution_time/supervisor.py | 33 ++++---------- .../execution_time/test_supervisor.py | 18 +++----- 13 files changed, 17 insertions(+), 197 deletions(-) diff --git a/airflow-core/src/airflow/executors/base_executor.py b/airflow-core/src/airflow/executors/base_executor.py index 84acfe04c73ae..8030e9c14455e 100644 --- a/airflow-core/src/airflow/executors/base_executor.py +++ b/airflow-core/src/airflow/executors/base_executor.py @@ -724,7 +724,6 @@ def run_workload( log_path=workload.log_path, subprocess_logs_to_stdout=subprocess_logs_to_stdout, sentry_integration=getattr(workload, "sentry_integration", ""), - execute_tasks_new_python_interpreter=workload.execute_tasks_new_python_interpreter, ) if isinstance(workload, ExecuteCallback): from airflow.sdk.execution_time.callback_supervisor import supervise_callback diff --git a/airflow-core/src/airflow/executors/workloads/task.py b/airflow-core/src/airflow/executors/workloads/task.py index c661e3d6cded5..3099fe1d77485 100644 --- a/airflow-core/src/airflow/executors/workloads/task.py +++ b/airflow-core/src/airflow/executors/workloads/task.py @@ -67,7 +67,6 @@ class ExecuteTask(BaseDagBundleWorkload): ti: TaskInstanceDTO sentry_integration: str = "" - execute_tasks_new_python_interpreter: bool | None = None type: Literal["ExecuteTask"] = Field(init=False, default="ExecuteTask") @@ -97,7 +96,6 @@ def make( generator: JWTGenerator | None = None, bundle_info: BundleInfo | None = None, sentry_integration: str = "", - execute_tasks_new_python_interpreter: bool | None = None, ) -> ExecuteTask: """Create an ExecuteTask workload from a TaskInstance ORM model.""" from airflow.utils.helpers import log_filename_template_renderer @@ -124,5 +122,4 @@ def make( log_path=fname, bundle_info=bundle_info, sentry_integration=sentry_integration, - execute_tasks_new_python_interpreter=execute_tasks_new_python_interpreter, ) diff --git a/airflow-core/src/airflow/jobs/scheduler_job_runner.py b/airflow-core/src/airflow/jobs/scheduler_job_runner.py index 2b62aef1b6c70..a9efa1c03d62c 100644 --- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py +++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py @@ -1112,15 +1112,6 @@ def _get_sentry_integration(executor: BaseExecutor) -> str: return "" return sentry_integration - def _get_execute_tasks_new_python_interpreter(ti: TI) -> bool | None: - serialized_dag = self.scheduler_dag_bag.get_dag_for_run(dag_run=ti.dag_run, session=session) - if not serialized_dag or not serialized_dag.has_task(ti.task_id): - return None - value = getattr(serialized_dag.get_task(ti.task_id), "execute_tasks_new_python_interpreter", None) - if value is None or isinstance(value, bool): - return value - return None - # actually enqueue them for ti in task_instances: if ti.dag_run.state in State.finished_dr_states: @@ -1146,7 +1137,6 @@ def _get_execute_tasks_new_python_interpreter(ti: TI) -> bool | None: ti, generator=executor.jwt_generator, sentry_integration=_get_sentry_integration(executor), - execute_tasks_new_python_interpreter=_get_execute_tasks_new_python_interpreter(ti), ) executor.queue_workload(workload, session=session) diff --git a/airflow-core/src/airflow/serialization/definitions/baseoperator.py b/airflow-core/src/airflow/serialization/definitions/baseoperator.py index 0f72d2870b228..6bafc5891235a 100644 --- a/airflow-core/src/airflow/serialization/definitions/baseoperator.py +++ b/airflow-core/src/airflow/serialization/definitions/baseoperator.py @@ -192,7 +192,6 @@ def get_serialized_fields(cls): "execution_timeout", "executor", "executor_config", - "execute_tasks_new_python_interpreter", "ignore_first_depends_on_past", "inlets", "is_setup", diff --git a/airflow-core/tests/unit/executors/test_workloads.py b/airflow-core/tests/unit/executors/test_workloads.py index 73acfc3e3534f..37fbcd96ce950 100644 --- a/airflow-core/tests/unit/executors/test_workloads.py +++ b/airflow-core/tests/unit/executors/test_workloads.py @@ -234,25 +234,6 @@ def _make_mock_ti( return ti - @pytest.mark.parametrize("execute_tasks_new_python_interpreter", [True, False, None]) - def test_make_uses_explicit_execute_tasks_new_python_interpreter( - self, execute_tasks_new_python_interpreter - ): - ti = self._make_mock_ti(bundle_version="abc123", version_data={}) - - workload = ExecuteTask.make( - ti, execute_tasks_new_python_interpreter=execute_tasks_new_python_interpreter - ) - - assert workload.execute_tasks_new_python_interpreter is execute_tasks_new_python_interpreter - - def test_make_defaults_execute_tasks_new_python_interpreter_to_none(self): - ti = self._make_mock_ti(bundle_version="abc123", version_data={}) - - workload = ExecuteTask.make(ti) - - assert workload.execute_tasks_new_python_interpreter is None - def test_pinned_run_populates_version_data(self): """When the run is pinned, version_data from the run's created_dag_version flows to BundleInfo.""" version_data = {"schema_version": 1, "files": {"dags/my_dag.py": "ver123"}} diff --git a/airflow-core/tests/unit/jobs/test_scheduler_job.py b/airflow-core/tests/unit/jobs/test_scheduler_job.py index 3b15f18be6ad4..4add965a6abda 100644 --- a/airflow-core/tests/unit/jobs/test_scheduler_job.py +++ b/airflow-core/tests/unit/jobs/test_scheduler_job.py @@ -3130,36 +3130,6 @@ def test_enqueue_task_instances_with_queued_state(self, dag_maker, session): assert mock_queue_workload.called session.rollback() - def test_enqueue_task_instances_sets_task_interpreter_override_from_serialized_task( - self, dag_maker, session - ): - from airflow.providers.standard.operators.python import PythonOperator - - dag_id = "SchedulerJobTest.test_enqueue_sets_task_interpreter_override" - task_id = "dummy" - session = settings.Session() - with dag_maker(dag_id=dag_id, start_date=DEFAULT_DATE, session=session): - PythonOperator( - task_id=task_id, - python_callable=lambda: None, - execute_tasks_new_python_interpreter=True, - ) - - scheduler_job = Job() - self.job_runner = SchedulerJobRunner(job=scheduler_job, executors=[self.null_exec]) - - dr = dag_maker.create_dagrun() - ti = dr.get_task_instance(task_id, session=session) - - with patch.object(BaseExecutor, "queue_workload") as mock_queue_workload: - self.job_runner._enqueue_task_instances_with_queued_state( - [ti], executor=self.job_runner.executor, session=session - ) - - workload = mock_queue_workload.call_args.args[0] - assert workload.execute_tasks_new_python_interpreter is True - session.rollback() - def test_executable_task_instances_to_queued_sets_external_executor_id(self, dag_maker, session): """external_executor_id is written to the DB in the same UPDATE that sets state=QUEUED.""" dag_id = "SchedulerJobTest.test_executable_sets_external_executor_id" diff --git a/providers/edge3/src/airflow/providers/edge3/worker_api/v2-edge-generated.yaml b/providers/edge3/src/airflow/providers/edge3/worker_api/v2-edge-generated.yaml index 708171f3f5c66..13cacba97be36 100644 --- a/providers/edge3/src/airflow/providers/edge3/worker_api/v2-edge-generated.yaml +++ b/providers/edge3/src/airflow/providers/edge3/worker_api/v2-edge-generated.yaml @@ -1136,11 +1136,6 @@ components: type: string title: Sentry Integration default: '' - execute_tasks_new_python_interpreter: - anyOf: - - type: boolean - - type: 'null' - title: Execute Tasks New Python Interpreter type: type: string const: ExecuteTask diff --git a/providers/standard/docs/operators/python.rst b/providers/standard/docs/operators/python.rst index 2b1a04d3c4480..30f7151e55bef 100644 --- a/providers/standard/docs/operators/python.rst +++ b/providers/standard/docs/operators/python.rst @@ -47,23 +47,6 @@ Use the :class:`~airflow.providers.standard.operators.python.PythonOperator` to :start-after: [START howto_operator_python] :end-before: [END howto_operator_python] -Running selected tasks in a new interpreter -^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ - -By default, task process creation follows the global -:ref:`execute_tasks_new_python_interpreter ` -configuration. Set ``execute_tasks_new_python_interpreter`` on a ``PythonOperator`` task to override -that global setting for a single task. This is useful when one task needs a fresh Python interpreter, -for example to pick up plugin changes immediately, without changing the execution mode for every task. - -.. code-block:: python - - PythonOperator( - task_id="refresh_plugin_state", - python_callable=refresh_plugin_state, - execute_tasks_new_python_interpreter=True, - ) - Passing in arguments ^^^^^^^^^^^^^^^^^^^^ diff --git a/providers/standard/src/airflow/providers/standard/operators/python.py b/providers/standard/src/airflow/providers/standard/operators/python.py index 6f6dad20bc6b1..8698ddc403b60 100644 --- a/providers/standard/src/airflow/providers/standard/operators/python.py +++ b/providers/standard/src/airflow/providers/standard/operators/python.py @@ -35,7 +35,7 @@ from itertools import chain from pathlib import Path from tempfile import TemporaryDirectory -from typing import TYPE_CHECKING, Any, ClassVar, NamedTuple, cast +from typing import TYPE_CHECKING, Any, NamedTuple, cast import lazy_object_proxy from packaging.requirements import InvalidRequirement, Requirement @@ -76,19 +76,6 @@ log = logging.getLogger(__name__) - -def has_execute_tasks_new_python_interpreter_serialization_support() -> bool: - try: - from airflow.serialization.definitions.baseoperator import SerializedBaseOperator - except ImportError: - from airflow.serialization.serialized_objects import SerializedBaseOperator - - get_serialized_fields = getattr(SerializedBaseOperator, "get_serialized_fields", None) - if get_serialized_fields is None: - return False - return "execute_tasks_new_python_interpreter" in get_serialized_fields() - - if TYPE_CHECKING: from typing import Literal @@ -182,8 +169,6 @@ def my_python_callable(**kwargs): logs. Defaults to True, which allows return value log output. It can be set to False to prevent log output of return value when you return huge data such as transmission a large amount of XCom to TaskAPI. - :param execute_tasks_new_python_interpreter: If set, overrides the global - ``[core] execute_tasks_new_python_interpreter`` setting for this task. """ template_fields: Sequence[str] = ("templates_dict", "op_args", "op_kwargs") @@ -191,8 +176,6 @@ def my_python_callable(**kwargs): BLUE = "#ffefeb" ui_color = BLUE - __serialized_fields: ClassVar[frozenset[str] | None] = None - # since we won't mutate the arguments, we should just do the shallow copy # there are some cases we can't deepcopy the objects(e.g protobuf). shallow_copy_attrs: Sequence[str] = ("python_callable", "op_kwargs") @@ -206,7 +189,6 @@ def __init__( templates_dict: dict[str, Any] | None = None, templates_exts: Sequence[str] | None = None, show_return_value_in_logs: bool = True, - execute_tasks_new_python_interpreter: bool | None = None, **kwargs, ) -> None: super().__init__(**kwargs) @@ -219,16 +201,6 @@ def __init__( if templates_exts: self.template_ext = templates_exts self.show_return_value_in_logs = show_return_value_in_logs - self.execute_tasks_new_python_interpreter = execute_tasks_new_python_interpreter - - @classmethod - def get_serialized_fields(cls): - if not cls.__serialized_fields: - serialized_fields = super().get_serialized_fields() - if has_execute_tasks_new_python_interpreter_serialization_support(): - serialized_fields = serialized_fields | {"execute_tasks_new_python_interpreter"} - cls.__serialized_fields = frozenset(serialized_fields) - return cls.__serialized_fields @property def is_async(self) -> bool: diff --git a/providers/standard/tests/unit/standard/operators/test_python.py b/providers/standard/tests/unit/standard/operators/test_python.py index 52cc46ff6f0d4..d8228b1032f30 100644 --- a/providers/standard/tests/unit/standard/operators/test_python.py +++ b/providers/standard/tests/unit/standard/operators/test_python.py @@ -64,7 +64,7 @@ from airflow.utils.state import DagRunState, State, TaskInstanceState from airflow.utils.types import DagRunType -from tests_common.test_utils.compat import SerializedBaseOperator, TriggerRule, timezone +from tests_common.test_utils.compat import TriggerRule, timezone from tests_common.test_utils.db import clear_db_runs from tests_common.test_utils.in_process_taskrun import pushed_xcom, run_task_no_db from tests_common.test_utils.taskinstance import get_template_context, run_task_instance @@ -89,13 +89,6 @@ from airflow.sdk import Context -def has_execute_tasks_new_python_interpreter_serialization_support() -> bool: - get_serialized_fields = getattr(SerializedBaseOperator, "get_serialized_fields", None) - if get_serialized_fields is None: - return False - return "execute_tasks_new_python_interpreter" in get_serialized_fields() - - AIRFLOW_ROOT_PATH = Path(__file__).parents[6] TI = TaskInstance @@ -298,41 +291,6 @@ def a_fn(): rendered_op_kwargs = task.op_kwargs assert rendered_op_kwargs["a_callable"] == a_fn - @pytest.mark.skipif( - not has_execute_tasks_new_python_interpreter_serialization_support(), - reason="Current Airflow core does not serialize execute_tasks_new_python_interpreter", - ) - @pytest.mark.parametrize("execute_tasks_new_python_interpreter", [True, False, None]) - def test_execute_tasks_new_python_interpreter_is_serialized(self, execute_tasks_new_python_interpreter): - from tests_common.test_utils.compat import OperatorSerialization - - task = PythonOperator( - python_callable=lambda: None, - task_id=self.task_id, - execute_tasks_new_python_interpreter=execute_tasks_new_python_interpreter, - ) - - serialized = OperatorSerialization.serialize_operator(task) - deserialized = OperatorSerialization.deserialize_operator(serialized) - - assert "execute_tasks_new_python_interpreter" in PythonOperator.get_serialized_fields() - assert serialized.get("execute_tasks_new_python_interpreter") is execute_tasks_new_python_interpreter - assert ( - getattr(deserialized, "execute_tasks_new_python_interpreter", None) - is execute_tasks_new_python_interpreter - ) - - @mock.patch( - "airflow.providers.standard.operators.python.has_execute_tasks_new_python_interpreter_serialization_support", - return_value=False, - ) - def test_execute_tasks_new_python_interpreter_serialized_field_requires_core_support(self, _): - PythonOperator._PythonOperator__serialized_fields = None - try: - assert "execute_tasks_new_python_interpreter" not in PythonOperator.get_serialized_fields() - finally: - PythonOperator._PythonOperator__serialized_fields = None - def test_python_operator_shallow_copy_attr(self): def not_callable(x): raise RuntimeError("Should not be triggered") diff --git a/task-sdk/src/airflow/sdk/execution_time/coordinator.py b/task-sdk/src/airflow/sdk/execution_time/coordinator.py index 7175a1f929ad2..1469b27888c51 100644 --- a/task-sdk/src/airflow/sdk/execution_time/coordinator.py +++ b/task-sdk/src/airflow/sdk/execution_time/coordinator.py @@ -98,7 +98,6 @@ def execute_task( logger: FilteringBoundLogger | None = None, sentry_integration: str = "", subprocess_logs_to_stdout: bool, - execute_tasks_new_python_interpreter: bool | None = None, **kwargs, ) -> ExecutionResult: """ @@ -172,7 +171,6 @@ def execute_task( logger: FilteringBoundLogger | None = None, sentry_integration: str = "", subprocess_logs_to_stdout: bool, - execute_tasks_new_python_interpreter: bool | None = None, **kwargs, ) -> BaseCoordinator.ExecutionResult: # TODO: Importing this at the top causes circular imports. @@ -193,7 +191,6 @@ def execute_task( logger=logger, bundle_info=bundle_info, subprocess_logs_to_stdout=subprocess_logs_to_stdout, - execute_tasks_new_python_interpreter=execute_tasks_new_python_interpreter, sentry_integration=sentry_integration, ) exit_code = process.wait() diff --git a/task-sdk/src/airflow/sdk/execution_time/supervisor.py b/task-sdk/src/airflow/sdk/execution_time/supervisor.py index c2344c3bd5824..27133bc6d4905 100644 --- a/task-sdk/src/airflow/sdk/execution_time/supervisor.py +++ b/task-sdk/src/airflow/sdk/execution_time/supervisor.py @@ -515,19 +515,10 @@ def exit(n: int) -> NoReturn: def _should_use_exec() -> bool: - """Whether forked children should ``exec`` a fresh interpreter on this platform.""" - return sys.platform in _FORK_EXEC_PLATFORMS - - -def _should_use_exec_for_task(execute_tasks_new_python_interpreter: bool | None = None) -> bool: - """Whether task execution should ``exec`` a fresh Python interpreter.""" - if sys.platform in _FORK_EXEC_PLATFORMS: - return True - if execute_tasks_new_python_interpreter is None: - execute_tasks_new_python_interpreter = conf.getboolean( - "core", "execute_tasks_new_python_interpreter", fallback=False - ) - return execute_tasks_new_python_interpreter + """Whether forked children should ``exec`` a fresh interpreter.""" + return sys.platform in _FORK_EXEC_PLATFORMS or conf.getboolean( + "core", "execute_tasks_new_python_interpreter", fallback=False + ) def _resolve_child_target(dotted: str) -> Callable[[], None]: @@ -699,7 +690,8 @@ def start( """ Fork and start a new subprocess with the specified target function. - :param use_exec: If True, immediately ``os.execv`` a fresh Python interpreter after ``os.fork``. + :param use_exec: If True, on platforms that need it (currently macOS), + immediately ``os.execv`` a fresh Python interpreter after ``os.fork``. This avoids macOS fork-safety issues with Objective-C frameworks. ``target`` is rehydrated in the exec'd child from its ``module:qualname``, so any importable entry point (task execution, DAG processor, triggerer) @@ -1362,16 +1354,13 @@ def start( # type: ignore[override] target: Callable[[], None] = _subprocess_main, logger: FilteringBoundLogger | None = None, sentry_integration: str = "", - execute_tasks_new_python_interpreter: bool | None = None, **kwargs, ) -> Self: """Fork and start a new subprocess to execute the given task.""" + # Opt in to fork+exec on platforms that need it (currently macOS). # Tests override `target` with a local stub to exercise the base - # infrastructure; keep bare fork for those. Otherwise task execution - # can opt in to fork+exec through platform safety or configuration. - use_exec = target is _subprocess_main and _should_use_exec_for_task( - execute_tasks_new_python_interpreter - ) + # infrastructure; keep bare fork for those. + use_exec = target is _subprocess_main and _should_use_exec() proc: Self = super().start( id=what.id, client=client, target=target, logger=logger, use_exec=use_exec, **kwargs ) @@ -2467,7 +2456,6 @@ def supervise_task( subprocess_logs_to_stdout: bool = False, client: Client | None = None, sentry_integration: str = "", - execute_tasks_new_python_interpreter: bool | None = None, ) -> int: """ Run a single task execution to completion. @@ -2483,8 +2471,6 @@ def supervise_task( :param client: Optional preconfigured client for communication with the server (Mostly for tests). :param sentry_integration: If the executor has a Sentry integration, import path to a callable to initialize it (empty means no integration). - :param execute_tasks_new_python_interpreter: Override ``[core] execute_tasks_new_python_interpreter`` - for this task. ``None`` keeps the global configuration. :return: Exit code of the process. :raises ValueError: If server URL is empty or invalid. :raises InvalidCoordinatorError: If the coordinator for the task is not @@ -2567,7 +2553,6 @@ def supervise_task( logger=logger, sentry_integration=sentry_integration, subprocess_logs_to_stdout=subprocess_logs_to_stdout, - execute_tasks_new_python_interpreter=execute_tasks_new_python_interpreter, ) end = time.monotonic() log.info( diff --git a/task-sdk/tests/task_sdk/execution_time/test_supervisor.py b/task-sdk/tests/task_sdk/execution_time/test_supervisor.py index 1de7737f5765e..9bf7c2428f6b9 100644 --- a/task-sdk/tests/task_sdk/execution_time/test_supervisor.py +++ b/task-sdk/tests/task_sdk/execution_time/test_supervisor.py @@ -186,20 +186,14 @@ @pytest.mark.parametrize( - ("task_override", "global_config", "expected"), - [(True, False, True), (False, True, False), (None, True, True), (None, False, False)], + ("platform", "config_value", "expected"), + [("darwin", False, True), ("linux", True, True), ("linux", False, False)], ) -def test_should_use_exec_for_task_honors_task_override(monkeypatch, task_override, global_config, expected): - monkeypatch.setattr(supervisor.sys, "platform", "linux") - monkeypatch.setattr(supervisor.conf, "getboolean", lambda *args, **kwargs: global_config) +def test_should_use_exec_honors_platform_and_config(monkeypatch, platform, config_value, expected): + monkeypatch.setattr(supervisor.sys, "platform", platform) + monkeypatch.setattr(supervisor.conf, "getboolean", lambda *args, **kwargs: config_value) - assert supervisor._should_use_exec_for_task(task_override) is expected - - -def test_should_use_exec_for_task_keeps_platform_exec_requirement(monkeypatch): - monkeypatch.setattr(supervisor.sys, "platform", "darwin") - - assert supervisor._should_use_exec_for_task(False) is True + assert supervisor._should_use_exec() is expected def lineno():