kaxil commented on code in PR #73692:
URL: https://github.com/apache/airflow/pull/73692#discussion_r4111084198
##########
airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py:
##########
@@ -199,6 +199,7 @@ class DAGDetailsResponse(DAGResponse):
"dag_run_timeout": "dagrun_timeout",
"last_parsed": "last_loaded",
"template_search_path": "template_searchpath",
+ "is_at_max_active_runs": "exceeds_max_non_backfill",
Review Comment:
This aliases a scheduler cache, and nothing refreshes that cache when a run
is triggered. `exceeds_max_non_backfill` is only written by the dag-processor
on parse (`collection.py:704`) and by `_set_exceeds_max_active_runs`, which the
scheduler calls on scheduled-run creation, dagrun timeout and run finish.
`trigger_dag_run` never touches it. So for a `schedule=None` Dag with
`max_active_runs=1`, triggering a run leaves this `false` until the next parse
(up to `min_file_process_interval`, 30s by default), and for the whole run if
it finishes before that parse. Marking a run failed through the API has the
reverse problem and leaves it `true`. That's the "check capacity before
triggering another run" case from the thread above.
`get_dag_details` already runs a count query for `active_runs_count`, so
computing this at request time from
`DagRun.active_runs_of_dags(dag_ids=[dag_id], exclude_backfill=True,
session=session)` against `dag_model.max_active_runs` would keep it accurate
without depending on the parse cycle. A test that triggers a run after parsing
and reads details without reparsing would cover it; the current one sets the
column directly.
##########
airflow-core/src/airflow/dag_processing/collection.py:
##########
@@ -170,15 +170,21 @@ class _RunInfo(NamedTuple):
num_active_runs: int
@classmethod
- def calculate(cls, dag: LazyDeserializedDAG, *, session: Session) -> Self:
+ def calculate(cls, dag: LazyDeserializedDAG, *, num_active_runs: int,
session: Session) -> Self:
Review Comment:
`calculate` now takes `num_active_runs` only to hand it back in the tuple,
and the comment below says it "must always be computed" even though this method
no longer computes it. `update_dags` is the only caller, so reading
`active_run_counts.get(dag_id, 0)` there at line 704 and leaving `calculate` to
the latest-run lookup would be simpler.
##########
airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py:
##########
@@ -221,6 +222,7 @@ class DAGDetailsResponse(DAGResponse):
owner_links: dict[str, str] | None = None
is_favorite: bool = False
active_runs_count: int = 0
+ is_at_max_active_runs: bool
Review Comment:
Could this get a `Field(description=...)` like `is_backfillable` has? It
counts differently from `active_runs_count` right above it. That one is RUNNING
only and includes backfill runs, while this is RUNNING + QUEUED with backfill
excluded. So a Dag with `max_active_runs=1` and one running backfill returns
`active_runs_count: 1` next to `is_at_max_active_runs: false`, and a Dag with
one queued run returns `0` next to `true`. Saying what's counted would help the
#73693 tooltip explain that.
--
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]