rodrigopimentel opened a new issue, #74239:
URL: https://github.com/apache/airflow/issues/74239

   ## Under which category would you file this issue?
   Airflow Core
   
   ## Apache Airflow version
   3.3.2 (the same code is on `main`)
   
   ## What happened and how to reproduce it?
   
   **Issue Description**
   
   Two response models of the public API query the metadata database during 
response serialization, and FastAPI serializes on the event loop. For a sync 
route, `serialize_response` validates in the threadpool but calls the 
serializer directly (`fastapi/routing.py` at 0.136.3, lines 292-306), so these 
queries run on the loop thread:
   
   - `DAGDetailsResponse.latest_dag_version` (`core_api/datamodels/dags.py`) is 
a `@computed_field` that calls `DagVersion.get_latest_version(...)` on a new 
session (`@provide_session`).
   - `PoolCollectionResponse.pools` (`core_api/datamodels/pools.py`) is typed 
`Iterable[PoolResponse]`. Pydantic validates `Iterable` fields lazily, during 
iteration 
(https://pydantic.dev/docs/validation/latest/api/pydantic/standard_library_types).
 So the `BeforeValidator(_call_function)` fields of each `PoolResponse` 
(`occupied_slots`, `running_slots`, `queued_slots`, `scheduled_slots`, 
`open_slots`, `deferred_slots`), each a `@provide_session` query, run inside 
`dump_json`, six per pool.
   
   With a free connection this costs a few milliseconds of loop time per 
request. When the SQLAlchemy pool is exhausted, the loop blocks in the pool 
checkout. The connections it waits for are held by `SessionDep` sessions of 
other requests, and their teardown is scheduled by the loop, so nothing is 
returned until `pool_timeout` (30 s) expires. During that time the process 
serves nothing: no UI requests, no Execution API calls, no health probe.
   
   The DAG list page sends one `/api/v2/dags/{dag_id}/details` request per row, 
so a single page load is enough to exhaust a small pool. In our deployment (two 
api-server replicas behind a load balancer, metadata DB behind PgBouncer), one 
page load sent about 30 parallel `/details` requests. Both replicas stopped 
answering for up to 210 s, the Helm chart's liveness probe restarted both, the 
UI returned 502, and running tasks failed with `httpx.ConnectError: [Errno 111] 
Connection refused` on the Execution API. The error logged at the end of each 
stall:
   
   ```
   pydantic_core._pydantic_core.PydanticSerializationError: Error serializing 
to JSON: TimeoutError: QueuePool limit of size 5 overflow 1 reached, connection 
timed out, timeout 30.00
     File "fastapi/routing.py", line 695, in app
     File "fastapi/routing.py", line 306, in serialize_response
     File "fastapi/_compat/v2.py", line 231, in serialize_json
     File "pydantic/type_adapter.py", line 677, in dump_json
   ```
   
   The same error appeared with the default pool (`size 5 overflow 10`), so a 
larger pool raises the threshold but does not prevent the stall.
   
   **Steps to reproduce**
   
   This counts connection checkouts made on the event loop thread while serving 
the two endpoints, on an unmodified `apache-airflow==3.3.2` (FastAPI 0.136.3, 
pydantic 2.13.5) with SQLite.
   
   1. Environment:
   
   ```bash
   export AIRFLOW_HOME=/tmp/loopcheck
   export 
AIRFLOW__DATABASE__SQL_ALCHEMY_CONN=sqlite:////tmp/loopcheck/airflow.db
   export AIRFLOW__CORE__DAGS_FOLDER=/tmp/loopcheck/dags
   export AIRFLOW__CORE__LOAD_EXAMPLES=False
   export 
AIRFLOW__CORE__AUTH_MANAGER=airflow.api_fastapi.auth.managers.simple.simple_auth_manager.SimpleAuthManager
   export AIRFLOW__CORE__SIMPLE_AUTH_MANAGER_ALL_ADMINS=True
   export AIRFLOW__API_AUTH__JWT_SECRET=$(python -c "import secrets; 
print(secrets.token_hex(32))")
   mkdir -p /tmp/loopcheck/dags
   ```
   
   2. A DAG, `/tmp/loopcheck/dags/loopcheck.py`:
   
   ```python
   from airflow.sdk import DAG, task
   
   with DAG("loopcheck", schedule=None):
   
       @task
       def noop():
           pass
   
       noop()
   ```
   
   3. `airflow db migrate && airflow dags reserialize`
   
   4. Run this script:
   
   ```python
   import asyncio
   import traceback
   
   from fastapi.testclient import TestClient
   from sqlalchemy import event
   from sqlalchemy.pool import Pool
   
   from airflow.api_fastapi.app import create_app, get_auth_manager
   from airflow.api_fastapi.auth.managers.simple.user import 
SimpleAuthManagerUser
   
   stacks = []
   
   
   def checkout(*_):
       # A running loop exists only on the event loop thread; AnyIO worker 
threads raise here.
       try:
           asyncio.get_running_loop()
       except RuntimeError:
           return
       stacks.append(
           [
               f"{frame.filename.rsplit('/airflow/', 1)[-1]}:{frame.lineno} 
{frame.name}"
               for frame in traceback.extract_stack()
               if "/airflow/" in frame.filename and "sqlalchemy" not in 
frame.filename
           ][-4:]
       )
   
   
   event.listen(Pool, "checkout", checkout)
   
   with TestClient(create_app("core")) as client:
       token = 
get_auth_manager().generate_jwt(SimpleAuthManagerUser(username="admin", 
role="admin"))
       for path in ["/api/v2/dags/loopcheck/details", "/api/v2/pools"]:
           stacks.clear()
           response = client.get(path, headers={"Authorization": f"Bearer 
{token}"})
           print(f"{path} -> {response.status_code}, checkouts on the event 
loop: {len(stacks)}")
           for stack in stacks:
               print("    " + " > ".join(stack))
   ```
   
   Output (with only `default_pool`):
   
   ```
   /api/v2/dags/loopcheck/details -> 200, checkouts on the event loop: 2
       api_fastapi/core_api/security.py:126 resolve_user_from_token > 
api_fastapi/auth/managers/base_auth_manager.py:154 get_user_from_token > 
utils/session.py:100 wrapper > models/revoked_token.py:61 is_revoked
       lib/python3.13/site-packages/pydantic/type_adapter.py:677 dump_json > 
api_fastapi/core_api/datamodels/dags.py:281 latest_dag_version > 
utils/session.py:100 wrapper > models/dag_version.py:202 get_latest_version
   /api/v2/pools -> 200, checkouts on the event loop: 7
       api_fastapi/core_api/security.py:126 resolve_user_from_token > 
api_fastapi/auth/managers/base_auth_manager.py:154 get_user_from_token > 
utils/session.py:100 wrapper > models/revoked_token.py:61 is_revoked
       lib/python3.13/site-packages/pydantic/type_adapter.py:677 dump_json > 
api_fastapi/core_api/datamodels/pools.py:35 _call_function > 
utils/session.py:100 wrapper > models/pool.py:281 occupied_slots
       lib/python3.13/site-packages/pydantic/type_adapter.py:677 dump_json > 
api_fastapi/core_api/datamodels/pools.py:35 _call_function > 
utils/session.py:100 wrapper > models/pool.py:307 running_slots
       lib/python3.13/site-packages/pydantic/type_adapter.py:677 dump_json > 
api_fastapi/core_api/datamodels/pools.py:35 _call_function > 
utils/session.py:100 wrapper > models/pool.py:326 queued_slots
       lib/python3.13/site-packages/pydantic/type_adapter.py:677 dump_json > 
api_fastapi/core_api/datamodels/pools.py:35 _call_function > 
utils/session.py:100 wrapper > models/pool.py:345 scheduled_slots
       utils/session.py:100 wrapper > models/pool.py:382 open_slots > 
utils/session.py:98 wrapper > models/pool.py:281 occupied_slots
       lib/python3.13/site-packages/pydantic/type_adapter.py:677 dump_json > 
api_fastapi/core_api/datamodels/pools.py:35 _call_function > 
utils/session.py:100 wrapper > models/pool.py:364 deferred_slots
   ```
   
   The `is_revoked` checkout is a separate occurrence of the same problem, on 
every authenticated request; it was reported in #66493 and addressed by #73422, 
which was closed without merging. This issue is about the serialization-time 
queries.
   
   ## What you think should happen instead?
   
   Response serialization should not touch the database, so that a saturated 
pool slows the requests that need a connection instead of stopping the event 
loop.
   
   We run this change as a local patch on 3.3.2, and with it the script above 
reports no serialization-time checkouts for either endpoint, with the same 
response bodies (apart from `last_parsed`, the DagBag load time):
   
   - `get_dag_details` loads the latest version through the request's own 
session and sets it on the model, as it already does for `is_favorite` and 
`active_runs_count`; `latest_dag_version` becomes a regular `DagVersionResponse 
| None` field. This also removes the second connection each `/details` request 
took.
   - `PoolCollectionResponse.pools` becomes `list[PoolResponse]`, which 
pydantic validates eagerly when the route builds the response, inside the 
threadpool. The route already passes a `ScalarResult`, which `list` accepts.
   
   Seventeen other collection responses in `core_api/datamodels` are also typed 
`Iterable[...]`. We have not observed them querying during serialization, but 
any whose item validation reads lazy attributes would behave the same way. 
#67799 tracks moving the API off blocking database calls and lists 
serialization in its per-endpoint checklist.
   
   ## Operating System
   Debian GNU/Linux 12 (bookworm)
   
   ## Deployment
   Official Apache Airflow Helm Chart
   
   ## Apache Airflow Provider(s)
   (empty)
   
   ## Versions of Apache Airflow Providers
   Not applicable
   
   ## Official Helm Chart version
   1.21.0
   
   ## Kubernetes Version
   v1.35.8-gke.1036000 (GKE Autopilot)
   
   ## Helm Chart configuration
   ```yaml
   apiServer:
     replicas: 2
     env:
       - name: AIRFLOW__DATABASE__SQL_ALCHEMY_POOL_SIZE
         value: "5"
       - name: AIRFLOW__DATABASE__SQL_ALCHEMY_MAX_OVERFLOW
         value: "1"
   pgbouncer:
     enabled: true  # pool_mode = transaction
   ```
   The stall also occurred with the default pool settings (5 + 10).
   
   ## Docker Image customizations
   The official image extended with our DAG dependencies and local patches to 
the api-server: the health endpoint served from a cached snapshot, the 
revoked-token check run in the threadpool, Sentry initialization, and the fix 
described above. The reproduction above runs on unmodified 3.3.2.
   
   ## Anything else?
   Our Sentry recorded this error 292 times between 2026-07-28 and 2026-10-02 
(117 on `get_dag_details`, 32 on `get_pools`, the rest duplicate log captures), 
on weekdays when the UI is in use. It happens whenever the pool is exhausted 
while one of these responses serializes.
   
   ## Are you willing to submit PR?
   - [x] Yes I am willing to submit a PR!
   
   ## Code of Conduct
   - [x] I agree to follow this project's Code of Conduct
   
   ---
   Drafted-by: Claude Code (Opus 5.5); reviewed by @rodrigopimentel before 
posting
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to