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 015e17648aba3..617b70ea8941c 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 @@ -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] @@ -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( @@ -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( 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 68120ace15380..30726b9524da3 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 @@ -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], @@ -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=[ 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 72d92fbef2c2f..ee81e2887c7f3 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 @@ -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 = [] 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 a5874035e27f8..c1738f7382ad9 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 dce7c88bad71c..c37887c8c4bd0 100644 --- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py +++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py @@ -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