Skip to content

Commit bb68de6

Browse files
committed
Preserve asset event ordering in bulk Dag run associations
Asset-triggered Dag runs should retain the established event selection order while avoiding ORM materialization. Scoping the ordering to scheduler association creation avoids adding sort costs to every consumed_asset_events relationship load. Signed-off-by: viiccwen <vicwen@apache.org>
1 parent f1963ec commit bb68de6

2 files changed

Lines changed: 23 additions & 19 deletions

File tree

airflow-core/src/airflow/jobs/scheduler_job_runner.py

Lines changed: 22 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -169,11 +169,11 @@
169169
def _associate_asset_events_with_dag_run(
170170
*,
171171
dag_run: DagRun,
172-
asset_event_ids: Select[tuple[int]],
172+
asset_event_ids_select: Select[tuple[int]],
173173
session: Session,
174174
) -> None:
175175
"""Associate selected asset events without materializing ORM objects."""
176-
selected_event_ids = asset_event_ids.subquery()
176+
selected_event_ids = asset_event_ids_select.subquery()
177177
session.execute(
178178
insert(association_table).from_select(
179179
["dag_run_id", "event_id"],
@@ -2362,13 +2362,13 @@ def _create_dagruns_for_partitioned_asset_dags(self, session: Session) -> set[st
23622362
creating_job_id=self.job.id,
23632363
session=session,
23642364
)
2365-
asset_event_ids = select(AssetEvent.id.label("event_id")).where(
2365+
asset_event_ids_select = select(AssetEvent.id.label("event_id")).where(
23662366
PartitionedAssetKeyLog.asset_partition_dag_run_id == apdr.id,
23672367
PartitionedAssetKeyLog.asset_event_id == AssetEvent.id,
23682368
)
23692369
_associate_asset_events_with_dag_run(
23702370
dag_run=dag_run,
2371-
asset_event_ids=asset_event_ids,
2371+
asset_event_ids_select=asset_event_ids_select,
23722372
session=session,
23732373
)
23742374
session.flush()
@@ -2642,21 +2642,25 @@ def _create_dag_runs_asset_triggered(
26422642
)
26432643
event_window_floor.append(date.min)
26442644

2645-
asset_event_ids = select(AssetEvent.id.label("event_id")).where(
2646-
or_(
2647-
AssetEvent.asset_id.in_(
2648-
select(DagScheduleAssetReference.asset_id).where(
2649-
DagScheduleAssetReference.dag_id == dag.dag_id
2645+
asset_event_ids_select = (
2646+
select(AssetEvent.id.label("event_id"))
2647+
.where(
2648+
or_(
2649+
AssetEvent.asset_id.in_(
2650+
select(DagScheduleAssetReference.asset_id).where(
2651+
DagScheduleAssetReference.dag_id == dag.dag_id
2652+
),
2653+
),
2654+
AssetEvent.source_aliases.any(
2655+
AssetAliasModel.scheduled_dags.any(
2656+
DagScheduleAssetAliasReference.dag_id == dag.dag_id
2657+
)
26502658
),
26512659
),
2652-
AssetEvent.source_aliases.any(
2653-
AssetAliasModel.scheduled_dags.any(
2654-
DagScheduleAssetAliasReference.dag_id == dag.dag_id
2655-
)
2656-
),
2657-
),
2658-
AssetEvent.timestamp <= triggered_date,
2659-
AssetEvent.timestamp > func.coalesce(*event_window_floor),
2660+
AssetEvent.timestamp <= triggered_date,
2661+
AssetEvent.timestamp > func.coalesce(*event_window_floor),
2662+
)
2663+
.order_by(AssetEvent.timestamp.asc(), AssetEvent.id.asc())
26602664
)
26612665

26622666
dag_run = dag.create_dagrun(
@@ -2680,7 +2684,7 @@ def _create_dag_runs_asset_triggered(
26802684
stats.incr("asset.triggered_dagruns", tags=prune_dict({"team_name": team_name}))
26812685
_associate_asset_events_with_dag_run(
26822686
dag_run=dag_run,
2683-
asset_event_ids=asset_event_ids,
2687+
asset_event_ids_select=asset_event_ids_select,
26842688
session=session,
26852689
)
26862690
consumed_asset_event_count = (

airflow-core/tests/unit/jobs/test_scheduler_job.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5596,7 +5596,7 @@ def test_associate_asset_events_with_dag_run_keeps_relationship_coherent(self, s
55965596
with assert_queries_count(1, session=session):
55975597
_associate_asset_events_with_dag_run(
55985598
dag_run=dag_run,
5599-
asset_event_ids=select(AssetEvent.id.label("event_id")).where(
5599+
asset_event_ids_select=select(AssetEvent.id.label("event_id")).where(
56005600
AssetEvent.id == asset_event_id
56015601
),
56025602
session=session,

0 commit comments

Comments
 (0)