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); }