From 8fc637fb14d33f1907f29b2c8c85781851cf0d82 Mon Sep 17 00:00:00 2001 From: houseme Date: Wed, 8 Jul 2026 21:38:00 +0800 Subject: [PATCH] perf(ecstore): run the short EC encode inline instead of block_in_place (#4484) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Each erasure block runs its Reed-Solomon encode through tokio::task::block_in_place on the multi-threaded runtime. That parks the worker and asks the scheduler to relocate other tasks, but the encode itself is only ~110µs per 1MiB block (p99 ~542µs) — the profiling in backlog#932 flagged the scheduling disturbance as comparable to the compute it guards. Call the encode closure inline on the multi-threaded runtime instead. The CurrentThread (and any other) flavor keeps spawn_blocking so the sole executor thread is never blocked and block_in_place's multi-thread-only requirement is respected. Applied to both ingest paths (encode_block / Vec and encode_block_bytes_mut / BytesMut); no change to encode output, quorum, shutdown, or error handling. Adds encode_works_on_multi_thread_runtime to cover the previously-untested multi-threaded arm for both ingest paths, asserting streaming and batched encode produce identical shard bytes. This is the low-risk, correctness-neutral item that backlog#932's adversarial verification recommended splitting out and doing first; the larger per-writer pipeline restructure it belongs to stays gated on a Linux multi-disk baseline. Refs: rustfs/backlog#932 (HP-11), rustfs/backlog#936 Co-authored-by: heihutu --- crates/ecstore/src/erasure/coding/encode.rs | 67 ++++++++++++++++++++- 1 file changed, 65 insertions(+), 2 deletions(-) diff --git a/crates/ecstore/src/erasure/coding/encode.rs b/crates/ecstore/src/erasure/coding/encode.rs index b0ce598c0..76d04dc99 100644 --- a/crates/ecstore/src/erasure/coding/encode.rs +++ b/crates/ecstore/src/erasure/coding/encode.rs @@ -336,7 +336,14 @@ impl Erasure { }; let (res, returned_buf) = match tokio::runtime::Handle::current().runtime_flavor() { - RuntimeFlavor::MultiThread => tokio::task::block_in_place(encode_once), + // EC encode is a short CPU burst (~110µs per 1MiB block, p99 ~542µs). + // On the multi-threaded runtime block_in_place parked the worker and + // churned the scheduler for a cost comparable to the encode itself + // (rustfs/backlog#932); at this duration a direct inline call is + // cheaper and equally correct. CurrentThread (and any other flavor) + // keep spawn_blocking so the sole executor thread is never blocked + // and block_in_place's multi-thread-only requirement is respected. + RuntimeFlavor::MultiThread => encode_once(), RuntimeFlavor::CurrentThread => tokio::task::spawn_blocking(encode_once) .await .map_err(|err| std::io::Error::other(format!("EC encode task failed: {err}")))?, @@ -354,7 +361,10 @@ impl Erasure { let encode_once = move || self.encode_data_bytes_mut(encode_buf, len); let res = match tokio::runtime::Handle::current().runtime_flavor() { - RuntimeFlavor::MultiThread => tokio::task::block_in_place(encode_once), + // Same rationale as encode_block: inline the short EC burst on the + // multi-threaded runtime instead of parking a worker via + // block_in_place; keep spawn_blocking on single-threaded flavors. + RuntimeFlavor::MultiThread => encode_once(), RuntimeFlavor::CurrentThread => tokio::task::spawn_blocking(encode_once) .await .map_err(|err| std::io::Error::other(format!("EC encode task failed: {err}")))?, @@ -1025,6 +1035,59 @@ mod tests { assert_eq!(written, b"current-thread payload".len()); } + // Covers the RuntimeFlavor::MultiThread arm of encode_block / + // encode_block_bytes_mut, which runs the EC compute inline instead of via + // block_in_place (rustfs/backlog#932). The other tests default to the + // current-thread runtime and only exercise the spawn_blocking arm. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn encode_works_on_multi_thread_runtime() { + const DATA_SHARDS: usize = 2; + const PARITY_SHARDS: usize = 2; + const TOTAL_SHARDS: usize = DATA_SHARDS + PARITY_SHARDS; + const BLOCK_SIZE: usize = 32; + + let payload = vec![9u8; BLOCK_SIZE * 3 + 7]; + + // Streaming encode() drives encode_block (Vec ingest). + let streamed: Vec>>> = (0..TOTAL_SHARDS).map(|_| Arc::new(Mutex::new(Vec::new()))).collect(); + let mut streamed_writers: Vec> = streamed + .iter() + .map(|c| Some(bitrot_writer(DeferredCommitWriter::new(c.clone()), BLOCK_SIZE / DATA_SHARDS))) + .collect(); + let erasure = Arc::new(Erasure::new(DATA_SHARDS, PARITY_SHARDS, BLOCK_SIZE)); + let (_r, streamed_total) = erasure + .clone() + .encode( + tokio::io::BufReader::new(Cursor::new(payload.clone())), + &mut streamed_writers, + DATA_SHARDS, + ) + .await + .expect("streaming encode should succeed on the multi-threaded runtime"); + assert_eq!(streamed_total, payload.len()); + + // Batched encode() drives encode_block_bytes_mut (BytesMut ingest). + let batched: Vec>>> = (0..TOTAL_SHARDS).map(|_| Arc::new(Mutex::new(Vec::new()))).collect(); + let mut batched_writers: Vec> = batched + .iter() + .map(|c| Some(bitrot_writer(DeferredCommitWriter::new(c.clone()), BLOCK_SIZE / DATA_SHARDS))) + .collect(); + let (_r2, batched_total) = erasure + .encode_batched(tokio::io::BufReader::new(Cursor::new(payload.clone())), &mut batched_writers, DATA_SHARDS) + .await + .expect("batched encode should succeed on the multi-threaded runtime"); + assert_eq!(batched_total, payload.len()); + + // Both ingest paths must produce identical shard bytes regardless of the + // runtime flavor / dispatch strategy. + for (index, (a, b)) in streamed.iter().zip(batched.iter()).enumerate() { + let a = a.lock().expect("streamed shard lockable").clone(); + let b = b.lock().expect("batched shard lockable").clone(); + assert!(!a.is_empty(), "shard {index} should receive committed data"); + assert_eq!(a, b, "shard {index} must match between streaming and batched encode"); + } + } + #[tokio::test] async fn encode_batched_writes_full_and_tail_batches() { const DATA_SHARDS: usize = 2;