From f8e6fc1f10096dfe2a8c3c2a3e00fc21ecf871ba Mon Sep 17 00:00:00 2001 From: Henry Guo Date: Thu, 28 May 2026 01:54:05 +0800 Subject: [PATCH] fix(ecstore): offload erasure encoding from async workers (#3099) * fix(ecstore): offload erasure encoding from async workers * fix(ecstore): reuse encode buffers in blocking tasks --------- Co-authored-by: Henry Guo Co-authored-by: houseme --- crates/ecstore/src/erasure_coding/encode.rs | 11 ++++++++++- crates/ecstore/src/erasure_coding/erasure.rs | 16 +++++++++++++++- 2 files changed, 25 insertions(+), 2 deletions(-) diff --git a/crates/ecstore/src/erasure_coding/encode.rs b/crates/ecstore/src/erasure_coding/encode.rs index aff98c934..d4781af0b 100644 --- a/crates/ecstore/src/erasure_coding/encode.rs +++ b/crates/ecstore/src/erasure_coding/encode.rs @@ -226,7 +226,16 @@ impl Erasure { match rustfs_utils::read_full(&mut reader, &mut buf).await { Ok(n) if n > 0 => { total += n; - let res = self.encode_data(&buf[..n])?; + let erasure = self.clone(); + let encode_buf = std::mem::take(&mut buf); + let (res, returned_buf) = tokio::task::spawn_blocking(move || { + let res = erasure.encode_data(&encode_buf[..n]); + (res, encode_buf) + }) + .await + .map_err(|err| std::io::Error::other(format!("EC encode task failed: {err}")))?; + buf = returned_buf; + let res = res?; let queued_bytes = queued_block_bytes(&res); rustfs_io_metrics::add_ec_encode_inflight_bytes(queued_bytes); if let Err(err) = tx.send(res).await { diff --git a/crates/ecstore/src/erasure_coding/erasure.rs b/crates/ecstore/src/erasure_coding/erasure.rs index 58e6229cf..2a07fb686 100644 --- a/crates/ecstore/src/erasure_coding/erasure.rs +++ b/crates/ecstore/src/erasure_coding/erasure.rs @@ -587,7 +587,21 @@ impl Erasure { Ok(n) if n > 0 => { warn!("encode_stream_callback_async read n={}", n); total += n; - let res = self.encode_data(&buf[..n]); + let erasure = self.clone(); + let encode_buf = std::mem::take(&mut buf); + let (res, returned_buf) = match tokio::task::spawn_blocking(move || { + let res = erasure.encode_data(&encode_buf[..n]); + (res, encode_buf) + }) + .await + { + Ok(result) => result, + Err(err) => { + on_block(Err(std::io::Error::other(format!("EC encode task failed: {err}")))).await?; + break; + } + }; + buf = returned_buf; on_block(res).await? } Ok(_) => {