Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -127,7 +127,7 @@ def attach_dag_versions_to_runs(dag_runs: Sequence[DagRun], *, session: Session)
.where(DagVersion.id.in_(all_version_ids))
.options(joinedload(DagVersion.bundle))
)
versions_by_id = {dv.id: dv for dv in session.scalars(dv_query).unique()}
versions_by_id = {dv.id: dv for dv in session.scalars(dv_query)}

versions_per_run: dict[tuple[str, str], dict[UUID, DagVersion]] = defaultdict(dict)
for row in rows:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -618,7 +618,7 @@ def get_dag_runs(
dag_run_select = order_by.to_orm(dag_run_select, reversed=True)
dag_run_select = apply_cursor_filter(dag_run_select, token, order_by, is_backward=is_backward)

fetched = list(session.scalars(dag_run_select).unique())
fetched = list(session.scalars(dag_run_select))
has_more = len(fetched) > page_limit
dag_runs = fetched[:page_limit]

Expand Down Expand Up @@ -648,7 +648,7 @@ def get_dag_runs(
limit=limit,
session=session,
)
dag_runs = list(session.scalars(dag_run_select).unique())
dag_runs = list(session.scalars(dag_run_select))
attach_dag_versions_to_runs(dag_runs, session=session)

return DAGRunCollectionResponse(
Expand Down Expand Up @@ -930,7 +930,7 @@ def get_list_dag_runs_batch(
session=session,
)

dag_runs = list(session.scalars(dag_runs_select).unique())
dag_runs = list(session.scalars(dag_runs_select))
attach_dag_versions_to_runs(dag_runs, session=session)

return DAGRunCollectionResponse(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -970,16 +970,12 @@ def _collect_relatives(run_id: str, direction: Literal["upstream", "downstream"]
# dag.clear() returns TIs without this relationship loaded; re-query with joinedload.
# populate_existing=True ensures the joinedload updates TIs already in the identity map.
if task_instances:
task_instances = (
session.scalars(
select(TI)
.options(joinedload(TI.rendered_task_instance_fields))
.where(TI.id.in_([ti.id for ti in task_instances]))
.execution_options(populate_existing=True)
)
.unique()
.all()
)
task_instances = session.scalars(
select(TI)
.options(joinedload(TI.rendered_task_instance_fields))
.where(TI.id.in_([ti.id for ti in task_instances]))
.execution_options(populate_existing=True)
).all()

return TaskInstanceCollectionResponse(
task_instances=[TaskInstanceResponse.model_validate(ti) for ti in task_instances],
Expand Down Expand Up @@ -1136,16 +1132,12 @@ def patch_task_instance_dry_run(
# set_task_instance_state() returns TIs without this relationship loaded; re-query with joinedload.
# populate_existing=True ensures the joinedload updates TIs already in the identity map.
if tis:
tis = (
session.scalars(
select(TI)
.options(joinedload(TI.rendered_task_instance_fields))
.where(TI.id.in_([ti.id for ti in tis]))
.execution_options(populate_existing=True)
)
.unique()
.all()
)
tis = session.scalars(
select(TI)
.options(joinedload(TI.rendered_task_instance_fields))
.where(TI.id.in_([ti.id for ti in tis]))
.execution_options(populate_existing=True)
).all()

return TaskInstanceCollectionResponse(
task_instances=[
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -129,7 +129,7 @@ def get_deadlines(
session=session,
)

deadlines = session.scalars(deadlines_select).unique()
deadlines = session.scalars(deadlines_select)

if dag_run_id != "~" and total_entries == 0:
dag_run = session.scalar(select(DagRun).where(DagRun.dag_id == dag_id, DagRun.run_id == dag_run_id))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -328,7 +328,7 @@ def get_grid_runs(
limit=limit,
return_total_entries=False,
)
results = session.execute(dag_runs_select_filter).unique().all()
results = session.execute(dag_runs_select_filter).all()
dag_runs = [run for run, _ in results]
attach_dag_versions_to_runs(dag_runs, session=session)
grid_runs = []
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -140,9 +140,7 @@ def _reload_tis_with_rendered_fields(tis: list[TI], session: Session) -> list[TI
.options(joinedload(TI.rendered_task_instance_fields))
.where(TI.id.in_([ti.id for ti in tis]))
.execution_options(populate_existing=True)
)
.unique()
.all()
).all()
)


Expand Down
18 changes: 7 additions & 11 deletions airflow-core/src/airflow/jobs/scheduler_job_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -2391,19 +2391,15 @@ def _create_dag_runs(self, dag_models: Collection[DagModel], session: Session) -
# as DagModel.dag_id and DagModel.next_dagrun
# This list is used to verify if the DagRun already exist so that we don't attempt to create
# duplicate DagRuns
existing_dagrun_objects = (
session.scalars(
select(DagRun)
.where(
tuple_(DagRun.dag_id, DagRun.logical_date).in_(
(dm.dag_id, dm.next_dagrun) for dm in dag_models
)
existing_dagrun_objects = session.scalars(
select(DagRun)
.where(
tuple_(DagRun.dag_id, DagRun.logical_date).in_(
(dm.dag_id, dm.next_dagrun) for dm in dag_models
)
.options(load_only(DagRun.dag_id, DagRun.logical_date))
)
.unique()
.all()
)
.options(load_only(DagRun.dag_id, DagRun.logical_date))
).all()
existing_dagruns = {(x.dag_id, x.logical_date): x for x in existing_dagrun_objects}

# backfill runs are not created by scheduler and their concurrency is separate
Expand Down