mirror of
https://github.com/qdrant/qdrant.git
synced 2026-10-02 19:07:48 -05:00
non blocking reading from wal (#8101)
* move wal.lock.read into a spawn_blocking * [AI] In the file with update_worker bloast the code, please try to move it out of the main function, so the function itself stay clean
This commit is contained in:
@@ -7,7 +7,7 @@ use segment::types::SeqNumberType;
|
||||
use shard::operations::CollectionUpdateOperations;
|
||||
use shard::segment_holder::locked::LockedSegmentHolder;
|
||||
use tokio::sync::mpsc::{Receiver, Sender};
|
||||
use tokio::sync::{Mutex as TokioMutex, watch};
|
||||
use tokio::sync::{Mutex as TokioMutex, oneshot, watch};
|
||||
|
||||
use crate::collection_manager::collection_updater::CollectionUpdater;
|
||||
use crate::common::stoppable_task::StoppableTaskHandle;
|
||||
@@ -24,6 +24,20 @@ use crate::wal_delta::LockedWal;
|
||||
|
||||
const BYTES_IN_KB: usize = 1024;
|
||||
|
||||
/// Sends the operation result through the feedback channel if present.
|
||||
/// Logs a debug message if the receiver is no longer waiting.
|
||||
fn send_feedback(
|
||||
sender: Option<oneshot::Sender<CollectionResult<usize>>>,
|
||||
result: CollectionResult<usize>,
|
||||
op_num: SeqNumberType,
|
||||
) {
|
||||
if let Some(feedback) = sender {
|
||||
feedback.send(result).unwrap_or_else(|_| {
|
||||
log::debug!("Can't report operation {op_num} result. Assume already not required");
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
impl UpdateWorkers {
|
||||
/// Main loop of the update worker.
|
||||
///
|
||||
@@ -71,25 +85,35 @@ impl UpdateWorkers {
|
||||
let operation = if let Some(operation) = operation {
|
||||
*operation
|
||||
} else {
|
||||
let record = wal.lock().await.read_single_record(op_num);
|
||||
let wal_clone = wal.clone();
|
||||
let record = match tokio::task::spawn_blocking(move || {
|
||||
wal_clone.blocking_lock().read_single_record(op_num)
|
||||
})
|
||||
.await
|
||||
{
|
||||
Ok(record) => record,
|
||||
Err(err) => {
|
||||
log::error!("Can't read operation {op_num} from WAL - {err}");
|
||||
send_feedback(sender, Err(CollectionError::from(err)), op_num);
|
||||
continue;
|
||||
}
|
||||
};
|
||||
|
||||
match record {
|
||||
Ok(Some(op)) => op.operation,
|
||||
Ok(None) => {
|
||||
if let Some(feedback) = sender {
|
||||
feedback.send(Err(CollectionError::service_error(
|
||||
format!("Operation {op_num} not found in WAL"),
|
||||
))).unwrap_or_else(|_| {
|
||||
log::debug!("Can't report operation {op_num} result. Assume already not required");
|
||||
});
|
||||
}
|
||||
send_feedback(
|
||||
sender,
|
||||
Err(CollectionError::service_error(format!(
|
||||
"Operation {op_num} not found in WAL"
|
||||
))),
|
||||
op_num,
|
||||
);
|
||||
continue;
|
||||
}
|
||||
Err(err) => {
|
||||
if let Some(feedback) = sender {
|
||||
feedback.send(Err(CollectionError::from(err))).unwrap_or_else(|_| {
|
||||
log::debug!("Can't report operation {op_num} result. Assume already not required");
|
||||
});
|
||||
}
|
||||
log::error!("Can't read operation {op_num} from WAL - {err}");
|
||||
send_feedback(sender, Err(CollectionError::from(err)), op_num);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
@@ -103,14 +127,10 @@ impl UpdateWorkers {
|
||||
)
|
||||
.await;
|
||||
|
||||
if let Err(err) = operation_result
|
||||
&& let Some(feedback) = sender
|
||||
{
|
||||
feedback.send(Err(err)).unwrap_or_else(|_| {
|
||||
log::debug!("Can't report operation {op_num} result. Assume already not required");
|
||||
});
|
||||
if let Err(err) = operation_result {
|
||||
send_feedback(sender, Err(err), op_num);
|
||||
continue;
|
||||
};
|
||||
}
|
||||
|
||||
let wait = sender.is_some();
|
||||
let operation_result = tokio::task::spawn_blocking(move || {
|
||||
@@ -142,11 +162,7 @@ impl UpdateWorkers {
|
||||
log::error!("Can't update last applied_seq {err}")
|
||||
}
|
||||
|
||||
if let Some(feedback) = sender {
|
||||
feedback.send(res).unwrap_or_else(|_| {
|
||||
log::debug!("Can't report operation {op_num} result. Assume already not required");
|
||||
});
|
||||
};
|
||||
send_feedback(sender, res, op_num);
|
||||
}
|
||||
UpdateSignal::Nop => optimize_sender
|
||||
.send(OptimizerSignal::Nop)
|
||||
|
||||
Reference in New Issue
Block a user