Skip to content
Merged
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
113 changes: 64 additions & 49 deletions quickwit/quickwit-common/src/thread_pool/with_priority.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,8 @@ struct ThreadPoolInner {
struct State {
/// High-priority tasks are served first, in FIFO order.
high_priority_tasks: VecDeque<Box<dyn RunnableTask>>,
/// Normal-priority tasks are served by priority, then in FIFO order.
/// Normal-priority tasks are served by priority, then by lowest estimated query cost,
/// then in FIFO order.
normal_priority_tasks: BinaryHeap<NormalQueueEntry>,
/// Monotonically increasing counter stamped on each task entering
/// `normal_priority_tasks`; used as a tie-breaker.
Expand All @@ -62,10 +63,12 @@ impl State {

/// A normal-priority task waiting to be scheduled on the thread pool.
///
/// The ordering is defined over `(priority, sequence)`.
/// The ordering is defined over `(priority, job_cost, sequence)`.
struct NormalQueueEntry {
/// Lower values have higher priority.
priority: i32,
/// Within a priority, cheaper tasks are scheduled first.
job_cost: usize,
sequence: u64,
task: Box<dyn RunnableTask>,
}
Expand All @@ -75,6 +78,7 @@ impl Ord for NormalQueueEntry {
other
.priority
.cmp(&self.priority)
.then_with(|| other.job_cost.cmp(&self.job_cost))
.then_with(|| other.sequence.cmp(&self.sequence))
}
}
Expand Down Expand Up @@ -159,17 +163,28 @@ impl Drop for Permit {
/// The priority of a task submitted to a [`ThreadPoolWithPriority`].
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum Priority {
/// The default priority, carrying the priority of the originating request.
Normal(i32),
/// The default priority, carrying the priority and the estimated cost of the originating
/// request.
Normal {
/// Lower values have higher priority.
priority: i32,
/// Within a priority, cheaper tasks are scheduled first.
job_cost: usize,
},
/// A high-priority task is scheduled before normal-priority tasks that are still pending.
/// This is reserved for short tasks sitting on the critical path of a request (e.g. merging
/// partial results).
High,
}

/// Default tasks have zero priority and cost so short operations, such as
/// list-fields result merging, run before split searches.
impl Default for Priority {
fn default() -> Self {
Priority::Normal(0)
Priority::Normal {
priority: 0,
job_cost: 0,
}
}
}

Expand Down Expand Up @@ -253,11 +268,12 @@ impl ThreadPoolInner {
fn enqueue(&self, priority: Priority, task: Box<dyn RunnableTask>) {
let mut state = self.state.lock().unwrap();
match priority {
Priority::Normal(priority) => {
Priority::Normal { priority, job_cost } => {
let sequence = state.next_normal_sequence;
state.next_normal_sequence += 1;
state.normal_priority_tasks.push(NormalQueueEntry {
priority,
job_cost,
sequence,
task,
});
Expand Down Expand Up @@ -292,7 +308,6 @@ impl ThreadPoolInner {

#[cfg(test)]
mod tests {
use std::collections::HashMap;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Duration;
Expand Down Expand Up @@ -353,13 +368,13 @@ mod tests {

let execution_order_clone = execution_order.clone();
let normal_task_1 =
thread_pool.run_cpu_intensive_with_priority(Priority::Normal(0), move || {
thread_pool.run_cpu_intensive_with_priority(Priority::default(), move || {
execution_order_clone.lock().unwrap().push(1);
});

let execution_order_clone = execution_order.clone();
let normal_task_2 =
thread_pool.run_cpu_intensive_with_priority(Priority::Normal(0), move || {
thread_pool.run_cpu_intensive_with_priority(Priority::default(), move || {
execution_order_clone.lock().unwrap().push(2);
});

Expand All @@ -380,13 +395,13 @@ mod tests {

/// Saturates a single-threaded pool, then enqueues `priorities` while
/// the pool is busy, so they all pile up in the pending queue. Returns the order in which
/// they were eventually executed, as `(priority, index_within_that_priority)`.
/// they were eventually executed, as indices into `priorities`.
async fn execution_order_for_pending_priorities(
name: &'static str,
priorities: &[i32],
) -> Vec<(i32, usize)> {
priorities: &[Priority],
) -> Vec<usize> {
let thread_pool = ThreadPoolWithPriority::new(name, Some(1));
let execution_order: Arc<std::sync::Mutex<Vec<(i32, usize)>>> = Default::default();
let execution_order: Arc<std::sync::Mutex<Vec<usize>>> = Default::default();
let (started_tx, started_rx) = oneshot::channel();
let (release_tx, release_rx) = std::sync::mpsc::channel();

Expand All @@ -397,19 +412,14 @@ mod tests {
});
started_rx.await.unwrap();

let mut seen_per_priority: HashMap<i32, usize> = HashMap::new();
let mut pending_tasks = Vec::new();
for &priority in priorities {
let index_within_priority = seen_per_priority.entry(priority).or_insert(0);
let expected = (priority, *index_within_priority);
*index_within_priority += 1;
for (index, &priority) in priorities.iter().enumerate() {
let execution_order_clone = execution_order.clone();
pending_tasks.push(thread_pool.run_cpu_intensive_with_priority(
Priority::Normal(priority),
move || {
execution_order_clone.lock().unwrap().push(expected);
},
));
pending_tasks.push(
thread_pool.run_cpu_intensive_with_priority(priority, move || {
execution_order_clone.lock().unwrap().push(index);
}),
);
}

release_tx.send(()).unwrap();
Expand All @@ -422,35 +432,40 @@ mod tests {

#[tokio::test]
async fn test_normal_priority_tasks_are_served_by_priority() {
let priorities = [5, -3, 0, 10, -8].map(|priority| Priority::Normal {
priority,
job_cost: 0,
});
let execution_order =
execution_order_for_pending_priorities("priority_heap_order_test", &[5, -3, 0, 10, -8])
.await;
assert_eq!(
execution_order,
vec![(-8, 0), (-3, 0), (0, 0), (5, 0), (10, 0)]
);
execution_order_for_pending_priorities("priority_heap_order_test", &priorities).await;
assert_eq!(execution_order, vec![4, 1, 2, 0, 3]);
}

#[tokio::test]
async fn test_normal_priority_interleaves_priority_then_fifo() {
let execution_order = execution_order_for_pending_priorities(
"priority_heap_fifo_test",
&[1, 0, 1, 0, 1, -1, 0, -1],
)
.await;
assert_eq!(
execution_order,
vec![
(-1, 0),
(-1, 1),
(0, 0),
(0, 1),
(0, 2),
(1, 0),
(1, 1),
(1, 2),
]
);
async fn test_normal_tasks_order_by_priority_cost_then_fifo() {
let priorities = [
Priority::Normal {
priority: 0,
job_cost: 100,
},
Priority::Normal {
priority: 0,
job_cost: 10,
},
Priority::Normal {
priority: -1,
job_cost: 1000,
},
Priority::Normal {
priority: 0,
job_cost: 10,
},
Priority::default(),
];
let execution_order =
execution_order_for_pending_priorities("priority_heap_fifo_test", &priorities).await;

assert_eq!(execution_order, vec![2, 4, 1, 3, 0]);
}

#[tokio::test]
Expand Down
23 changes: 12 additions & 11 deletions quickwit/quickwit-janitor/src/actors/delete_task_planner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,10 @@ use quickwit_proto::metastore::{
};
use quickwit_proto::search::SearchRequest;
use quickwit_proto::types::IndexUid;
use quickwit_search::{IndexMetasForLeafSearch, SearchJob, SearchJobPlacer, jobs_to_leaf_request};
use quickwit_search::{
IndexMetasForLeafSearch, SearchJob, SearchJobPlacer, compute_query_complexity_factor,
jobs_to_leaf_request,
};
use serde::Serialize;
use tantivy::Inventory;
use tracing::{debug, info};
Expand Down Expand Up @@ -298,11 +301,6 @@ impl DeleteTaskPlanner {
index_uri: &str,
ctx: &ActorContext<Self>,
) -> anyhow::Result<bool> {
let search_job = SearchJob::from(&stale_split.split_metadata);
let mut search_client = self
.search_job_placer
.assign_job(search_job.clone(), &HashSet::new())
.await?;
for delete_task in delete_tasks {
let delete_query = delete_task
.delete_query
Expand All @@ -316,6 +314,12 @@ impl DeleteTaskPlanner {
end_timestamp: delete_query.end_timestamp,
..Default::default()
};
let query_complexity_factor = compute_query_complexity_factor(&search_request)?;
let search_job = SearchJob::new(&stale_split.split_metadata, query_complexity_factor);
let mut search_client = self
.search_job_placer
.assign_job(search_job.clone(), &HashSet::new())
.await?;
let mut search_indexes_metas = HashMap::new();
let index_uri = Uri::from_str(index_uri).context("invalid index URI")?;
search_indexes_metas.insert(
Expand All @@ -325,11 +329,8 @@ impl DeleteTaskPlanner {
index_uri,
},
);
let leaf_search_request = jobs_to_leaf_request(
&search_request,
&search_indexes_metas,
vec![search_job.clone()],
)?;
let leaf_search_request =
jobs_to_leaf_request(&search_request, &search_indexes_metas, vec![search_job])?;
let response = search_client.leaf_search(leaf_search_request).await?;
ctx.record_progress();
if response.num_hits > 0 {
Expand Down
Loading
Loading