Skip to content

Commit 67b6015

Browse files
[v3-3-test] Reduce memory used when deleting queued asset events (#71917) (#71937)
delete_asset_queued_events and delete_dag_asset_queued_event forced SQLAlchemy's "fetch" synchronize_session strategy on their AssetDagRunQueue deletes. That strategy reads the primary key of every deleted row back from the database to update the ORM session's identity map, but neither endpoint loads any AssetDagRunQueue objects into the session beforehand, so the read-back keys are matched against an empty map and discarded. The sibling delete_dag_asset_queued_events endpoint already used the default "auto" strategy; the other two now match it. (cherry picked from commit 49c5d51) Co-authored-by: Jyun-An Chen <jun930436@gmail.com>
1 parent 3a7d6cf commit 67b6015

2 files changed

Lines changed: 62 additions & 4 deletions

File tree

  • airflow-core

airflow-core/src/airflow/api_fastapi/core_api/routes/public/assets.py

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -667,7 +667,7 @@ def delete_asset_queued_events(
667667
where_clause = _generate_queued_event_where_clause(
668668
asset_id=asset_id, before=before, permitted_dag_ids=readable_dags_filter.value
669669
)
670-
delete_stmt = delete(AssetDagRunQueue).where(*where_clause).execution_options(synchronize_session="fetch")
670+
delete_stmt = delete(AssetDagRunQueue).where(*where_clause)
671671
result = cast("CursorResult", session.execute(delete_stmt))
672672
if result.rowcount == 0:
673673
raise HTTPException(
@@ -734,9 +734,7 @@ def delete_dag_asset_queued_event(
734734
where_clause = _generate_queued_event_where_clause(
735735
dag_id=dag_id, before=before, asset_id=asset_id, permitted_dag_ids=readable_dags_filter.value
736736
)
737-
delete_statement = (
738-
delete(AssetDagRunQueue).where(*where_clause).execution_options(synchronize_session="fetch")
739-
)
737+
delete_statement = delete(AssetDagRunQueue).where(*where_clause)
740738
result = cast("CursorResult", session.execute(delete_statement))
741739
if result.rowcount == 0:
742740
raise HTTPException(

airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_assets.py

Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2045,6 +2045,36 @@ def test_should_respond_404(self, test_client):
20452045
assert response.status_code == 404
20462046
assert response.json()["detail"] == "Queue event with asset_id: `1` was not found"
20472047

2048+
def test_delete_does_not_read_back_deleted_row_keys(self, test_client, session, create_dummy_dag):
2049+
from sqlalchemy import event
2050+
2051+
import airflow.settings
2052+
2053+
dag, _ = create_dummy_dag()
2054+
dag_id = dag.dag_id
2055+
(asset,) = self.create_assets(session=session, num=1)
2056+
self._create_asset_dag_run_queues(dag_id, asset.id, session)
2057+
2058+
executed_statements: list[str] = []
2059+
2060+
def capture(_conn, _cursor, statement, _parameters, _context, _executemany):
2061+
executed_statements.append(" ".join(statement.split()).upper())
2062+
2063+
event.listen(airflow.settings.engine, "before_cursor_execute", capture)
2064+
try:
2065+
response = test_client.delete(f"/assets/{asset.id}/queuedEvents")
2066+
finally:
2067+
event.remove(airflow.settings.engine, "before_cursor_execute", capture)
2068+
2069+
assert response.status_code == 204
2070+
deletes = [s for s in executed_statements if s.startswith("DELETE")]
2071+
assert deletes, "Expected the endpoint to issue a DELETE statement"
2072+
assert [s for s in deletes if "RETURNING" in s] == [], "DELETE must not read back deleted keys"
2073+
after_first_delete = executed_statements[executed_statements.index(deletes[0]) :]
2074+
assert [s for s in after_first_delete if s.startswith("SELECT")] == [], (
2075+
"No SELECT may precede a DELETE to collect the keys it is about to remove"
2076+
)
2077+
20482078

20492079
class TestDeleteDagAssetQueuedEvent(TestQueuedEventEndpoint):
20502080
def test_delete_should_respond_204(self, test_client, session, create_dummy_dag):
@@ -2073,6 +2103,36 @@ def test_should_respond_403(self, unauthorized_test_client):
20732103
response = unauthorized_test_client.delete("/dags/random/assets/random/queuedEvents")
20742104
assert response.status_code == 403
20752105

2106+
def test_delete_does_not_read_back_deleted_row_keys(self, test_client, session, create_dummy_dag):
2107+
from sqlalchemy import event
2108+
2109+
import airflow.settings
2110+
2111+
dag, _ = create_dummy_dag()
2112+
dag_id = dag.dag_id
2113+
(asset,) = self.create_assets(session=session, num=1)
2114+
self._create_asset_dag_run_queues(dag_id, asset.id, session)
2115+
2116+
executed_statements: list[str] = []
2117+
2118+
def capture(_conn, _cursor, statement, _parameters, _context, _executemany):
2119+
executed_statements.append(" ".join(statement.split()).upper())
2120+
2121+
event.listen(airflow.settings.engine, "before_cursor_execute", capture)
2122+
try:
2123+
response = test_client.delete(f"/dags/{dag_id}/assets/{asset.id}/queuedEvents")
2124+
finally:
2125+
event.remove(airflow.settings.engine, "before_cursor_execute", capture)
2126+
2127+
assert response.status_code == 204
2128+
deletes = [s for s in executed_statements if s.startswith("DELETE")]
2129+
assert deletes, "Expected the endpoint to issue a DELETE statement"
2130+
assert [s for s in deletes if "RETURNING" in s] == [], "DELETE must not read back deleted keys"
2131+
after_first_delete = executed_statements[executed_statements.index(deletes[0]) :]
2132+
assert [s for s in after_first_delete if s.startswith("SELECT")] == [], (
2133+
"No SELECT may precede a DELETE to collect the keys it is about to remove"
2134+
)
2135+
20762136
def test_should_respond_404(self, test_client):
20772137
dag_id = "not_exists"
20782138
asset_id = 1

0 commit comments

Comments
 (0)