Skip to content

Async support for /task-instances/states and /task-instances/count API routes - #73966

Open
zach-overflow wants to merge 4 commits into
apache:mainfrom
DataDog:zach.gottesman/async-get-ti-states
Open

zach-overflow wants to merge 4 commits into
apache:mainfrom
DataDog:zach.gottesman/async-get-ti-states

Conversation

@zach-overflow

Copy link
Copy Markdown
Contributor

What?

  • Converts the following two routes functions to async coroutines using the AsyncSessionDep
    • /task-instances/state (get_task_instance_states)
    • /task-instances/count (get_task_instance_count)

Why?


Was generative AI tooling used to co-author this PR?
  • Yes

Generated-by: Codex GPT-Astra following the guidelines

PR description written by human, some manual de-slopification of docstrings, test structure also done by hand.


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

@boring-cyborg boring-cyborg Bot added area:API Airflow's REST/HTTP API area:task-sdk labels Sep 30, 2026
@zach-overflow zach-overflow changed the title Zach.gottesman/async get ti states zo/async get ti states Sep 30, 2026
@zach-overflow zach-overflow changed the title zo/async get ti states Async support for /task-instances/states and /task-instances/count API routes Sep 30, 2026
@zach-overflow
zach-overflow marked this pull request as ready for review September 30, 2026 16:24

@Dev-iL Dev-iL left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The migration correctly reuses Airflow’s async session infrastructure, explicitly detaches serialized Dag rows before threaded deserialization, and preserves the tested endpoint contracts across SQLite, PostgreSQL, and MySQL.

The main comment is related to a performance regression (see evidence in the section below). I'll update the migration skill to be careful about similar mechanical replacements.

Breeze verification by Codex

Head: 33c54041cb54f640975630db239a26e0ea1e5e21.
Comparison base (merge base): 508301ab263dc6df7a492b8f284229259c889fff.
Current target branch SHA when fetched: f7a276337416bc3ba97a7bfb5836144568d46410.
The final read of GitHub confirmed the head was unchanged: final-head.json.

Review checkout: /tmp/airflow-review-73966.
Base checkout: /tmp/airflow-review-73966-base.
Tests ran through Breeze, using the cached Python 3.10 CI image and the checkout's mounted sources. Pytest reported Python 3.10.21 and function-scoped asyncio fixtures. No production source file was changed. Temporary instrumentation was confined to dev/review_73966.py in those review checkouts; its final source is retained as verification-harness.py.

Contract and helper tests

Revision Backend Async driver Scope Result Log
Head SQLite, initially selected default aiosqlite Both changed endpoint classes and full common/model DagBag test files 85 passed SQLite head
Head PostgreSQL 14 psycopg async dialect Same scope 85 passed PostgreSQL head
Head MySQL 8.0 aiomysql Same scope 85 passed MySQL head
Base SQLite, initially selected default aiosqlite Both changed endpoint classes 40 passed SQLite base
Head SQLite, explicitly selected aiosqlite Both classes × API versions 2025-04-11 and 2026-10-30 × compression false/true 160 passed SQLite matrix
Head PostgreSQL 14 psycopg async dialect Same version/compression matrix 160 passed PostgreSQL matrix
Head MySQL 8.0 aiomysql Same version/compression matrix 160 passed MySQL matrix

The configured drivers were traced in settings.py; PostgreSQL runtime introspection reported postgresql psycopg psycopg, MySQL introspection reported mysqldb aiomysql, and the explicitly selected SQLite probes reported sqlite aiosqlite. The PostgreSQL result used the configured psycopg asynchronous implementation, rather than assuming asyncpg. Database drivers/configuration were not changed by this PR; PgBouncer and driver overrides were not tested.

Reproduce the 85-test head suite from the head checkout, substituting sqlite, postgres, or mysql:

breeze run --backend sqlite --skip-image-upgrade-check --answer no pytest \
  airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py::TestGetCount \
  airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py::TestGetTaskStates \
  airflow-core/tests/unit/api_fastapi/common/test_dagbag.py \
  airflow-core/tests/unit/models/test_dagbag.py -q

The original SQLite invocation omitted --backend after Breeze initialized its fresh worktree with SQLite. PostgreSQL and MySQL were explicitly selected. The 40-test base command uses only the first two test targets above.

The version/compression matrix invokes the existing tests without changing their assertions. The temporary pytest plugin sets the requested Airflow-API-Version and compression setting before test Dag creation:

breeze run --backend sqlite --skip-image-upgrade-check --answer no \
  python dev/review_73966.py versions

Only the oldest and newest negotiated versions, plus the unversioned head tests, were exercised; intermediate version headers were not separately tested. Existing count/state filters, map indexes, task groups, missing Dag/group responses, and response shapes passed within this scope.

Ruff lint and formatting checks passed on all six changed Python files through Breeze: see the end of postgres-extra-verification.log. This was a review, not a full CI run: no mypy or complete prek suite was run.

Finding 1: event-loop stalls during full TaskInstance loading

The harness creates a real Dag run with one ordinary task instance and 1,000 or 5,000 mapped instances of the same task, with no inflated executor configuration or Dag-run payload. It sends actual HTTP requests to /execution/task-instances/states with both dag_id and run_ids, verifies HTTP 200 and the returned state count, and measures a 1 ms ticker running on the same TestClient event loop. The existing test authentication override remains in use.

Each dataset receives three sequential requests. The first request for each dataset is excluded from the warm comparison. Cold calls also show substantial stalls at the base revision, so they are not used as PR-specific evidence. These are small local samples, not a production throughput benchmark.

Revision/query Returned states Warm ticker gaps (ms) Warm HTTP duration (ms) Warm awaited query duration (ms)
Base, synchronous handler 1,001 15.20 / 11.92 228.27 / 289.30 Not instrumented
Head, first reproduction 1,001 158.85 / 109.94 189.49 / 139.37 Not instrumented
Base, synchronous handler 5,001 15.24 / 12.46 1,051.24 / 1,226.62 Not instrumented
Head, first reproduction 5,001 638.35 / 758.60 720.87 / 868.52 Not instrumented
Head, instrumented full-entity query 5,001 754.84 / 693.19 903.55 / 842.09 816.52 / 766.61
Head, temporary four-column projection 5,001 36.26 / 37.52 56.19 / 62.71 17.50 / 20.96

Logs: base probe, first head probe, head/projection comparison.

The source-level mechanism is the full select(TI) at task_instances.py:1325, executed inside the async handler at line 1339. TaskInstance.dag_run has lazy="joined"; ORM result processing loads the entire TaskInstance/DagRun graph, although the response consumes only run_id, task_id, map_index, and state. Moving the handler from a synchronous FastAPI worker to the event loop also moves this CPU work there. Awaiting database I/O does not offload ORM result processing.

The projection experiment replaces only TaskInstance statements inside the temporary AsyncSession.scalars wrapper with statement.with_only_columns(...), awaits execute on the same session, and supplies buffered row objects to the original route. Production code is untouched. It tests a mitigation for this reproduction; it is not a completed production patch and does not independently establish compatibility for every filter/group combination. _get_group_tasks also selects complete TaskInstances, so the draft asks that its results use the same required column set.

The measured regression is event-loop availability. The original head route returned faster than the base in the initial comparison; the finding does not claim that the endpoint itself got slower.

# Run from the base checkout and from the head checkout.
breeze run --backend sqlite --skip-image-upgrade-check --answer no \
  python dev/review_73966.py states

# Run from the head checkout for the local projection experiment.
breeze run --backend sqlite --skip-image-upgrade-check --answer no \
  python dev/review_73966.py projected-states

The base copy of the harness predates query-duration/projection instrumentation; its ticker and HTTP request construction match the initial head reproduction.

Finding 2: committed private workspace instructions

The end of sqlite-states-comparison.log contains the actual committed CLAUDE.local.md and AGENTS.local.md, read through Breeze using git show HEAD:<path> with the mounted main Git directory. It verifies:

  • CLAUDE.local.md:20 imports AGENTS.local.md.
  • AGENTS.local.md:22–27 declares a Workforest and a fixed author-specific branch.
  • AGENTS.local.md:31–34 requires branch registration with forest.
  • command -v forest fails in the standard Breeze CI image, producing forest-unavailable-in-breeze.

The applicability concern follows from those committed instructions: their workspace facts do not describe arbitrary Airflow checkouts. The draft asks to remove the two private setup files. The tool-availability observation is limited to the tested Breeze image.

Investigated leads and setup limitations

The detached serialized-Dag row contains the scalar columns _read_dag accesses. Existing model tests validate worker execution, detached state, compressed/uncompressed deserialization, missing rows, and worker exceptions; the version/compression matrix also exercises actual database-backed group routes. Cache writes in the API's CachedDBDagBag are protected by an existing RLock. These leads did not produce another finding.

Comment thread airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py Outdated
Comment thread AGENTS.local.md Outdated
@zach-overflow
zach-overflow force-pushed the zach.gottesman/async-get-ti-states branch from 33c5404 to b9b8f90 Compare October 1, 2026 13:29
@zach-overflow
zach-overflow requested a review from Dev-iL October 1, 2026 13:42
@zach-overflow
zach-overflow force-pushed the zach.gottesman/async-get-ti-states branch from b9b8f90 to 026f787 Compare October 1, 2026 15:14

@Dev-iL Dev-iL left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Please see #73966 (comment)

Task-group lookups can decompress and deserialize serialized Dags and
initialize operator links. This work must stay off the API event loop
when database reads use an async session.
The async states route ran ORM entity hydration for every matching task
instance and its eagerly joined Dag run on the API event loop, stalling
other requests on large mapped runs. Only four columns are consumed.

Per-developer agent instruction files must never be committed.
Building the response row by row and processing the full result in one
piece still blocked the event loop longer than the old synchronous
handler on large mapped runs. The Dag-wide query also ran when a task
group without task IDs was requested, although its result was discarded.
@zach-overflow
zach-overflow force-pushed the zach.gottesman/async-get-ti-states branch from be60f19 to 66be871 Compare October 5, 2026 13:05
Pooled async connections are bound to the event loop that opened them.
The in-process Execution API used by dag.test() starts its own loop but
shared the async pool with earlier loops, so the first async route
request could pick up a stale connection and fail. On MySQL this broke
mapped task groups, whose upstream map-index resolution now goes through
the async task-instance count route.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:API Airflow's REST/HTTP API area:task-sdk

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants