Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,11 +5,13 @@
### New features

* The server scheduler now contains a safety limit for computation, configurable via `--scheduler-time-limit` (default: 5s)
* Better scheduling policy (prefill) for heterogenous clusters

### Fixes

* Fixed some occasional greedy backfilling in server scheduler
* Fixed server crash in a specific situation when an unschedulable high-priority task occurs

* Fixed server crash caused by invalid handling of prefill

## v0.26.2

Expand Down
71 changes: 50 additions & 21 deletions crates/tako/src/internal/scheduler/mapping.rs
Original file line number Diff line number Diff line change
Expand Up @@ -161,10 +161,15 @@ fn process_proactive_filling(core: &mut Core, mapping: &mut WorkerTaskMapping) {
task_map,
worker_map,
task_queues,
request_map: _,
request_map,
scheduler_state,
..
} = core.split_mut();
let max_prefill = scheduler_state.config.proactive_filling_max as u64;
if max_prefill == 0 {
// Prefill explicitly disabled.
return;
}
let top_priority = task_queues.top_priority();
for queue in task_queues.iter_mut() {
if queue.top_priority() != Some(top_priority) {
Expand All @@ -176,6 +181,13 @@ fn process_proactive_filling(core: &mut Core, mapping: &mut WorkerTaskMapping) {
if size == 0 {
continue;
}
let rqv = request_map.get(queue.resource_rq_id);
let max_capacity = worker_map
.get_workers()
.map(|w| w.resources.task_max_count(rqv))
.max()
.unwrap_or(1)
.max(1) as u64;
let workers: Vec<_> = worker_map
.values_mut()
.filter(|worker| {
Expand Down Expand Up @@ -207,28 +219,45 @@ fn process_proactive_filling(core: &mut Core, mapping: &mut WorkerTaskMapping) {
if workers.is_empty() {
continue;
}
let prefill_size =
(size / workers.len() as u32).min(scheduler_state.config.proactive_filling_max);
if prefill_size == 0 {
continue;
}
for worker in workers {
let tasks = queue.take_tasks_for_prefill(prefill_size);
for task_id in &tasks {
let capacities: Vec<u64> = workers
.iter()
.map(|w| w.resources.task_max_count(rqv).max(1) as u64)
.collect();
let total_capacity: u64 = capacities.iter().sum();
let mut return_back = Vec::new();
for (worker, capacity) in workers.into_iter().zip(capacities) {
// The shares sum to at most `size`, so the queue entry we are drawing from is
// never exhausted before the last worker.
let share = size as u64 * capacity / total_capacity;
let depth = (max_prefill * capacity / max_capacity).max(1);
let prefill_size = share.min(depth);
if prefill_size == 0 {
continue;
}
let prefills = &mut mapping.workers.entry(worker.id).or_default().prefills;
for _ in 0..prefill_size {
let task_id = queue.take_one().unwrap();
log::debug!("Prefiling task={task_id} to worker={}", worker.id);
let task = task_map.get_task_mut(*task_id);
assert!(task.is_waiting());
task.state = TaskRuntimeState::Prefilled {
worker_id: worker.id,
};
worker.insert_prefill_task(*task_id);
let task = task_map.get_task_mut(task_id);
if task.is_waiting() {
task.state = TaskRuntimeState::Prefilled {
worker_id: worker.id,
};
worker.insert_prefill_task(task_id);
queue.insert_prefill(task_id, top_priority, prefill_size as usize);
prefills.push(task_id);
} else {
// This can happen when task is in retracting, and it should be queite rare
log::debug!(
"Task is not in waiting state ({:?}) back to the queue.",
task.state
);
return_back.push(task_id);
}
}
mapping
.workers
.entry(worker.id)
.or_default()
.prefills
.extend(tasks);
}
for task_id in return_back {
queue.return_back(task_id, top_priority);
}
}
}
Expand Down
13 changes: 8 additions & 5 deletions crates/tako/src/internal/scheduler/solver.rs
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,7 @@ pub(crate) fn run_scheduling_solver(
task_queues: _,
request_map,
worker_groups,
scheduler_state: scheduler_cache,
scheduler_state,
..
} = core.split();
if request_map.is_empty() {
Expand Down Expand Up @@ -156,7 +156,7 @@ pub(crate) fn run_scheduling_solver(
&& worker.is_capable_to_run_rqv(rqv, now)
&& let Some(a) = worker.sn_assignment()
{
let weight = w_idx as f64 / (n_workers * 100) as f64;
let weight = -((n_workers - w_idx) as f64 / (n_workers * 1024) as f64);
solver.set_name(|| format!("R{}:{}", worker.id, batch.resource_rq_id));
let v = solver.add_bool_variable(weight);
tasks_count_vars
Expand Down Expand Up @@ -292,13 +292,16 @@ pub(crate) fn run_scheduling_solver(
if !w.is_capable_to_run_rqv(blocker_rqv, now) {
continue;
}
let gap = scheduler_cache.gap_cache.get_gap(
let gap = scheduler_state.gap_cache.get_gap(
*blocker_rq_id,
batch.resource_rq_id,
&w.resources,
sn_assignment.assigned_tasks.iter().map(|task_id| {
let t = task_map.get_task(*task_id);
(t.resource_rq_id, t.rv_id().unwrap())
(
t.resource_rq_id,
t.assigned_placement(&scheduler_state.redirects).unwrap().1,
)
}),
request_map,
);
Expand Down Expand Up @@ -430,7 +433,7 @@ pub(crate) fn run_scheduling_solver(
}

let mut result = SchedulingSolution::default();
let Some((solution, is_optimal)) = solver.solve_bounded(scheduler_cache.config.mip_time_limit)
let Some((solution, is_optimal)) = solver.solve_bounded(scheduler_state.config.mip_time_limit)
else {
return result;
};
Expand Down
51 changes: 40 additions & 11 deletions crates/tako/src/internal/scheduler/taskqueue.rs
Original file line number Diff line number Diff line change
Expand Up @@ -151,6 +151,11 @@ impl TaskQueue {
}
}

#[inline]
pub fn return_back(&mut self, task_id: TaskId, priority: Priority) {
self.add(task_id, priority);
}

fn add(&mut self, task_id: TaskId, priority: Priority) {
match self.queue.entry(Reverse(priority)) {
Entry::Vacant(e) => {
Expand Down Expand Up @@ -216,6 +221,36 @@ impl TaskQueue {
}
}

/// Removes the task if it is in the queue and returns whether it was removed.
/// A retracting task may have been already taken out of the queue by the scheduler.
pub fn remove_if_queued(&mut self, task_id: TaskId, priority: Priority) -> bool {
if let Some((p, ts)) = &mut self.prefill
&& priority == *p
&& ts.remove(&task_id)
{
return true;
}
match self.queue.entry(Reverse(priority)) {
Entry::Vacant(_) => false,
Entry::Occupied(mut e) => match e.get_mut() {
OneOrMoreTaskIds::One(v) => {
if *v != task_id {
return false;
}
e.remove();
true
}
OneOrMoreTaskIds::More(tasks) => {
let found = tasks.remove(&task_id);
if tasks.is_empty() {
e.remove();
}
found
}
},
}
}

pub fn shrink_to_fit(&mut self) {
// Do nothing
}
Expand Down Expand Up @@ -301,20 +336,14 @@ impl TaskQueue {
}
}

pub fn take_tasks_for_prefill(&mut self, mut count: u32) -> Vec<TaskId> {
let entry = self.queue.first_entry().unwrap();
let mut result = Vec::with_capacity(count as usize);
let priority = entry.key().0;
take_from_entry(entry, &mut count, &mut result);
pub fn insert_prefill(&mut self, task_id: TaskId, priority: Priority, max_prefill: usize) {
if let Some(prefill) = &mut self.prefill {
assert_eq!(prefill.0, priority);
for task_id in &result {
prefill.1.insert(*task_id);
}
prefill.1.insert(task_id);
} else {
self.prefill = Some((priority, result.iter().copied().collect()))
let mut v = Set::with_capacity(max_prefill);
v.insert(task_id);
self.prefill = Some((priority, v))
}
result
}

pub fn take_tasks(&mut self, mut count: u32) -> Vec<TaskId> {
Expand Down
35 changes: 28 additions & 7 deletions crates/tako/src/internal/server/reactor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -320,13 +320,20 @@ fn task_running(
// By removing redirections first, we unassign the task so we can later assign it back
// In theory, we could optimize this special case by doing nothing, but it should be quite rare
// So I prefer to keep the code simple.
try_remove_redirection(
if !try_remove_redirection(
worker_map,
scheduler_state,
request_map,
task_id,
task.resource_rq_id,
);
) {
// We have tried to retract the task, but it has already started.
// Without a redirection, the task is still in the ready queue and it has
// to be removed, otherwise the queue keeps an id of an already removed task.
task_queues
.get_mut(task.resource_rq_id)
.remove_if_queued(task.id, task.priority());
}
let rqv = request_map.get(task.resource_rq_id);
worker_map
.get_worker_mut(worker_id)
Expand Down Expand Up @@ -546,13 +553,17 @@ fn task_finished(
}
TaskRuntimeState::Retracting { worker_id: w_id } => {
assert_eq!(*w_id, worker_id);
try_remove_redirection(
if !try_remove_redirection(
worker_map,
scheduler_state,
request_map,
task_id,
task.resource_rq_id,
);
) {
task_queues
.get_mut(task.resource_rq_id)
.remove_if_queued(task.id, task.priority());
}
}
TaskRuntimeState::Prefilled { .. }
| TaskRuntimeState::Waiting { .. }
Expand Down Expand Up @@ -589,17 +600,23 @@ fn task_finished(
true
}

/// Returns true if a redirection was found (and removed).
/// A task with a redirection was already taken out of the ready queue by the scheduler,
/// a retracting task without a redirection is still waiting in the ready queue.
fn try_remove_redirection(
worker_map: &mut WorkerMap,
scheduler_state: &mut SchedulerState,
request_map: &ResourceRqMap,
task_id: TaskId,
resource_rq_id: ResourceRqId,
) {
) -> bool {
if let Some((worker_id, rv_id)) = scheduler_state.redirects.remove(&task_id) {
let worker = worker_map.get_worker_mut(worker_id);
let rq = request_map.get(resource_rq_id).get(rv_id);
worker.remove_sn_task(task_id, rq);
true
} else {
false
}
}

Expand Down Expand Up @@ -653,13 +670,17 @@ fn task_failed(
}
TaskRuntimeState::Retracting { worker_id: w } => {
assert_eq!(worker_id, *w);
try_remove_redirection(
if !try_remove_redirection(
worker_map,
scheduler_state,
request_map,
task_id,
task.resource_rq_id,
);
) {
task_queues
.get_mut(task.resource_rq_id)
.remove_if_queued(task.id, task.priority());
}
}
_ => {}
}
Expand Down
14 changes: 14 additions & 0 deletions crates/tako/src/internal/server/task.rs
Original file line number Diff line number Diff line change
Expand Up @@ -247,6 +247,7 @@ impl Task {
}
}

#[cfg(test)]
pub(crate) fn rv_id(&self) -> Option<ResourceVariantId> {
match self.state {
TaskRuntimeState::Running { rv_id, .. } | TaskRuntimeState::Assigned { rv_id, .. } => {
Expand All @@ -256,6 +257,19 @@ impl Task {
}
}

#[inline]
pub(crate) fn assigned_placement(
&self,
redirects: &Map<TaskId, (WorkerId, ResourceVariantId)>,
) -> Option<(WorkerId, ResourceVariantId)> {
match self.state {
TaskRuntimeState::Assigned { worker_id, rv_id }
| TaskRuntimeState::Running { worker_id, rv_id } => Some((worker_id, rv_id)),
TaskRuntimeState::Retracting { .. } => redirects.get(&self.id).copied(),
_ => None,
}
}

pub(crate) fn increment_instance_id(&mut self) {
self.instance_id = InstanceId::new(self.instance_id.as_num() + 1);
}
Expand Down
12 changes: 3 additions & 9 deletions crates/tako/src/internal/server/worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -243,15 +243,9 @@ impl Worker {
let mut resources = self.resources.clone();
for task_id in a.assigned_tasks.iter() {
let task = task_map.get_task(*task_id);
let (worker_id, rv_id) = match &task.state {
TaskRuntimeState::Assigned { worker_id, rv_id }
| TaskRuntimeState::Running { worker_id, rv_id } => (*worker_id, *rv_id),
TaskRuntimeState::Retracting { .. } => {
let (worker_id, rv_id) = transfers.get(task_id).unwrap();
(*worker_id, *rv_id)
}
s => panic!("Invalid state {s:?}"),
};
let (worker_id, rv_id) = task
.assigned_placement(transfers)
.unwrap_or_else(|| panic!("Invalid state {:?}", task.state));
assert_eq!(self.id, worker_id);
let rq = request_map.get(task.resource_rq_id).get(rv_id);
assert!(resources.is_capable_to_run_request(rq));
Expand Down
Loading
Loading