This is an automated email from the ASF dual-hosted git repository.
jason810496 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 7c597e2830b Speed up the event logs API by loading only display names
(#73021)
7c597e2830b is described below
commit 7c597e2830b1fa1a8a3130f81945366ba7547512
Author: PoAn Yang <[email protected]>
AuthorDate: Sun Sep 27 13:01:23 2026 +0800
Speed up the event logs API by loading only display names (#73021)
---
.../core_api/routes/public/event_logs.py | 31 ++++++++++++++++------
.../core_api/routes/public/test_event_logs.py | 29 +++++++++++++++++++-
2 files changed, 51 insertions(+), 9 deletions(-)
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/event_logs.py
b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/event_logs.py
index 524b66b8275..481209b0883 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/event_logs.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/event_logs.py
@@ -17,7 +17,7 @@
from __future__ import annotations
from datetime import datetime
-from typing import Annotated
+from typing import TYPE_CHECKING, Annotated
from fastapi import Depends, HTTPException, status
from sqlalchemy import select
@@ -50,11 +50,30 @@ from airflow.api_fastapi.core_api.security import (
requires_access_event_log,
)
from airflow.api_fastapi.core_api.services.public.event_logs import
event_log_to_response
-from airflow.models import Log
+from airflow.models import DagModel, Log, TaskInstance
+
+if TYPE_CHECKING:
+ from sqlalchemy.orm.interfaces import LoaderOption
event_logs_router = AirflowRouter(tags=["Event Log"], prefix="/eventLogs")
+def _eager_load_display_names() -> tuple[LoaderOption, ...]:
+ """
+ Load only the columns that EventLogResponse reads from an event log's Dag
and task instance.
+
+ The query also loads ``task_id`` because ``task_display_name`` falls back
to it.
+ Without ``raiseload``, it would still join ``dag_run`` and select all of
its columns, since
+ ``TaskInstance.dag_run`` is ``lazy="joined"``.
+ """
+ return (
+ joinedload(Log.task_instance)
+ .load_only(TaskInstance._task_display_property_value,
TaskInstance.task_id)
+ .raiseload(TaskInstance.dag_run),
+
joinedload(Log.dag_model).load_only(DagModel._dag_display_property_value),
+ )
+
+
@event_logs_router.get(
"/{event_log_id}",
responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]),
@@ -71,9 +90,7 @@ def get_event_log(
# that bypass Log.__init__ (which always sets dttm =
timezone.utcnow()).
# Making EventLogResponse.when nullable would be a breaking API
contract change for
# clients that currently rely on `when` always being present.
- select(Log)
- .where(Log.id == event_log_id, Log.dttm.is_not(None))
- .options(joinedload(Log.task_instance), joinedload(Log.dag_model))
+ select(Log).where(Log.id == event_log_id,
Log.dttm.is_not(None)).options(*_eager_load_display_names())
)
if event_log is None:
raise HTTPException(status.HTTP_404_NOT_FOUND, f"The Event Log with
id: `{event_log_id}` not found")
@@ -187,9 +204,7 @@ def get_event_logs(
# that bypass Log.__init__ (which always sets dttm =
timezone.utcnow()).
# Making EventLogResponse.when nullable would be a breaking API
contract change for
# clients that currently rely on `when` always being present.
- select(Log)
- .where(Log.dttm.is_not(None))
- .options(joinedload(Log.task_instance), joinedload(Log.dag_model))
+
select(Log).where(Log.dttm.is_not(None)).options(*_eager_load_display_names())
)
event_logs_select, total_entries = paginated_select(
statement=query,
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_event_logs.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_event_logs.py
index 240c95a46a2..c5224745505 100644
---
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_event_logs.py
+++
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_event_logs.py
@@ -16,6 +16,7 @@
# under the License.
from __future__ import annotations
+import re
from datetime import datetime, timezone
from unittest import mock
@@ -30,7 +31,7 @@ from
airflow.api_fastapi.auth.managers.models.resource_details import (
from airflow.models.log import Log
from airflow.utils.session import NEW_SESSION, provide_session
-from tests_common.test_utils.asserts import assert_queries_count
+from tests_common.test_utils.asserts import assert_queries_count,
capture_orm_selects
from tests_common.test_utils.db import clear_db_logs, clear_db_runs
from tests_common.test_utils.format_datetime import from_datetime_to_zulu,
from_datetime_to_zulu_without_ms
@@ -58,6 +59,18 @@ TEAM_EVENT = "TEAM_EVENT"
TEAM_NAME = "TEST_TEAM"
+def _assert_selects_only_display_name_columns(statements: list[str]) -> None:
+ (sql,) = [sql for sql in statements if "task_instance_1" in sql]
+ select_clause = sql.split(" FROM ", 1)[0]
+ assert set(re.findall(r"\bdag_1\.(\w+)", select_clause)) == {"dag_id",
"dag_display_name"}
+ assert set(re.findall(r"\btask_instance_1\.(\w+)", select_clause)) == {
+ "id",
+ "task_id",
+ "task_display_name",
+ }
+ assert "dag_run" not in sql
+
+
class TestEventLogsEndpoint:
"""Common class for /eventLogs related unit tests."""
@@ -199,6 +212,13 @@ class TestGetEventLog(TestEventLogsEndpoint):
assert response.json() == expected_json
+ def test_get_event_log_selects_only_display_name_columns(self,
test_client, setup):
+ with capture_orm_selects("log") as statements:
+ response =
test_client.get(f"/eventLogs/{setup[TASK_INSTANCE_EVENT].id}")
+
+ assert response.status_code == 200
+ _assert_selects_only_display_name_columns(statements)
+
def test_get_event_log_returns_the_recorded_team(self, test_client,
session):
event_log = Log(event="cli_triggerer", team_name=TEAM_NAME)
session.add(event_log)
@@ -426,6 +446,13 @@ class TestGetEventLogs(TestEventLogsEndpoint):
for event_log, expected_event in zip(resp_json["event_logs"],
expected_events):
assert event_log["event"] == expected_event
+ def test_get_event_logs_selects_only_display_name_columns(self,
test_client):
+ with capture_orm_selects("log") as statements:
+ response = test_client.get("/eventLogs")
+
+ assert response.status_code == 200
+ _assert_selects_only_display_name_columns(statements)
+
@provide_session
def test_get_event_logs_excludes_logs_without_dttm(
self,