diff --git a/airflow-core/src/airflow/api_fastapi/common/db/dag_runs.py b/airflow-core/src/airflow/api_fastapi/common/db/dag_runs.py index f9ae3d6c3b6c6..1508ea8528ebb 100644 --- a/airflow-core/src/airflow/api_fastapi/common/db/dag_runs.py +++ b/airflow-core/src/airflow/api_fastapi/common/db/dag_runs.py @@ -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: diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_run.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_run.py index 59b3593b9ed63..81d47539229ce 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_run.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_run.py @@ -621,7 +621,7 @@ def get_dag_runs( dag_run_select, token, order_by, session.get_bind().dialect.name, 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] @@ -651,7 +651,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( @@ -933,7 +933,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( diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py index 77f68611c3e7f..f72cd3801c65c 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py @@ -973,16 +973,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], @@ -1139,16 +1135,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=[ diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/deadlines.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/deadlines.py index d9cfea10d6d94..06eda42ed8957 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/deadlines.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/deadlines.py @@ -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)) diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/grid.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/grid.py index a16f9336fd57e..8a0c55d315ec5 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/grid.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/grid.py @@ -334,7 +334,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 = [] diff --git a/airflow-core/src/airflow/api_fastapi/core_api/services/public/task_instances.py b/airflow-core/src/airflow/api_fastapi/core_api/services/public/task_instances.py index 24ae587680a86..0a1885f0e26fe 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/services/public/task_instances.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/services/public/task_instances.py @@ -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() ) diff --git a/airflow-core/src/airflow/jobs/scheduler_job_runner.py b/airflow-core/src/airflow/jobs/scheduler_job_runner.py index 7d5f7f58020c7..f6d6a4787f629 100644 --- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py +++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py @@ -2417,19 +2417,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