From 978a6f725a00cb49bc31754c74ede81dd4ee6b2f Mon Sep 17 00:00:00 2001 From: Jarek Potiuk Date: Fri, 19 May 2023 19:19:36 +0200 Subject: [PATCH 1/2] Revert "Revert "Save scheduler execution time by caching dags (#30704)" (#31413)" This reverts commit e6f21174ab42092b4ccc9f264627aa43b48b5a02. --- airflow/jobs/scheduler_job_runner.py | 18 +++++++++++++++--- 1 file changed, 15 insertions(+), 3 deletions(-) diff --git a/airflow/jobs/scheduler_job_runner.py b/airflow/jobs/scheduler_job_runner.py index 128a82d7eedb3..775ac759ffb64 100644 --- a/airflow/jobs/scheduler_job_runner.py +++ b/airflow/jobs/scheduler_job_runner.py @@ -28,8 +28,9 @@ from collections import Counter from dataclasses import dataclass from datetime import datetime, timedelta +from functools import lru_cache, partial from pathlib import Path -from typing import TYPE_CHECKING, Any, Collection, Iterable, Iterator +from typing import TYPE_CHECKING, Any, Callable, Collection, Iterable, Iterator from sqlalchemy import and_, func, not_, or_, text from sqlalchemy.exc import OperationalError @@ -1052,8 +1053,13 @@ def _do_scheduling(self, session: Session) -> int: callback_tuples = self._schedule_all_dag_runs(guard, dag_runs, session) # Send the callbacks after we commit to ensure the context is up to date when it gets run + # cache saves time during scheduling of many dag_runs for same dag + cached_get_dag: Callable[[str], DAG | None] = lru_cache()( + partial(self.dagbag.get_dag, session=session) + ) for dag_run, callback_to_run in callback_tuples: - dag = self.dagbag.get_dag(dag_run.dag_id, session=session) + dag = cached_get_dag(dag_run.dag_id) + if not dag: self.log.error("DAG '%s' not found in serialized_dag table", dag_run.dag_id) continue @@ -1317,8 +1323,14 @@ def _update_state(dag: DAG, dag_run: DagRun): tags={"dag_id": dag.dag_id}, ) + # cache saves time during scheduling of many dag_runs for same dag + cached_get_dag: Callable[[str], DAG | None] = lru_cache()( + partial(self.dagbag.get_dag, session=session) + ) + for dag_run in dag_runs: - dag = dag_run.dag = self.dagbag.get_dag(dag_run.dag_id, session=session) + dag = dag_run.dag = cached_get_dag(dag_run.dag_id) + if not dag: self.log.error("DAG '%s' not found in serialized_dag table", dag_run.dag_id) continue From 50cc7021eaaa634f5af6407dc18b20b17fe3f10c Mon Sep 17 00:00:00 2001 From: Jarek Potiuk Date: Fri, 19 May 2023 19:20:05 +0200 Subject: [PATCH 2/2] Revert "Save scheduler execution time by adding new Index idea for dag_run (#30827)" This reverts commit c63b7774cdba29394ec746b381f45e816dcb0830. --- ..._add_index_on_last_scheduling_decision_.py | 56 ------------------- airflow/models/dagrun.py | 9 --- docs/apache-airflow/img/airflow_erd.sha256 | 2 +- docs/apache-airflow/img/airflow_erd.svg | 4 +- docs/apache-airflow/migrations-ref.rst | 5 +- 5 files changed, 4 insertions(+), 72 deletions(-) delete mode 100644 airflow/migrations/versions/0126_2_7_0_add_index_on_last_scheduling_decision_.py diff --git a/airflow/migrations/versions/0126_2_7_0_add_index_on_last_scheduling_decision_.py b/airflow/migrations/versions/0126_2_7_0_add_index_on_last_scheduling_decision_.py deleted file mode 100644 index 6f94ce27dd137..0000000000000 --- a/airflow/migrations/versions/0126_2_7_0_add_index_on_last_scheduling_decision_.py +++ /dev/null @@ -1,56 +0,0 @@ -# -# Licensed to the Apache Software Foundation (ASF) under one -# or more contributor license agreements. See the NOTICE file -# distributed with this work for additional information -# regarding copyright ownership. The ASF licenses this file -# to you under the Apache License, Version 2.0 (the -# "License"); you may not use this file except in compliance -# with the License. You may obtain a copy of the License at -# -# http://www.apache.org/licenses/LICENSE-2.0 -# -# Unless required by applicable law or agreed to in writing, -# software distributed under the License is distributed on an -# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY -# KIND, either express or implied. See the License for the -# specific language governing permissions and limitations -# under the License. - -"""Add index on last_scheduling_decision NULLS FIRST, execution_date, state for queued dagrun - -Revision ID: 14db5484317e -Revises: 937cbd173ca1 -Create Date: 2023-05-14 21:16:42.399167 - -""" -from __future__ import annotations - -from alembic import op -from sqlalchemy import text - -# revision identifiers, used by Alembic. -revision = "14db5484317e" -down_revision = "937cbd173ca1" -branch_labels = None -depends_on = None -airflow_version = "2.7.0" - - -def upgrade(): - """Apply Add index on last_scheduling_decision NULLS FIRST, execution_date, state for queued dagrun""" - conn = op.get_bind() - if conn.dialect.name == "postgresql": - with op.batch_alter_table("dag_run") as batch_op: - batch_op.create_index( - "idx_last_scheduling_decision_queued", - [text("last_scheduling_decision NULLS FIRST"), "execution_date", "state"], - postgresql_where=text("state='queued'"), - ) - - -def downgrade(): - """Unapply Add index on last_scheduling_decision NULLS FIRST, execution_date, state for queued dagrun""" - conn = op.get_bind() - if conn.dialect.name == "postgresql": - with op.batch_alter_table("dag_run") as batch_op: - batch_op.drop_index("idx_last_scheduling_decision_queued") diff --git a/airflow/models/dagrun.py b/airflow/models/dagrun.py index dba6fbf748376..42845b34bdd0e 100644 --- a/airflow/models/dagrun.py +++ b/airflow/models/dagrun.py @@ -144,15 +144,6 @@ class DagRun(Base, LoggingMixin): UniqueConstraint("dag_id", "execution_date", name="dag_run_dag_id_execution_date_key"), UniqueConstraint("dag_id", "run_id", name="dag_run_dag_id_run_id_key"), Index("idx_last_scheduling_decision", last_scheduling_decision), - Index( - "idx_last_scheduling_decision_queued", - # Not possible to add .nulls_first(), because only postgresql can handle Index like that. - # Migration script which contains postgres dialect check adds NULLS FIST to index. - last_scheduling_decision, - execution_date, - _state, - postgresql_where=text("state='queued'"), - ), Index("idx_dag_run_dag_id", dag_id), Index( "idx_dag_run_running_dags", diff --git a/docs/apache-airflow/img/airflow_erd.sha256 b/docs/apache-airflow/img/airflow_erd.sha256 index 8cde30d727163..1f2c2c1f3419a 100644 --- a/docs/apache-airflow/img/airflow_erd.sha256 +++ b/docs/apache-airflow/img/airflow_erd.sha256 @@ -1 +1 @@ -811b1c45f8fa985feacacffafc30c82a6049bb33948c33bb218c13c48f971097 \ No newline at end of file +4987842fd67d29e194f1117e127d3291ba60d3fbc3e81cba75ce93884c263321 \ No newline at end of file diff --git a/docs/apache-airflow/img/airflow_erd.svg b/docs/apache-airflow/img/airflow_erd.svg index b531759ac4ec7..8439c226f8b98 100644 --- a/docs/apache-airflow/img/airflow_erd.svg +++ b/docs/apache-airflow/img/airflow_erd.svg @@ -1225,7 +1225,7 @@ task_instance--xcom -0..N +1 1 @@ -1239,7 +1239,7 @@ task_instance--xcom -1 +0..N 1 diff --git a/docs/apache-airflow/migrations-ref.rst b/docs/apache-airflow/migrations-ref.rst index c29c0fbeb86bd..2fdc682025a07 100644 --- a/docs/apache-airflow/migrations-ref.rst +++ b/docs/apache-airflow/migrations-ref.rst @@ -39,10 +39,7 @@ Here's the list of all the Database Migrations that are executed via when you ru +---------------------------------+-------------------+-------------------+--------------------------------------------------------------+ | Revision ID | Revises ID | Airflow Version | Description | +=================================+===================+===================+==============================================================+ -| ``14db5484317e`` (head) | ``937cbd173ca1`` | ``2.7.0`` | Add index on last_scheduling_decision NULLS FIRST, | -| | | | execution_date, state for queued dagrun | -+---------------------------------+-------------------+-------------------+--------------------------------------------------------------+ -| ``937cbd173ca1`` | ``98ae134e6fff`` | ``2.7.0`` | Add index to task_instance table | +| ``937cbd173ca1`` (head) | ``98ae134e6fff`` | ``2.7.0`` | Add index to task_instance table | +---------------------------------+-------------------+-------------------+--------------------------------------------------------------+ | ``98ae134e6fff`` | ``6abdffdd4815`` | ``2.6.0`` | Increase length of user identifier columns in ``ab_user`` | | | | | and ``ab_register_user`` tables |