ashb commented on code in PR #73689:
URL: https://github.com/apache/airflow/pull/73689#discussion_r4103245440
##########
airflow-core/src/airflow/serialization/definitions/dag.py:
##########
@@ -712,6 +712,24 @@ def create_dagrun(
if params_dag.deadline:
self._process_dagrun_deadline_alerts(orm_dagrun, session)
+ if state == DagRunState.QUEUED and backfill_id is None and
self.max_active_runs:
+ num_running = (
+ session.scalar(
+ select(func.count())
+ .select_from(DagRun)
+ .where(DagRun.dag_id == self.dag_id, DagRun.state ==
DagRunState.RUNNING)
+ )
+ or 0
+ )
+ if num_running >= self.max_active_runs:
+ log.info(
+ "created DagRun will not be scheduled yet, dag is at
max_active_runs",
+ dag_id=self.dag_id,
+ run_id=run_id,
+ active_runs=num_running,
+ max_active_runs=self.max_active_runs,
+ )
+
Review Comment:
I don't know about this living in here. My gut says this belongs with the
caller, not here. (create_dagrun did what it was told, it created it. That it
won't be scheduled is nothing to do with creating a dagrun)
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]