From 5bc3144c22e5fe545b512a7630381663a70d78bb Mon Sep 17 00:00:00 2001 From: Andrey Vasnetsov Date: Wed, 11 Feb 2026 14:38:19 +0100 Subject: [PATCH] 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 --- .../src/update_workers/update_worker.rs | 68 ++++++++++++------- 1 file changed, 42 insertions(+), 26 deletions(-) diff --git a/lib/collection/src/update_workers/update_worker.rs b/lib/collection/src/update_workers/update_worker.rs index 0285450780..bbb48cadc5 100644 --- a/lib/collection/src/update_workers/update_worker.rs +++ b/lib/collection/src/update_workers/update_worker.rs @@ -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>>, + result: CollectionResult, + 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)