hussein-awala commented on code in PR #70307:
URL: https://github.com/apache/airflow/pull/70307#discussion_r4207563752
##########
airflow-core/src/airflow/jobs/scheduler_job_runner.py:
##########
@@ -968,12 +982,26 @@ def _select_task_instances_to_queue(
or len(starved_tasks_task_dagrun_concurrency) >
num_starved_tasks_task_dagrun_concurrency
)
- if is_done or not found_new_filters:
+ if len(executable_tis) >= max_tis or not found_new_filters:
break
+ if executable_tis:
+ if num_refill_queries >= MAX_TI_REFILL_QUERIES_PER_LOOP:
+ self.log.debug(
+ "Selected %s of %s task instances; reached the limit
of %s refill queries.",
+ len(executable_tis),
+ max_tis,
+ MAX_TI_REFILL_QUERIES_PER_LOOP,
+ )
+
stats.incr("scheduler.critical_section_refill_limit_reached")
+ break
+ num_refill_queries += 1
Review Comment:
The new refill path makes the existing `starved_dags` exclusion too
broad: the concurrency check is per Dag run, but reaching that limit excludes
every run of the Dag from subsequent queries.
For example, with `max_tis=2`, two free pool slots, and one task already
running in A/run1 (`max_active_tasks=2`), consider these scheduled candidates:
| Task | Run | Priority |
|---|---|---:|
| `a1` | A/run1 | 100 |
| `a2` | A/run1 | 99 |
| `b` | A/run2 | 98 |
| `c` | Another Dag | 1 |
The first query selects `a1`, rejects `a2` because run1 is now virtually
full, and adds A to `starved_dags`. The refill then excludes runnable `b` and
queues `c`, consuming the remaining pool slot.
Although the broad exclusion predates this change, this consequence is
new: previously, selection stopped after `a1`. On the next scheduler loop, the
SQL concurrency check excluded only saturated run1, allowing `b` to receive the
remaining slot before `c`.
Could we track saturated runs by `(dag_id, run_id)` rather than excluding
the entire Dag, and add a regression covering this multiple-run scenario?
To confirm this, I added the following regression test to
`TestSchedulerJob` in `airflow-core/tests/unit/jobs/test_scheduler_job.py` and
ran the identical test against upstream `main` and this PR.
The test allows up to two scheduling passes, so it checks **priority
preservation**, rather than requiring the scheduler to fill the batch in one
pass. Both available pool slots should go to the higher-priority tasks from
`multiple_runs`, even though they belong to
different runs.
```python
def
test_select_task_instances_to_queue_preserves_priority_across_dag_runs(
self, dag_maker, mock_executors, session
):
with dag_maker(dag_id="multiple_runs", max_active_tasks=2,
session=session):
EmptyOperator(task_id="running")
EmptyOperator(task_id="fills_run", priority_weight=100)
EmptyOperator(task_id="blocked_in_run", priority_weight=99)
EmptyOperator(task_id="other_run", priority_weight=98)
first_run = dag_maker.create_dagrun(run_type=DagRunType.SCHEDULED)
second_run = dag_maker.create_dagrun_after(first_run,
run_type=DagRunType.SCHEDULED)
with dag_maker(dag_id="lower_priority", session=session):
EmptyOperator(task_id="competitor", priority_weight=1)
competing_run =
dag_maker.create_dagrun(run_type=DagRunType.SCHEDULED)
first_run.get_task_instance("running", session=session).state =
State.RUNNING
fills_run_ti = first_run.get_task_instance("fills_run",
session=session)
fills_run_ti.state = State.SCHEDULED
first_run.get_task_instance("blocked_in_run",
session=session).state = State.SCHEDULED
other_run_ti = second_run.get_task_instance("other_run",
session=session)
other_run_ti.state = State.SCHEDULED
competing_run.get_task_instance("competitor",
session=session).state = State.SCHEDULED
session.flush()
expected_keys = [fills_run_ti.key, other_run_ti.key]
pools = make_pool_stats(total=3, running=1)
self.job_runner = SchedulerJobRunner(job=Job())
queued_tis = []
# Allow another pass so this checks priority, not whether one
pass fills the batch.
for _ in range(2):
max_tis = int(pools["default_pool"]["open"])
if not max_tis:
break
queued_tis.extend(
self.job_runner._select_task_instances_to_queue(max_tis,
pools, set(), session=session)
)
assert [ti.key for ti in queued_tis] == expected_keys
```
| Revision | Result |
|---|---|
| Upstream `main` at `73a4ffa3a5` | **1 passed** in 12.80s |
| PR head at `3b15392fbf` | **1 failed** in 13.22s |
The PR run fails at the final assertion. Summarizing the returned
task-instance keys:
```text
Expected:
multiple_runs.fills_run
multiple_runs.other_run
Actual:
multiple_runs.fills_run
lower_priority.competitor
```
This confirms that saturating the first run causes the refill to exclude
the entire Dag, skipping the runnable higher-priority task in its second run
and consuming the remaining pool slot with the lower-priority competitor.
--
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]