From 08d0021cd6c1806a1f5fae1f9e48bdd78d0062c9 Mon Sep 17 00:00:00 2001 From: junxiang Mu <1948535941@qq.com> Date: Thu, 24 Apr 2025 03:04:05 +0000 Subject: [PATCH 1/2] support async calculate etag Signed-off-by: junxiang Mu <1948535941@qq.com> --- ecstore/src/io.rs | 50 ++++++++++++++++++++++++++++------------- ecstore/src/set_disk.rs | 4 ++-- 2 files changed, 37 insertions(+), 17 deletions(-) diff --git a/ecstore/src/io.rs b/ecstore/src/io.rs index 3ca27fe0d..7ecf8f84c 100644 --- a/ecstore/src/io.rs +++ b/ecstore/src/io.rs @@ -1,6 +1,8 @@ +use bytes::Bytes; use futures::TryStreamExt; use md5::Digest; use md5::Md5; +use tokio::sync::mpsc; use std::pin::Pin; use std::task::Context; use std::task::Poll; @@ -125,33 +127,51 @@ impl AsyncRead for HttpFileReader { pub struct EtagReader { inner: R, - md5: Md5, + bytes_tx: mpsc::Sender, + md5_rx: oneshot::Receiver, } impl EtagReader { pub fn new(inner: R) -> Self { - EtagReader { inner, md5: Md5::new() } + let (bytes_tx, mut bytes_rx) = mpsc::channel::(8); + let (md5_tx, md5_rx) = oneshot::channel::(); + + tokio::task::spawn_blocking(move || { + let mut md5 = Md5::new(); + while let Some(bytes) = bytes_rx.blocking_recv() { + md5.update(&bytes); + } + let digest = md5.finalize(); + let etag = hex_simd::encode_to_string(digest, hex_simd::AsciiCase::Lower); + let _ = md5_tx.send(etag); + }); + + EtagReader { inner, bytes_tx, md5_rx } } - pub fn etag(self) -> String { - hex_simd::encode_to_string(self.md5.finalize(), hex_simd::AsciiCase::Lower) + pub async fn etag(self) -> String { + drop(self.inner); + drop(self.bytes_tx); + self.md5_rx.await.unwrap() } } impl AsyncRead for EtagReader { + #[tracing::instrument(level = "debug", skip_all)] fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { - let befor_size = buf.filled().len(); - - match Pin::new(&mut self.inner).poll_read(cx, buf) { - Poll::Ready(Ok(())) => { - if buf.filled().len() > befor_size { - let bytes = &buf.filled()[befor_size..]; - self.md5.update(bytes); - } - - Poll::Ready(Ok(())) + let poll = Pin::new(&mut self.inner).poll_read(cx, buf); + if let Poll::Ready(Ok(())) = &poll { + if buf.remaining() == 0 { + let bytes = buf.filled(); + let bytes = Bytes::copy_from_slice(bytes); + let tx = self.bytes_tx.clone(); + tokio::spawn(async move { + if let Err(e) = tx.send(bytes).await { + warn!("EtagReader send error: {:?}", e); + } + }); } - other => other, } + poll } } diff --git a/ecstore/src/set_disk.rs b/ecstore/src/set_disk.rs index c64bc189b..b35c6ad89 100644 --- a/ecstore/src/set_disk.rs +++ b/ecstore/src/set_disk.rs @@ -3802,7 +3802,7 @@ impl ObjectIO for SetDisks { error!("close_bitrot_writers err {:?}", err); } - let etag = etag_stream.etag(); + let etag = etag_stream.etag().await; //TODO: userDefined user_defined.insert("etag".to_owned(), etag.clone()); @@ -4393,7 +4393,7 @@ impl StorageAPI for SetDisks { error!("close_bitrot_writers err {:?}", err); } - let mut etag = etag_stream.etag(); + let mut etag = etag_stream.etag().await; if let Some(ref tag) = opts.preserve_etag { etag = tag.clone(); From b2da8148bd4d19f12722eef2403c96c29d1024d9 Mon Sep 17 00:00:00 2001 From: junxiang Mu <1948535941@qq.com> Date: Thu, 24 Apr 2025 08:12:57 +0000 Subject: [PATCH 2/2] support concurrency write Signed-off-by: junxiang Mu <1948535941@qq.com> --- ecstore/src/bitrot.rs | 19 +++++++++++++------ ecstore/src/erasure.rs | 31 ++++++++++++++++--------------- ecstore/src/io.rs | 2 +- ecstore/src/set_disk.rs | 39 ++++++++++++++++++++++----------------- 4 files changed, 52 insertions(+), 39 deletions(-) diff --git a/ecstore/src/bitrot.rs b/ecstore/src/bitrot.rs index 757401557..68ccff71e 100644 --- a/ecstore/src/bitrot.rs +++ b/ecstore/src/bitrot.rs @@ -6,6 +6,7 @@ use crate::{ }; use blake2::Blake2b512; use blake2::Digest as _; +use bytes::Bytes; use common::error::{Error, Result}; use highway::{HighwayHash, HighwayHasher, Key}; use lazy_static::lazy_static; @@ -533,20 +534,26 @@ impl Writer for BitrotFileWriter { self } - async fn write(&mut self, buf: &[u8]) -> Result<()> { + async fn write(&mut self, buf: Bytes) -> Result<()> { if buf.is_empty() { return Ok(()); } - self.hasher.reset(); - self.hasher.update(buf); - let hash_bytes = self.hasher.clone().finalize(); + let mut hasher = self.hasher.clone(); + let h_buf = buf.clone(); + let hash_bytes = tokio::spawn(async move { + hasher.reset(); + hasher.update(h_buf); + hasher.finalize() + }) + .await + .unwrap(); if let Some(f) = self.inner.as_mut() { f.write_all(&hash_bytes).await?; - f.write_all(buf).await?; + f.write_all(&buf).await?; } else { self.inline_data.extend_from_slice(&hash_bytes); - self.inline_data.extend_from_slice(buf); + self.inline_data.extend_from_slice(&buf); } Ok(()) diff --git a/ecstore/src/erasure.rs b/ecstore/src/erasure.rs index d99c7eea5..b4c45ff11 100644 --- a/ecstore/src/erasure.rs +++ b/ecstore/src/erasure.rs @@ -86,7 +86,6 @@ impl Erasure { } self.buf.resize(new_len, 0u8); - match reader.read_exact(&mut self.buf).await { Ok(res) => res, Err(e) => { @@ -101,20 +100,22 @@ impl Erasure { } self.encode_data(&self.buf, &mut blocks)?; - let mut errs = Vec::new(); - // TODO: 并发写入 - for (i, w_op) in writers.iter_mut().enumerate() { - if let Some(w) = w_op { - match w.write(blocks[i].as_ref()).await { - Ok(_) => errs.push(None), - Err(e) => errs.push(Some(e)), + let write_futures = writers.iter_mut().enumerate().map(|(i, w_op)| { + let i_inner = i.clone(); + let blocks_inner = blocks.clone(); + async move { + if let Some(w) = w_op { + match w.write(blocks_inner[i_inner].clone()).await { + Ok(_) => None, + Err(e) => Some(e), + } + } else { + Some(Error::new(DiskError::DiskNotFound)) } - } else { - errs.push(Some(Error::new(DiskError::DiskNotFound))); } - } - + }); + let errs = join_all(write_futures).await; let none_count = errs.iter().filter(|&x| x.is_none()).count(); if none_count >= write_quorum { if total_size == 0 { @@ -472,7 +473,7 @@ impl Erasure { self.encoder.as_ref().unwrap().reconstruct(&mut bufs)?; } - let shards: Vec> = bufs.into_iter().flatten().collect::>(); + let shards = bufs.into_iter().flatten().collect::>(); if shards.len() != self.parity_shards + self.data_shards { return Err(Error::from_string("can not reconstruct data")); } @@ -481,7 +482,7 @@ impl Erasure { if w.is_none() { continue; } - match w.as_mut().unwrap().write(shards[i].as_ref()).await { + match w.as_mut().unwrap().write(shards[i].clone().into()).await { Ok(_) => {} Err(e) => { info!("write failed, err: {:?}", e); @@ -501,7 +502,7 @@ impl Erasure { #[async_trait::async_trait] pub trait Writer { fn as_any(&self) -> &dyn Any; - async fn write(&mut self, buf: &[u8]) -> Result<()>; + async fn write(&mut self, buf: Bytes) -> Result<()>; async fn close(&mut self) -> Result<()> { Ok(()) } diff --git a/ecstore/src/io.rs b/ecstore/src/io.rs index 7ecf8f84c..f2affe8cb 100644 --- a/ecstore/src/io.rs +++ b/ecstore/src/io.rs @@ -2,13 +2,13 @@ use bytes::Bytes; use futures::TryStreamExt; use md5::Digest; use md5::Md5; -use tokio::sync::mpsc; use std::pin::Pin; use std::task::Context; use std::task::Poll; use tokio::io::AsyncRead; use tokio::io::AsyncWrite; use tokio::io::ReadBuf; +use tokio::sync::mpsc; use tokio::sync::oneshot; use tokio_util::io::ReaderStream; use tokio_util::io::StreamReader; diff --git a/ecstore/src/set_disk.rs b/ecstore/src/set_disk.rs index b35c6ad89..77eafe01a 100644 --- a/ecstore/src/set_disk.rs +++ b/ecstore/src/set_disk.rs @@ -325,6 +325,7 @@ impl SetDisks { } } + let mut futures = Vec::with_capacity(disks.len()); if let Some(err) = reduce_write_quorum_errs(&errs, object_op_ignored_errs().as_ref(), write_quorum) { // TODO: 并发 for (i, err) in errs.iter().enumerate() { @@ -334,26 +335,30 @@ impl SetDisks { if let Some(disk) = disks[i].as_ref() { let fi = file_infos[i].clone(); - let _ = disk - .delete_version( - src_bucket, - src_object, - fi, - false, - DeleteOptions { - undo_write: true, - old_data_dir: data_dirs[i], - ..Default::default() - }, - ) - .await - .map_err(|e| { - debug!("rename_data delete_version err {:?}", e); - e - }); + let old_data_dir = data_dirs[i]; + futures.push(async move { + let _ = disk + .delete_version( + src_bucket, + src_object, + fi, + false, + DeleteOptions { + undo_write: true, + old_data_dir: old_data_dir, + ..Default::default() + }, + ) + .await + .map_err(|e| { + debug!("rename_data delete_version err {:?}", e); + e + }); + }); } } + let _ = join_all(futures).await; return Err(err); }