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]

Reply via email to