From bb9ed07ef5c30064e7633374b3ea8add3ed162e0 Mon Sep 17 00:00:00 2001 From: Aleksej Kirilishin Date: Sat, 29 Jan 2022 20:56:48 +0300 Subject: [PATCH 1/4] Removing the reset next_dagrun_create_after --- airflow/jobs/scheduler_job.py | 1 - airflow/models/dag.py | 4 +- tests/jobs/test_scheduler_job.py | 63 ++++++++++++++++++++++++-------- tests/models/test_dag.py | 37 ------------------- 4 files changed, 48 insertions(+), 57 deletions(-) diff --git a/airflow/jobs/scheduler_job.py b/airflow/jobs/scheduler_job.py index cbda16eb49e02..e8f4531882b57 100644 --- a/airflow/jobs/scheduler_job.py +++ b/airflow/jobs/scheduler_job.py @@ -972,7 +972,6 @@ def _should_update_dag_next_dagruns(self, dag, dag_model: DagModel, total_active total_active_runs, dag.max_active_runs, ) - dag_model.next_dagrun_create_after = None return False return True diff --git a/airflow/models/dag.py b/airflow/models/dag.py index 5cf1731dc6c8c..a300f97c94e58 100644 --- a/airflow/models/dag.py +++ b/airflow/models/dag.py @@ -2438,9 +2438,7 @@ def bulk_write_to_db(cls, dags: Collection["DAG"], session=NEW_SESSION): data_interval = None else: data_interval = dag.get_run_data_interval(run) - if num_active_runs.get(dag.dag_id, 0) >= orm_dag.max_active_runs: - orm_dag.next_dagrun_create_after = None - else: + if num_active_runs.get(dag.dag_id, 0) < orm_dag.max_active_runs: orm_dag.calculate_dagrun_date_fields(dag, data_interval) for orm_tag in list(orm_dag.tags): diff --git a/tests/jobs/test_scheduler_job.py b/tests/jobs/test_scheduler_job.py index 707f587223f0a..afcabf376c390 100644 --- a/tests/jobs/test_scheduler_job.py +++ b/tests/jobs/test_scheduler_job.py @@ -1255,7 +1255,6 @@ def test_queued_dagruns_stops_creating_when_max_active_is_reached(self, dag_make assert session.query(DagRun).count() == 10 assert session.query(DagRun.state).filter(DagRun.state == State.RUNNING).count() == 10 assert session.query(DagRun.state).filter(DagRun.state == State.QUEUED).count() == 0 - assert orm_dag.next_dagrun_create_after is None def test_dagrun_timeout_verify_max_active_runs(self, dag_maker): """ @@ -1287,7 +1286,6 @@ def test_dagrun_timeout_verify_max_active_runs(self, dag_maker): assert len(drs) == 1 dr = drs[0] - assert orm_dag.next_dagrun_create_after is None # But we should record the date of _what run_ it would be assert isinstance(orm_dag.next_dagrun, datetime.datetime) assert isinstance(orm_dag.next_dagrun_data_interval_start, datetime.datetime) @@ -2926,6 +2924,46 @@ def test_do_schedule_max_active_runs_task_removed(self, session, dag_maker): ti.refresh_from_db(session=session) assert ti.state == State.QUEUED + def test_runs_are_created_after_max_active_runs_was_reached(self, dag_maker, caplog): + """ + This tests that more DagRuns are created after max_active_runs had reached + and some DagRun had finished. + """ + with dag_maker(max_active_runs=1): + DummyOperator(task_id='task') + self.scheduler_job = SchedulerJob(subdir=os.devnull) + self.scheduler_job.executor = MockExecutor(do_update=False) + self.scheduler_job.processor_agent = mock.MagicMock(spec=DagFileProcessorAgent) + session = settings.Session() + + # schedule cycle when DR created + self.scheduler_job._do_scheduling(session) + session.flush() + assert session.query(DagRun).count() == 1 + dm = session.query(DagModel).one() + assert dm.next_dagrun == DEFAULT_DATE + + # schedule cycle when max_active_runs reached + self.scheduler_job._do_scheduling(session) + session.flush() + + # set dagrun to success + dr = session.query(DagRun).one() + dr.state = DagRunState.SUCCESS + ti = dr.get_task_instance('task', session) + ti.state = TaskInstanceState.SUCCESS + session.merge(ti) + session.merge(dr) + session.flush() + + # few cycle for calculate date and create new DR + for _ in range(2): + self.scheduler_job._do_scheduling(session) + session.flush() + assert session.query(DagRun).count() == 2 + dm = session.query(DagModel).one() + assert dm.next_dagrun == DEFAULT_DATE + timedelta(days=1) + def test_more_runs_are_not_created_when_max_active_runs_is_reached(self, dag_maker, caplog): """ This tests that when max_active_runs is reached, _create_dag_runs doesn't create @@ -2937,20 +2975,12 @@ def test_more_runs_are_not_created_when_max_active_runs_is_reached(self, dag_mak self.scheduler_job.executor = MockExecutor(do_update=False) self.scheduler_job.processor_agent = mock.MagicMock(spec=DagFileProcessorAgent) session = settings.Session() - assert session.query(DagRun).count() == 0 - dag_models = DagModel.dags_needing_dagruns(session).all() - self.scheduler_job._create_dag_runs(dag_models, session) - dr = session.query(DagRun).one() - dr.state == DagRunState.QUEUED - assert session.query(DagRun).count() == 1 - assert dag_maker.dag_model.next_dagrun_create_after is None + self.scheduler_job._do_scheduling(session) session.flush() - # dags_needing_dagruns query should not return any value - assert len(DagModel.dags_needing_dagruns(session).all()) == 0 - self.scheduler_job._create_dag_runs(dag_models, session) assert session.query(DagRun).count() == 1 - assert dag_maker.dag_model.next_dagrun_create_after is None - assert dag_maker.dag_model.next_dagrun == DEFAULT_DATE + dm = session.query(DagModel).one() + assert dm.next_dagrun == DEFAULT_DATE + # set dagrun to success dr = session.query(DagRun).one() dr.state = DagRunState.SUCCESS @@ -2959,12 +2989,14 @@ def test_more_runs_are_not_created_when_max_active_runs_is_reached(self, dag_mak session.merge(ti) session.merge(dr) session.flush() + # check that next_dagrun is set properly by Schedulerjob._update_dag_next_dagruns self.scheduler_job._schedule_dag_run(dr, session) session.flush() assert len(DagModel.dags_needing_dagruns(session).all()) == 1 # assert next_dagrun has been updated correctly - assert dag_maker.dag_model.next_dagrun == DEFAULT_DATE + timedelta(days=1) + dm = session.query(DagModel).one() + assert dm.next_dagrun == DEFAULT_DATE + timedelta(days=1) # assert no dagruns is created yet assert ( session.query(DagRun).filter(DagRun.state.in_([DagRunState.RUNNING, DagRunState.QUEUED])).count() @@ -3009,7 +3041,6 @@ def complete_one_dagrun(): assert DagRun.active_runs_of_dags(session=session) == {'test_dag': 3} assert model.next_dagrun == timezone.DateTime(2016, 1, 3, tzinfo=UTC) - assert model.next_dagrun_create_after is None complete_one_dagrun() diff --git a/tests/models/test_dag.py b/tests/models/test_dag.py index f7e40ef99cf70..514aba1dc0386 100644 --- a/tests/models/test_dag.py +++ b/tests/models/test_dag.py @@ -786,43 +786,6 @@ def test_bulk_write_to_db(self): for row in session.query(DagModel.last_parsed_time).all(): assert row[0] is not None - @parameterized.expand([State.RUNNING, State.QUEUED]) - def test_bulk_write_to_db_max_active_runs(self, state): - """ - Test that DagModel.next_dagrun_create_after is set to NULL when the dag cannot be created due to max - active runs being hit. - """ - dag = DAG(dag_id='test_scheduler_verify_max_active_runs', start_date=DEFAULT_DATE) - dag.max_active_runs = 1 - - DummyOperator(task_id='dummy', dag=dag, owner='airflow') - - session = settings.Session() - dag.clear() - DAG.bulk_write_to_db([dag], session) - - model = session.query(DagModel).get((dag.dag_id,)) - - assert model.next_dagrun == DEFAULT_DATE - assert model.next_dagrun_create_after == DEFAULT_DATE + timedelta(days=1) - - dr = dag.create_dagrun( - state=state, - execution_date=model.next_dagrun, - run_type=DagRunType.SCHEDULED, - session=session, - ) - assert dr is not None - DAG.bulk_write_to_db([dag]) - - model = session.query(DagModel).get((dag.dag_id,)) - # We signal "at max active runs" by saying this run is never eligible to be created - assert model.next_dagrun_create_after is None - # test that bulk_write_to_db again doesn't update next_dagrun_create_after - DAG.bulk_write_to_db([dag]) - model = session.query(DagModel).get((dag.dag_id,)) - assert model.next_dagrun_create_after is None - def test_bulk_write_to_db_has_import_error(self): """ Test that DagModel.has_import_error is set to false if no import errors. From 913ea9ef60befc62dd710b27d1d4fe4c4a0bd900 Mon Sep 17 00:00:00 2001 From: Aleksej Kirilishin Date: Fri, 4 Feb 2022 01:11:17 +0300 Subject: [PATCH 2/4] Improve tests --- tests/jobs/test_scheduler_job.py | 106 +++++++++++++------------------ 1 file changed, 45 insertions(+), 61 deletions(-) diff --git a/tests/jobs/test_scheduler_job.py b/tests/jobs/test_scheduler_job.py index afcabf376c390..04a8f8423e0fe 100644 --- a/tests/jobs/test_scheduler_job.py +++ b/tests/jobs/test_scheduler_job.py @@ -52,7 +52,7 @@ from airflow.utils.callback_requests import DagCallbackRequest from airflow.utils.file import list_py_file_paths from airflow.utils.session import create_session, provide_session -from airflow.utils.state import DagRunState, State, TaskInstanceState +from airflow.utils.state import DagRunState, State from airflow.utils.types import DagRunType from tests.test_utils.asserts import assert_queries_count from tests.test_utils.config import conf_vars, env_vars @@ -2924,85 +2924,69 @@ def test_do_schedule_max_active_runs_task_removed(self, session, dag_maker): ti.refresh_from_db(session=session) assert ti.state == State.QUEUED - def test_runs_are_created_after_max_active_runs_was_reached(self, dag_maker, caplog): + def test_runs_are_created_after_max_active_runs_was_reached(self, dag_maker, session): """ - This tests that more DagRuns are created after max_active_runs had reached - and some DagRun had finished. + Test that when creating runs once max_active_runs is reached the runs does not stick """ - with dag_maker(max_active_runs=1): - DummyOperator(task_id='task') self.scheduler_job = SchedulerJob(subdir=os.devnull) - self.scheduler_job.executor = MockExecutor(do_update=False) + self.scheduler_job.executor = MockExecutor(do_update=True) self.scheduler_job.processor_agent = mock.MagicMock(spec=DagFileProcessorAgent) - session = settings.Session() - # schedule cycle when DR created - self.scheduler_job._do_scheduling(session) - session.flush() - assert session.query(DagRun).count() == 1 - dm = session.query(DagModel).one() - assert dm.next_dagrun == DEFAULT_DATE + with dag_maker(max_active_runs=1, session=session) as dag: + # Need to use something that doesn't immediately get marked as success by the scheduler + BashOperator(task_id='task', bash_command='true') - # schedule cycle when max_active_runs reached - self.scheduler_job._do_scheduling(session) - session.flush() + dag_run = dag_maker.create_dagrun( + state=State.RUNNING, + session=session, + ) - # set dagrun to success - dr = session.query(DagRun).one() - dr.state = DagRunState.SUCCESS - ti = dr.get_task_instance('task', session) - ti.state = TaskInstanceState.SUCCESS - session.merge(ti) - session.merge(dr) - session.flush() + # Reach max_active_runs + for _ in range(3): + self.scheduler_job._do_scheduling(session) - # few cycle for calculate date and create new DR - for _ in range(2): + # Complete dagrun + # Add dag_run back in to the session (_do_scheduling does an expunge_all) + dag_run = session.merge(dag_run) + session.refresh(dag_run) + dag_run.get_task_instance(task_id='task', session=session).state = State.SUCCESS + + # create new run + for _ in range(3): self.scheduler_job._do_scheduling(session) - session.flush() - assert session.query(DagRun).count() == 2 - dm = session.query(DagModel).one() - assert dm.next_dagrun == DEFAULT_DATE + timedelta(days=1) - def test_more_runs_are_not_created_when_max_active_runs_is_reached(self, dag_maker, caplog): + # Assert that new runs has created + dag_runs = DagRun.find(dag_id=dag.dag_id, session=session) + assert len(dag_runs) == 2 + + def test_more_runs_are_not_created_when_max_active_runs_is_reached(self, dag_maker, session): """ - This tests that when max_active_runs is reached, _create_dag_runs doesn't create - more dagruns + This tests that when max_active_runs is reached, _create_dag_runs doesn't create more dagruns """ - with dag_maker(max_active_runs=1): - DummyOperator(task_id='task') self.scheduler_job = SchedulerJob(subdir=os.devnull) self.scheduler_job.executor = MockExecutor(do_update=False) self.scheduler_job.processor_agent = mock.MagicMock(spec=DagFileProcessorAgent) - session = settings.Session() - self.scheduler_job._do_scheduling(session) - session.flush() - assert session.query(DagRun).count() == 1 - dm = session.query(DagModel).one() - assert dm.next_dagrun == DEFAULT_DATE - # set dagrun to success - dr = session.query(DagRun).one() - dr.state = DagRunState.SUCCESS - ti = dr.get_task_instance('task', session) - ti.state = TaskInstanceState.SUCCESS - session.merge(ti) - session.merge(dr) - session.flush() + with dag_maker(max_active_runs=1) as dag: + # Need to use something that doesn't immediately get marked as success by the scheduler + BashOperator(task_id='dummy1', bash_command='true') - # check that next_dagrun is set properly by Schedulerjob._update_dag_next_dagruns - self.scheduler_job._schedule_dag_run(dr, session) - session.flush() - assert len(DagModel.dags_needing_dagruns(session).all()) == 1 - # assert next_dagrun has been updated correctly - dm = session.query(DagModel).one() - assert dm.next_dagrun == DEFAULT_DATE + timedelta(days=1) - # assert no dagruns is created yet - assert ( - session.query(DagRun).filter(DagRun.state.in_([DagRunState.RUNNING, DagRunState.QUEUED])).count() - == 0 + dag_maker.create_dagrun( + state=State.RUNNING, + session=session, ) + for _ in range(3): + self.scheduler_job._do_scheduling(session) + + # Assert that no more dagruns has been created + dag_runs = DagRun.find(dag_id=dag.dag_id, session=session) + assert len(dag_runs) == 1 + + # Assert that the next_dagrun hasn't been updated + dm = DagModel.get_current(dag_id=dag.dag_id, session=session) + assert dm.next_dagrun == DEFAULT_DATE + def test_max_active_runs_creation_phasing(self, dag_maker, session): """ Test that when creating runs once max_active_runs is reached that the runs come in the right order From 095394333a4fd1876695c3ab9bf0fee1abf5ee2e Mon Sep 17 00:00:00 2001 From: Aleksej Kirilishin Date: Sat, 5 Feb 2022 16:35:01 +0300 Subject: [PATCH 3/4] Reduce CPU usage --- airflow/config_templates/config.yml | 9 ++ airflow/config_templates/default_airflow.cfg | 3 + airflow/jobs/scheduler_job.py | 24 +++- airflow/models/dag.py | 6 +- tests/jobs/test_scheduler_job.py | 131 ++++++++++++------- 5 files changed, 122 insertions(+), 51 deletions(-) diff --git a/airflow/config_templates/config.yml b/airflow/config_templates/config.yml index 8ef738f4f4ad6..3b6aa30870421 100644 --- a/airflow/config_templates/config.yml +++ b/airflow/config_templates/config.yml @@ -1961,6 +1961,15 @@ type: string example: ~ default: "15" + - name: min_active_runs_check_interval + description: | + The number of seconds after which the dag that has reached ``max_active_run`` + will be checked again for the need to create new DagRuns. Keeping this + number low will increase CPU usage. + version_added: ~ + type: integer + example: ~ + default: "10" - name: triggerer description: ~ options: diff --git a/airflow/config_templates/default_airflow.cfg b/airflow/config_templates/default_airflow.cfg index 520ab4442850d..fbc325945b94b 100644 --- a/airflow/config_templates/default_airflow.cfg +++ b/airflow/config_templates/default_airflow.cfg @@ -984,6 +984,9 @@ dependency_detector = airflow.serialization.serialized_objects.DependencyDetecto # How often to check for expired trigger requests that have not run yet. trigger_timeout_check_interval = 15 +# How often to check the need to create new DagRuns after ``max_active_run`` has been reached +min_active_runs_check_interval = 10 + [triggerer] # How many triggers a single Triggerer will run at once, by default. default_capacity = 1000 diff --git a/airflow/jobs/scheduler_job.py b/airflow/jobs/scheduler_job.py index e8f4531882b57..795f38bbb9ca8 100644 --- a/airflow/jobs/scheduler_job.py +++ b/airflow/jobs/scheduler_job.py @@ -139,6 +139,8 @@ def __init__( self.dagbag = DagBag(dag_folder=self.subdir, read_dags_from_db=True, load_op_links=False) + self.last_check_time = {} + if conf.getboolean('smart_sensor', 'use_smart_sensor'): compatible_sensors = set( map(lambda l: l.strip(), conf.get('smart_sensor', 'sensors_enabled').split(',')) @@ -884,11 +886,28 @@ def _get_next_dagruns_to_examine(self, state: DagRunState, session: Session): """Get Next DagRuns to Examine with retries""" return DagRun.next_dagruns_to_examine(state, session) + def _get_recently_checked_dags(self, current_time=None) -> List: + """Get dags for which `max_active_runs` has been reached recently""" + if current_time is None: + current_time = timezone.utcnow() + check_interval = timedelta( + seconds=conf.getint('scheduler', 'min_active_runs_check_interval', fallback=10) + ) + + skip_dags = [] + for dag_id in self.last_check_time: + if self.last_check_time[dag_id] + check_interval > current_time: + skip_dags.append(dag_id) + return skip_dags + @retry_db_transaction def _create_dagruns_for_dags(self, guard, session): """Find Dag Models needing DagRuns and Create Dag Runs with retries in case of OperationalError""" - query = DagModel.dags_needing_dagruns(session) - self._create_dag_runs(query.all(), session) + # Reduce the number of dags due to high CPU usage + skip_dags = self._get_recently_checked_dags() + self.log.debug("Skipping dags after max_active_runs has been reached: %s", skip_dags) + dag_models = DagModel.dags_needing_dagruns(session, skip_dags).all() + self._create_dag_runs(dag_models, session) # commit the session - Release the write lock on DagModel table. guard.commit() @@ -972,6 +991,7 @@ def _should_update_dag_next_dagruns(self, dag, dag_model: DagModel, total_active total_active_runs, dag.max_active_runs, ) + self.last_check_time[dag_model.dag_id] = timezone.utcnow() return False return True diff --git a/airflow/models/dag.py b/airflow/models/dag.py index a300f97c94e58..7385414b8cc8b 100644 --- a/airflow/models/dag.py +++ b/airflow/models/dag.py @@ -2849,7 +2849,7 @@ def deactivate_deleted_dags(cls, alive_dag_filelocs: List[str], session=NEW_SESS continue @classmethod - def dags_needing_dagruns(cls, session: Session): + def dags_needing_dagruns(cls, session: Session, skip_dags=None): """ Return (and lock) a list of Dag objects that are due to create a new DagRun. @@ -2861,6 +2861,9 @@ def dags_needing_dagruns(cls, session: Session): # We limit so that _one_ scheduler doesn't try to do all the creation # of dag runs + if skip_dags is None: + skip_dags = [] + query = ( session.query(cls) .filter( @@ -2868,6 +2871,7 @@ def dags_needing_dagruns(cls, session: Session): cls.is_active == expression.true(), cls.has_import_errors == expression.false(), cls.next_dagrun_create_after <= func.now(), + cls.dag_id.notin_(skip_dags), ) .order_by(cls.next_dagrun_create_after) .limit(cls.NUM_DAGS_PER_DAGRUN_QUERY) diff --git a/tests/jobs/test_scheduler_job.py b/tests/jobs/test_scheduler_job.py index 04a8f8423e0fe..844d375aa31a4 100644 --- a/tests/jobs/test_scheduler_job.py +++ b/tests/jobs/test_scheduler_job.py @@ -2928,36 +2928,37 @@ def test_runs_are_created_after_max_active_runs_was_reached(self, dag_maker, ses """ Test that when creating runs once max_active_runs is reached the runs does not stick """ - self.scheduler_job = SchedulerJob(subdir=os.devnull) - self.scheduler_job.executor = MockExecutor(do_update=True) - self.scheduler_job.processor_agent = mock.MagicMock(spec=DagFileProcessorAgent) + with conf_vars({('scheduler', 'min_active_runs_check_interval'): '0'}): + self.scheduler_job = SchedulerJob(subdir=os.devnull) + self.scheduler_job.executor = MockExecutor(do_update=True) + self.scheduler_job.processor_agent = mock.MagicMock(spec=DagFileProcessorAgent) - with dag_maker(max_active_runs=1, session=session) as dag: - # Need to use something that doesn't immediately get marked as success by the scheduler - BashOperator(task_id='task', bash_command='true') + with dag_maker(max_active_runs=1, session=session) as dag: + # Need to use something that doesn't immediately get marked as success by the scheduler + BashOperator(task_id='task', bash_command='true') - dag_run = dag_maker.create_dagrun( - state=State.RUNNING, - session=session, - ) + dag_run = dag_maker.create_dagrun( + state=State.RUNNING, + session=session, + ) - # Reach max_active_runs - for _ in range(3): - self.scheduler_job._do_scheduling(session) + # Reach max_active_runs + for _ in range(3): + self.scheduler_job._do_scheduling(session) - # Complete dagrun - # Add dag_run back in to the session (_do_scheduling does an expunge_all) - dag_run = session.merge(dag_run) - session.refresh(dag_run) - dag_run.get_task_instance(task_id='task', session=session).state = State.SUCCESS + # Complete dagrun + # Add dag_run back in to the session (_do_scheduling does an expunge_all) + dag_run = session.merge(dag_run) + session.refresh(dag_run) + dag_run.get_task_instance(task_id='task', session=session).state = State.SUCCESS - # create new run - for _ in range(3): - self.scheduler_job._do_scheduling(session) + # create new run + for _ in range(3): + self.scheduler_job._do_scheduling(session) - # Assert that new runs has created - dag_runs = DagRun.find(dag_id=dag.dag_id, session=session) - assert len(dag_runs) == 2 + # Assert that new runs has created + dag_runs = DagRun.find(dag_id=dag.dag_id, session=session) + assert len(dag_runs) == 2 def test_more_runs_are_not_created_when_max_active_runs_is_reached(self, dag_maker, session): """ @@ -3007,39 +3008,40 @@ def complete_one_dagrun(): self.clean_db() - with dag_maker(max_active_runs=3, session=session) as dag: - # Need to use something that doesn't immediately get marked as success by the scheduler - BashOperator(task_id='task', bash_command='true') + with conf_vars({('scheduler', 'min_active_runs_check_interval'): '0'}): + with dag_maker(max_active_runs=3, session=session) as dag: + # Need to use something that doesn't immediately get marked as success by the scheduler + BashOperator(task_id='task', bash_command='true') - self.scheduler_job = SchedulerJob(subdir=os.devnull) - self.scheduler_job.executor = MockExecutor(do_update=True) - self.scheduler_job.processor_agent = mock.MagicMock(spec=DagFileProcessorAgent) + self.scheduler_job = SchedulerJob(subdir=os.devnull) + self.scheduler_job.executor = MockExecutor(do_update=True) + self.scheduler_job.processor_agent = mock.MagicMock(spec=DagFileProcessorAgent) - DagModel.dags_needing_dagruns(session).all() - for _ in range(3): - self.scheduler_job._do_scheduling(session) + DagModel.dags_needing_dagruns(session).all() + for _ in range(3): + self.scheduler_job._do_scheduling(session) - model: DagModel = session.query(DagModel).get(dag.dag_id) + model: DagModel = session.query(DagModel).get(dag.dag_id) - # Pre-condition - assert DagRun.active_runs_of_dags(session=session) == {'test_dag': 3} + # Pre-condition + assert DagRun.active_runs_of_dags(session=session) == {'test_dag': 3} - assert model.next_dagrun == timezone.DateTime(2016, 1, 3, tzinfo=UTC) + assert model.next_dagrun == timezone.DateTime(2016, 1, 3, tzinfo=UTC) - complete_one_dagrun() + complete_one_dagrun() - assert DagRun.active_runs_of_dags(session=session) == {'test_dag': 3} + assert DagRun.active_runs_of_dags(session=session) == {'test_dag': 3} - for _ in range(5): - self.scheduler_job._do_scheduling(session) - complete_one_dagrun() - model: DagModel = session.query(DagModel).get(dag.dag_id) + for _ in range(5): + self.scheduler_job._do_scheduling(session) + complete_one_dagrun() + model: DagModel = session.query(DagModel).get(dag.dag_id) - expected_execution_dates = [datetime.datetime(2016, 1, d, tzinfo=timezone.utc) for d in range(1, 6)] - dagrun_execution_dates = [ - dr.execution_date for dr in session.query(DagRun).order_by(DagRun.execution_date).all() - ] - assert dagrun_execution_dates == expected_execution_dates + expected_execution_dates = [datetime.datetime(2016, 1, d, tzinfo=timezone.utc) for d in range(1, 6)] + dagrun_execution_dates = [ + dr.execution_date for dr in session.query(DagRun).order_by(DagRun.execution_date).all() + ] + assert dagrun_execution_dates == expected_execution_dates def test_do_schedule_max_active_runs_and_manual_trigger(self, dag_maker): """ @@ -3828,3 +3830,36 @@ def test_catchup_works_correctly(self, dag_maker): .filter(DagRun.execution_date != DEFAULT_DATE) # exclude the first run .scalar() ) > (timezone.utcnow() - timedelta(days=2)) + + def test_get_recently_checked_dags(self, dag_maker): + with conf_vars({('scheduler', 'min_active_runs_check_interval'): '10'}): + self.scheduler_job = SchedulerJob(subdir=os.devnull) + + self.scheduler_job.last_check_time = {} + assert self.scheduler_job._get_recently_checked_dags() == [] + + self.scheduler_job.last_check_time = { + 'test_dag': timezone.datetime(2016, 1, 1, 20, 59, 59) + } + assert self.scheduler_job._get_recently_checked_dags( + timezone.datetime(2016, 1, 1, 21, 1, 0) + ) == [] + + self.scheduler_job.last_check_time = { + 'test_dag': timezone.datetime(2016, 1, 1, 20, 59, 59) + } + assert self.scheduler_job._get_recently_checked_dags( + timezone.datetime(2016, 1, 1, 21, 0, 5) + ) == ['test_dag'] + + self.scheduler_job.last_check_time = { + 'test_dag0': timezone.datetime(2016, 1, 1, 20, 59, 50), + 'test_dag1': timezone.datetime(2016, 1, 1, 20, 59, 51), + 'test_dag2': timezone.datetime(2016, 1, 1, 0, 0, 0), + 'test_dag3': timezone.datetime(2016, 1, 1, 20, 59, 55), + 'test_dag4': timezone.datetime(2015, 1, 1, 20, 59, 59), + 'test_dag5': timezone.datetime(2016, 1, 1, 20, 59, 59), + } + assert sorted(self.scheduler_job._get_recently_checked_dags( + timezone.datetime(2016, 1, 1, 21, 0, 0) + )) == ['test_dag1', 'test_dag3', 'test_dag5'] From 039ba528b870e3935780d99f02370d816feb1b61 Mon Sep 17 00:00:00 2001 From: Aleksej Kirilishin Date: Sat, 5 Feb 2022 16:43:03 +0300 Subject: [PATCH 4/4] Files were modified by black --- airflow/config_templates/default_airflow.cfg | 4 ++- airflow/jobs/scheduler_job.py | 4 +-- tests/jobs/test_scheduler_job.py | 30 +++++++++----------- 3 files changed, 19 insertions(+), 19 deletions(-) diff --git a/airflow/config_templates/default_airflow.cfg b/airflow/config_templates/default_airflow.cfg index fbc325945b94b..f6f68283d3a18 100644 --- a/airflow/config_templates/default_airflow.cfg +++ b/airflow/config_templates/default_airflow.cfg @@ -984,7 +984,9 @@ dependency_detector = airflow.serialization.serialized_objects.DependencyDetecto # How often to check for expired trigger requests that have not run yet. trigger_timeout_check_interval = 15 -# How often to check the need to create new DagRuns after ``max_active_run`` has been reached +# The number of seconds after which the dag that has reached ``max_active_run`` +# will be checked again for the need to create new DagRuns. Keeping this +# number low will increase CPU usage. min_active_runs_check_interval = 10 [triggerer] diff --git a/airflow/jobs/scheduler_job.py b/airflow/jobs/scheduler_job.py index 795f38bbb9ca8..f7338366e2b80 100644 --- a/airflow/jobs/scheduler_job.py +++ b/airflow/jobs/scheduler_job.py @@ -25,7 +25,7 @@ import time import warnings from collections import defaultdict -from datetime import timedelta +from datetime import datetime, timedelta from typing import Collection, DefaultDict, Dict, Iterator, List, Optional, Tuple from sqlalchemy import and_, func, not_, or_, text, tuple_ @@ -139,7 +139,7 @@ def __init__( self.dagbag = DagBag(dag_folder=self.subdir, read_dags_from_db=True, load_op_links=False) - self.last_check_time = {} + self.last_check_time: Dict[str, datetime] = {} if conf.getboolean('smart_sensor', 'use_smart_sensor'): compatible_sensors = set( diff --git a/tests/jobs/test_scheduler_job.py b/tests/jobs/test_scheduler_job.py index 844d375aa31a4..e12960225242f 100644 --- a/tests/jobs/test_scheduler_job.py +++ b/tests/jobs/test_scheduler_job.py @@ -3037,7 +3037,9 @@ def complete_one_dagrun(): complete_one_dagrun() model: DagModel = session.query(DagModel).get(dag.dag_id) - expected_execution_dates = [datetime.datetime(2016, 1, d, tzinfo=timezone.utc) for d in range(1, 6)] + expected_execution_dates = [ + datetime.datetime(2016, 1, d, tzinfo=timezone.utc) for d in range(1, 6) + ] dagrun_execution_dates = [ dr.execution_date for dr in session.query(DagRun).order_by(DagRun.execution_date).all() ] @@ -3838,19 +3840,15 @@ def test_get_recently_checked_dags(self, dag_maker): self.scheduler_job.last_check_time = {} assert self.scheduler_job._get_recently_checked_dags() == [] - self.scheduler_job.last_check_time = { - 'test_dag': timezone.datetime(2016, 1, 1, 20, 59, 59) - } - assert self.scheduler_job._get_recently_checked_dags( - timezone.datetime(2016, 1, 1, 21, 1, 0) - ) == [] + self.scheduler_job.last_check_time = {'test_dag': timezone.datetime(2016, 1, 1, 20, 59, 59)} + assert ( + self.scheduler_job._get_recently_checked_dags(timezone.datetime(2016, 1, 1, 21, 1, 0)) == [] + ) - self.scheduler_job.last_check_time = { - 'test_dag': timezone.datetime(2016, 1, 1, 20, 59, 59) - } - assert self.scheduler_job._get_recently_checked_dags( - timezone.datetime(2016, 1, 1, 21, 0, 5) - ) == ['test_dag'] + self.scheduler_job.last_check_time = {'test_dag': timezone.datetime(2016, 1, 1, 20, 59, 59)} + assert self.scheduler_job._get_recently_checked_dags(timezone.datetime(2016, 1, 1, 21, 0, 5)) == [ + 'test_dag' + ] self.scheduler_job.last_check_time = { 'test_dag0': timezone.datetime(2016, 1, 1, 20, 59, 50), @@ -3860,6 +3858,6 @@ def test_get_recently_checked_dags(self, dag_maker): 'test_dag4': timezone.datetime(2015, 1, 1, 20, 59, 59), 'test_dag5': timezone.datetime(2016, 1, 1, 20, 59, 59), } - assert sorted(self.scheduler_job._get_recently_checked_dags( - timezone.datetime(2016, 1, 1, 21, 0, 0) - )) == ['test_dag1', 'test_dag3', 'test_dag5'] + assert sorted( + self.scheduler_job._get_recently_checked_dags(timezone.datetime(2016, 1, 1, 21, 0, 0)) + ) == ['test_dag1', 'test_dag3', 'test_dag5']