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]
