This is an automated email from the ASF dual-hosted git repository.
ashb 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 8d815efdbc3 Own the BaseAsyncOperator event loop through
asyncio.Runner (#73560)
8d815efdbc3 is described below
commit 8d815efdbc3dad1ad5cdf592934d1b4de7458164
Author: David Blain <[email protected]>
AuthorDate: Tue Sep 29 15:49:20 2026 +0200
Own the BaseAsyncOperator event loop through asyncio.Runner (#73560)
Every task that runs through BaseAsyncOperator.execute logs
DeprecationWarning: There is no current event loop as its first line on Python
3.12 and 3.13:
```
[2026-09-22 16:34:41] INFO - Worker startup parse complete
bundle_name=dags-folder bundle_version= dag_file=test_msgraph.py
bundle_prepare_ms=1 dag_file_parse_ms=17
[2026-09-22 16:34:41] WARNING - There is no current event loop
category=DeprecationWarning
```
The task-runner process starts without an event loop set, and the
event_loop() helper relied on asyncio.get_event_loop() to create one. That is
exactly the case Python deprecated in 3.12, and Python 3.14 raises RuntimeError
there instead of creating a loop. The helper also never closed a loop that
asyncio created implicitly on its behalf, so that loop leaked for the rest of
the process.
event_loop() now owns its loop through asyncio.Runner, as the asyncio docs
recommend in place of policies. Plain asyncio.run() does not cover the helper's
purpose: it creates a loop for one coroutine and closes it, while the helper
hands out a loop that outlives a single coroutine so callers can drive it
through several run_until_complete() calls (the upcoming iterable operator does
exactly that between batches of sub-tasks). Runner gives that, and on exit also
cancels leftover tasks [...]
---
task-sdk/src/airflow/sdk/bases/operator.py | 69 ++++++++--
task-sdk/tests/task_sdk/bases/test_operator.py | 171 +++++++++++++++++++++++++
2 files changed, 226 insertions(+), 14 deletions(-)
diff --git a/task-sdk/src/airflow/sdk/bases/operator.py
b/task-sdk/src/airflow/sdk/bases/operator.py
index 940fee84b87..7634443b119 100644
--- a/task-sdk/src/airflow/sdk/bases/operator.py
+++ b/task-sdk/src/airflow/sdk/bases/operator.py
@@ -200,25 +200,66 @@ def coerce_resources(resources: dict[str, Any] | None) ->
Resources | None:
return Resources(**resources)
+def _cancel_all_tasks(loop: AbstractEventLoop) -> None:
+ """Cancel every task still pending on ``loop`` and let it run its cleanup,
as ``asyncio.run()`` does."""
+ to_cancel = asyncio.all_tasks(loop)
+ if not to_cancel:
+ return
+ for task in to_cancel:
+ task.cancel()
+ loop.run_until_complete(asyncio.gather(*to_cancel, return_exceptions=True))
+ for task in to_cancel:
+ if not task.cancelled() and task.exception() is not None:
+ loop.call_exception_handler(
+ {
+ "message": "unhandled exception during event_loop()
shutdown",
+ "exception": task.exception(),
+ "task": task,
+ }
+ )
+
+
@contextlib.contextmanager
def event_loop() -> Generator[AbstractEventLoop]:
- new_event_loop = False
- loop = None
+ """
+ Own an event loop for the duration of the block, to drive coroutines from
synchronous code.
+
+ Unlike ``asyncio.run()``, which creates a loop for one coroutine and
closes it, the loop yielded here
+ outlives any number of ``run_until_complete()`` calls: tasks created on it
in one call are still there
+ for the next, which is what code that pauses and resumes a loop between
batches of work needs.
+
+ On Python 3.11+ this is :class:`asyncio.Runner`; Python 3.10 has no
``Runner`` and gets a fallback that
+ replicates ``Runner.close()``. In both cases the loop is set as the
current loop of the thread while the
+ block runs and unset afterwards, leftover tasks are cancelled, async
generators and the default executor
+ are shut down and the loop is closed. Nothing goes through
``asyncio.get_event_loop()``, which since
+ Python 3.12 warns and since 3.14 raises when no loop is set. Entering the
block from a running loop is
+ refused: a loop cannot drive another loop on the same thread.
+ """
+ try:
+ asyncio.get_running_loop()
+ except RuntimeError:
+ pass
+ else:
+ raise RuntimeError("event_loop() cannot be used from a running event
loop")
+
+ if Runner := getattr(asyncio, "Runner", None):
+ with Runner() as runner:
+ yield runner.get_loop()
+ return
+
+ # Python 3.10 has no asyncio.Runner: replicate what Runner.close() does.
Drop with Python 3.10 support.
+ loop = asyncio.new_event_loop()
+ asyncio.set_event_loop(loop)
try:
- try:
- loop = asyncio.get_event_loop()
- if loop.is_closed():
- raise RuntimeError
- except RuntimeError:
- loop = asyncio.new_event_loop()
- asyncio.set_event_loop(loop)
- new_event_loop = True
yield loop
finally:
- if new_event_loop and loop is not None:
- with contextlib.suppress(AttributeError):
- loop.close()
- asyncio.set_event_loop(None)
+ try:
+ _cancel_all_tasks(loop)
+ loop.run_until_complete(loop.shutdown_asyncgens())
+ loop.run_until_complete(loop.shutdown_default_executor())
+ finally:
+ asyncio.set_event_loop(None)
+ loop.close()
class _PartialDescriptor:
diff --git a/task-sdk/tests/task_sdk/bases/test_operator.py
b/task-sdk/tests/task_sdk/bases/test_operator.py
index dcb5240a83d..9041e404b12 100644
--- a/task-sdk/tests/task_sdk/bases/test_operator.py
+++ b/task-sdk/tests/task_sdk/bases/test_operator.py
@@ -17,6 +17,7 @@
from __future__ import annotations
+import asyncio
import copy
import logging
import uuid
@@ -32,12 +33,14 @@ import structlog
from airflow.sdk import DAG, Label, TaskGroup, task as task_decorator
from airflow.sdk._shared.secrets_masker import _secrets_masker, mask_secret
from airflow.sdk.bases.operator import (
+ BaseAsyncOperator,
BaseOperator,
BaseOperatorMeta,
ExecutorSafeguard,
chain,
chain_linear,
cross_downstream,
+ event_loop,
)
from airflow.sdk.definitions.param import ParamsDict
from airflow.sdk.definitions.template import literal
@@ -1137,3 +1140,171 @@ def test_partial_default_args():
assert op.arg2 == "b"
assert op.arg3 == 3
assert op.queue == "THIS"
+
+
[email protected](params=["asyncio.Runner", "fallback"])
+def event_loop_runner(request, monkeypatch):
+ """
+ Run a test against both ways ``event_loop()`` can own its loop.
+
+ ``asyncio.Runner`` exists from Python 3.11; the fallback replicates it for
Python 3.10 and is
+ exercised on every Python by hiding ``asyncio.Runner`` for the test.
+ """
+ if request.param == "fallback":
+ monkeypatch.delattr(asyncio, "Runner", raising=False)
+ elif not hasattr(asyncio, "Runner"):
+ pytest.skip("asyncio.Runner needs Python 3.11+")
+ return request.param
+
+
[email protected]
+def fresh_process_loop_state():
+ """
+ Mimic a freshly started task-runner process: no loop set and
``set_event_loop()`` never called.
+
+ In that state ``asyncio.get_event_loop()`` emits ``DeprecationWarning:
There is no current event
+ loop`` when it has to create a loop (Python 3.12+) or raises (Python
3.14+).
+ """
+ with warnings.catch_warnings():
+ # The policy API is deprecated on Python 3.14; it is still the only
way to reset this state.
+ warnings.simplefilter("ignore", DeprecationWarning)
+ previous = asyncio.get_event_loop_policy()
+ asyncio.set_event_loop_policy(asyncio.DefaultEventLoopPolicy())
+ try:
+ yield
+ finally:
+ asyncio.set_event_loop_policy(previous)
+
+
+def _deprecation_warnings(caught: list[warnings.WarningMessage]) -> list[str]:
+ return [str(w.message) for w in caught if issubclass(w.category,
DeprecationWarning)]
+
+
[email protected]("event_loop_runner")
+class TestBaseAsyncOperator:
+ class AsyncOperator(BaseAsyncOperator):
+ async def aexecute(self, context):
+ await asyncio.sleep(0)
+ return "done"
+
+ def test_is_async(self):
+ assert self.AsyncOperator(task_id="async_op").is_async is True
+
+ def test_execute_requires_aexecute(self):
+ with pytest.raises(NotImplementedError):
+ BaseAsyncOperator(task_id="async_op").execute({})
+
+ @pytest.mark.usefixtures("fresh_process_loop_state")
+ @pytest.mark.parametrize("execution_timeout", [None,
timedelta(seconds=5)], ids=["no-timeout", "timeout"])
+ def test_execute_runs_aexecute_without_deprecation_warning(self,
execution_timeout):
+ op = self.AsyncOperator(task_id="async_op",
execution_timeout=execution_timeout)
+
+ with warnings.catch_warnings(record=True) as caught:
+ warnings.simplefilter("always")
+ assert op.execute({}) == "done"
+
+ assert _deprecation_warnings(caught) == []
+ # The loop execute() ran on is closed and not left behind as the
thread's current loop.
+ with pytest.raises(RuntimeError):
+ asyncio.get_event_loop()
+
+ def test_execute_enforces_execution_timeout(self):
+ class SlowOperator(BaseAsyncOperator):
+ async def aexecute(self, context):
+ await asyncio.sleep(60)
+
+ op = SlowOperator(task_id="slow_op",
execution_timeout=timedelta(milliseconds=50))
+
+ with pytest.raises(asyncio.TimeoutError):
+ op.execute({})
+
+
[email protected]("event_loop_runner")
+class TestEventLoop:
+ """``event_loop()`` hands synchronous code a loop it owns, through
``asyncio.Runner`` where Python has it."""
+
+ @pytest.mark.usefixtures("fresh_process_loop_state")
+ def test_creates_and_closes_own_loop_without_deprecation_warning(self):
+ with warnings.catch_warnings(record=True) as caught:
+ warnings.simplefilter("always")
+ with event_loop() as loop:
+ assert not loop.is_closed()
+ assert loop.run_until_complete(asyncio.sleep(0, result="ran"))
== "ran"
+
+ assert _deprecation_warnings(caught) == []
+ assert loop.is_closed()
+ # Nothing is left behind as the thread's current loop.
+ with pytest.raises(RuntimeError):
+ asyncio.get_event_loop()
+
+ @pytest.mark.usefixtures("fresh_process_loop_state")
+ def test_second_use_after_cleanup_still_works(self):
+ with event_loop() as first:
+ pass
+ with warnings.catch_warnings(record=True) as caught:
+ warnings.simplefilter("always")
+ with event_loop() as second:
+ assert second.run_until_complete(asyncio.sleep(0,
result="ran")) == "ran"
+
+ assert _deprecation_warnings(caught) == []
+ assert first.is_closed()
+ assert second.is_closed()
+
+ def test_loop_outlives_run_until_complete_calls(self):
+ """Unlike ``asyncio.run()``, work scheduled in one run is still there
for the next."""
+ with event_loop() as loop:
+ release = asyncio.Event()
+
+ async def wait_for_release():
+ await release.wait()
+ return "released"
+
+ pending = loop.create_task(wait_for_release())
+ assert loop.run_until_complete(asyncio.sleep(0, result="first
run")) == "first run"
+ assert not pending.done()
+
+ release.set()
+ assert loop.run_until_complete(pending) == "released"
+
+ def test_pending_tasks_are_cancelled_on_exit(self):
+ cleaned_up = False
+
+ async def never_finishes():
+ nonlocal cleaned_up
+ try:
+ await asyncio.sleep(3600)
+ except asyncio.CancelledError:
+ cleaned_up = True
+ raise
+
+ with event_loop() as loop:
+ pending = loop.create_task(never_finishes())
+ # Let the task start so cancellation reaches its body.
+ loop.run_until_complete(asyncio.sleep(0))
+
+ assert pending.cancelled()
+ assert cleaned_up
+ assert loop.is_closed()
+
+ def test_leaves_loop_already_set_open_but_does_not_reuse_it(self):
+ existing = asyncio.new_event_loop()
+ asyncio.set_event_loop(existing)
+ try:
+ with event_loop() as loop:
+ assert loop is not existing
+ assert loop.run_until_complete(asyncio.sleep(0, result="ran"))
== "ran"
+ assert loop.is_closed()
+ assert not existing.is_closed()
+ finally:
+ existing.close()
+ asyncio.set_event_loop(None)
+
+ @pytest.mark.asyncio
+ async def test_refuses_running_loop(self):
+ running = asyncio.get_running_loop()
+
+ with pytest.raises(RuntimeError, match="running event loop"):
+ with event_loop():
+ pass
+
+ assert not running.is_closed()