From eca916ae4a5427305d50f5cf429b506308e06244 Mon Sep 17 00:00:00 2001 From: Ephraim Anierobi Date: Tue, 8 Feb 2022 15:35:39 +0100 Subject: [PATCH 1/4] Fix max_active_runs=1 not scheduling runs when min_file_process_interval is high The finished dagrun was still being seen as running when we call dag.get_num_active_runs because the session was not flushed. This PR fixes it --- airflow/jobs/scheduler_job.py | 1 + tests/jobs/test_scheduler_job.py | 35 ++++++++++++++++++++++++++++++++ 2 files changed, 36 insertions(+) diff --git a/airflow/jobs/scheduler_job.py b/airflow/jobs/scheduler_job.py index 7f7fcf05e7c81..dd0b40b66b129 100644 --- a/airflow/jobs/scheduler_job.py +++ b/airflow/jobs/scheduler_job.py @@ -1096,6 +1096,7 @@ def _schedule_dag_run( # TODO[HA]: Rename update_state -> schedule_dag_run, ?? something else? schedulable_tis, callback_to_run = dag_run.update_state(session=session, execute_callbacks=False) if dag_run.state in State.finished: + session.flush() # to update the dag_run state active_runs = dag.get_num_active_runs(only_running=False, session=session) # Work out if we should allow creating a new DagRun now? if self._should_update_dag_next_dagruns(dag, dag_model, active_runs): diff --git a/tests/jobs/test_scheduler_job.py b/tests/jobs/test_scheduler_job.py index 72d60d2e22ee1..d959e3fba1500 100644 --- a/tests/jobs/test_scheduler_job.py +++ b/tests/jobs/test_scheduler_job.py @@ -1259,6 +1259,41 @@ def test_queued_dagruns_stops_creating_when_max_active_is_reached(self, dag_make assert session.query(DagRun.state).filter(DagRun.state == State.QUEUED).count() == 0 assert orm_dag.next_dagrun_create_after is None + def test_runs_are_created_after_max_active_runs_was_reached(self, dag_maker, session): + """ + 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 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, + ) + + # 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 + + # 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 + def test_dagrun_timeout_verify_max_active_runs(self, dag_maker): """ Test if a a dagrun will not be scheduled if max_dag_runs From d512e2d2ecec898eabdf8378d4ef9bc65e14c178 Mon Sep 17 00:00:00 2001 From: Ephraim Anierobi Date: Mon, 14 Feb 2022 04:39:37 +0100 Subject: [PATCH 2/4] flush session when dagrun is completed --- airflow/jobs/scheduler_job.py | 1 - airflow/models/dagrun.py | 1 + 2 files changed, 1 insertion(+), 1 deletion(-) diff --git a/airflow/jobs/scheduler_job.py b/airflow/jobs/scheduler_job.py index dd0b40b66b129..7f7fcf05e7c81 100644 --- a/airflow/jobs/scheduler_job.py +++ b/airflow/jobs/scheduler_job.py @@ -1096,7 +1096,6 @@ def _schedule_dag_run( # TODO[HA]: Rename update_state -> schedule_dag_run, ?? something else? schedulable_tis, callback_to_run = dag_run.update_state(session=session, execute_callbacks=False) if dag_run.state in State.finished: - session.flush() # to update the dag_run state active_runs = dag.get_num_active_runs(only_running=False, session=session) # Work out if we should allow creating a new DagRun now? if self._should_update_dag_next_dagruns(dag, dag_model, active_runs): diff --git a/airflow/models/dagrun.py b/airflow/models/dagrun.py index 5170ad376a592..a06a2120c0e3a 100644 --- a/airflow/models/dagrun.py +++ b/airflow/models/dagrun.py @@ -608,6 +608,7 @@ def update_state( self.data_interval_end, self.dag_hash, ) + session.flush() self._emit_true_scheduling_delay_stats_for_finished_state(finished_tis) self._emit_duration_stats_for_finished_state() From 7b45fa3fd715fe2462e98862b72c22c2661bd967 Mon Sep 17 00:00:00 2001 From: Ephraim Anierobi Date: Tue, 22 Feb 2022 09:27:27 +0100 Subject: [PATCH 3/4] move session.flush down to after merge --- airflow/models/dagrun.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/airflow/models/dagrun.py b/airflow/models/dagrun.py index a06a2120c0e3a..b2723355335fb 100644 --- a/airflow/models/dagrun.py +++ b/airflow/models/dagrun.py @@ -608,12 +608,12 @@ def update_state( self.data_interval_end, self.dag_hash, ) - session.flush() self._emit_true_scheduling_delay_stats_for_finished_state(finished_tis) self._emit_duration_stats_for_finished_state() session.merge(self) + session.flush() return schedulable_tis, callback From d7bf47b404aa84e1a0ce48717d44d44c1e526b45 Mon Sep 17 00:00:00 2001 From: Ephraim Anierobi Date: Wed, 23 Feb 2022 18:54:26 +0100 Subject: [PATCH 4/4] add reason why we don't flush down after merge --- airflow/models/dagrun.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/airflow/models/dagrun.py b/airflow/models/dagrun.py index b2723355335fb..36014fd7cf426 100644 --- a/airflow/models/dagrun.py +++ b/airflow/models/dagrun.py @@ -608,12 +608,13 @@ def update_state( self.data_interval_end, self.dag_hash, ) + session.flush() self._emit_true_scheduling_delay_stats_for_finished_state(finished_tis) self._emit_duration_stats_for_finished_state() session.merge(self) - session.flush() + # We do not flush here for performance reasons(It increases queries count by +20) return schedulable_tis, callback