This is an automated email from the ASF dual-hosted git repository. ashb pushed a commit to branch store-historic-ti-ownership-data in repository https://gitbox.apache.org/repos/asf/airflow.git
commit 9bd17a790f6cac9b7d75a192c820ac384828f00f Author: Ash Berlin-Taylor <[email protected]> AuthorDate: Mon Oct 5 18:49:21 2026 +0100 Reduced diffs now we added with_loader_options --- airflow-core/src/airflow/api/common/delete_dag.py | 6 ++---- .../api_fastapi/core_api/routes/public/task_instances.py | 7 +------ .../src/airflow/api_fastapi/core_api/routes/ui/grid.py | 6 ++---- .../api_fastapi/execution_api/routes/task_instances.py | 6 +----- airflow-core/src/airflow/models/dagrun.py | 6 +----- airflow-core/src/airflow/models/taskinstance.py | 5 +---- airflow-core/src/airflow/models/trigger.py | 13 ++++--------- airflow-core/src/airflow/ti_deps/deps/trigger_rule_dep.py | 15 +++------------ 8 files changed, 15 insertions(+), 49 deletions(-) diff --git a/airflow-core/src/airflow/api/common/delete_dag.py b/airflow-core/src/airflow/api/common/delete_dag.py index 5df54993c08..935b8faf078 100644 --- a/airflow-core/src/airflow/api/common/delete_dag.py +++ b/airflow-core/src/airflow/api/common/delete_dag.py @@ -55,10 +55,8 @@ def delete_dag(dag_id: str, keep_records_in_log: bool = True, *, session: Sessio log.info("Deleting Dag: %s", dag_id) running_tis = session.scalar( select(models.TaskInstance.state) - .where( - models.TaskInstance.dag_id == dag_id, - models.TaskInstance.state == TaskInstanceState.RUNNING, - ) + .where(models.TaskInstance.dag_id == dag_id) + .where(models.TaskInstance.state == TaskInstanceState.RUNNING) .limit(1) ) if running_tis: 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 1b77e22f512..1f40a0c73e7 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 @@ -439,12 +439,7 @@ def get_mapped_task_instance( """Get task instance.""" query = ( select(TI) - .where( - TI.dag_id == dag_id, - TI.run_id == dag_run_id, - TI.task_id == task_id, - TI.map_index == map_index, - ) + .where(TI.dag_id == dag_id, TI.run_id == dag_run_id, TI.task_id == task_id, TI.map_index == map_index) .options(joinedload(TI.rendered_task_instance_fields)) .options(joinedload(TI.dag_version)) .options(joinedload(TI.dag_run).options(joinedload(DagRun.dag_model))) 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 2dc91fa546b..e4f9b02f7d1 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 @@ -538,10 +538,8 @@ def get_grid_ti_summaries_stream( has_note_subq, ) .outerjoin(DagVersion, TaskInstance.dag_version_id == DagVersion.id) - .where( - TaskInstance.dag_id == dag_id, - TaskInstance.run_id == run_id, - ) + .where(TaskInstance.dag_id == dag_id) + .where(TaskInstance.run_id == run_id) .order_by(TaskInstance.task_id) .execution_options(yield_per=1000) ) diff --git a/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py b/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py index 0ea58145387..b05dd9e3bdc 100644 --- a/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py +++ b/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py @@ -1422,11 +1422,7 @@ async def get_task_instance_breadcrumbs( result = ( await session.execute( select(TI.task_id, TI.map_index, TI.state, TI.operator, TI.duration) - .where( - TI.dag_id == dag_id, - TI.run_id == run_id, - TI.state.in_(TerminalTIState), - ) + .where(TI.dag_id == dag_id, TI.run_id == run_id, TI.state.in_(TerminalTIState)) .order_by(TI.task_id, TI.map_index) ) ).mappings() diff --git a/airflow-core/src/airflow/models/dagrun.py b/airflow-core/src/airflow/models/dagrun.py index 09c62db8026..a23b0d8602f 100644 --- a/airflow-core/src/airflow/models/dagrun.py +++ b/airflow-core/src/airflow/models/dagrun.py @@ -585,11 +585,7 @@ class DagRun(Base, LoggingMixin): def check_version_id_exists_in_dr(self, dag_version_id: UUID, *, session: Session = NEW_SESSION): select_stmt = ( select(TI.dag_version_id) - .where( - TI.dag_id == self.dag_id, - TI.dag_version_id == dag_version_id, - TI.run_id == self.run_id, - ) + .where(TI.dag_id == self.dag_id, TI.dag_version_id == dag_version_id, TI.run_id == self.run_id) .limit(1) .execution_options(include_all_attempts=True) ) diff --git a/airflow-core/src/airflow/models/taskinstance.py b/airflow-core/src/airflow/models/taskinstance.py index ce187eb613d..d801fd7e747 100644 --- a/airflow-core/src/airflow/models/taskinstance.py +++ b/airflow-core/src/airflow/models/taskinstance.py @@ -2353,10 +2353,7 @@ class TaskInstance(Base, LoggingMixin, BaseWorkload): stmt = ( select(func.count()) .select_from(TaskInstance) - .where( - TaskInstance.dag_id == self.dag_id, - TaskInstance.task_id == self.task_id, - ) + .where(TaskInstance.dag_id == self.dag_id, TaskInstance.task_id == self.task_id) ) if states: stmt = stmt.where(or_(*(TaskInstance.state == s for s in states))) diff --git a/airflow-core/src/airflow/models/trigger.py b/airflow-core/src/airflow/models/trigger.py index cbc9a67a109..30fdc7d3117 100644 --- a/airflow-core/src/airflow/models/trigger.py +++ b/airflow-core/src/airflow/models/trigger.py @@ -245,8 +245,7 @@ class Trigger(Base): session.execute( update(TaskInstance) .where( - TaskInstance.state != TaskInstanceState.DEFERRED, - TaskInstance.trigger_id.is_not(None), + TaskInstance.state != TaskInstanceState.DEFERRED, TaskInstance.trigger_id.is_not(None) ) .values(trigger_id=None) ) @@ -281,8 +280,7 @@ class Trigger(Base): # Resume deferred tasks for task_instance in session.scalars( select(TaskInstance).where( - TaskInstance.trigger_id == trigger_id, - TaskInstance.state == TaskInstanceState.DEFERRED, + TaskInstance.trigger_id == trigger_id, TaskInstance.state == TaskInstanceState.DEFERRED ) ): handle_event_submit(event, task_instance=task_instance, session=session) @@ -321,8 +319,7 @@ class Trigger(Base): """ for task_instance in session.scalars( select(TaskInstance).where( - TaskInstance.trigger_id == trigger_id, - TaskInstance.state == TaskInstanceState.DEFERRED, + TaskInstance.trigger_id == trigger_id, TaskInstance.state == TaskInstanceState.DEFERRED ) ): # Add the error and set the next_method to the fail state @@ -471,9 +468,7 @@ class Trigger(Base): select(cls.id) .prefix_with("STRAIGHT_JOIN", dialect="mysql") .join(TaskInstance, cls.id == TaskInstance.trigger_id, isouter=False) - .where( - or_(cls.triggerer_id.is_(None), cls.triggerer_id.not_in(alive_triggerer_ids)), - ) + .where(or_(cls.triggerer_id.is_(None), cls.triggerer_id.not_in(alive_triggerer_ids))) .order_by(coalesce(TaskInstance.priority_weight, 0).desc(), cls.created_date), # Asset triggers select(cls.id) diff --git a/airflow-core/src/airflow/ti_deps/deps/trigger_rule_dep.py b/airflow-core/src/airflow/ti_deps/deps/trigger_rule_dep.py index 24b77000555..f49b71f7ca9 100644 --- a/airflow-core/src/airflow/ti_deps/deps/trigger_rule_dep.py +++ b/airflow-core/src/airflow/ti_deps/deps/trigger_rule_dep.py @@ -302,10 +302,7 @@ class TriggerRuleDep(BaseTIDep): else: task_id_counts = session.execute( select(TaskInstance.task_id, func.count(TaskInstance.task_id)) - .where( - TaskInstance.dag_id == ti.dag_id, - TaskInstance.run_id == ti.run_id, - ) + .where(TaskInstance.dag_id == ti.dag_id, TaskInstance.run_id == ti.run_id) .where(or_(*_iter_upstream_conditions(relevant_tasks=indirect_setups))) .group_by(TaskInstance.task_id) ).all() @@ -413,10 +410,7 @@ class TriggerRuleDep(BaseTIDep): (task_id, count) for task_id, count in session.execute( select(TaskInstance.task_id, func.count(TaskInstance.task_id)) - .where( - TaskInstance.dag_id == ti.dag_id, - TaskInstance.run_id == ti.run_id, - ) + .where(TaskInstance.dag_id == ti.dag_id, TaskInstance.run_id == ti.run_id) .where(or_(*_iter_upstream_conditions(relevant_tasks=upstream_tasks))) .group_by(TaskInstance.task_id) ) @@ -712,10 +706,7 @@ class TriggerRuleDep(BaseTIDep): expected = ( session.scalar( select(func.count(TaskInstance.task_id)) - .where( - TaskInstance.dag_id == ti.dag_id, - TaskInstance.run_id == ti.run_id, - ) + .where(TaskInstance.dag_id == ti.dag_id, TaskInstance.run_id == ti.run_id) .where(or_(*_iter_upstream_conditions(relevant_tasks=in_scope_tasks))) ) or 0
