support concurrency write

Signed-off-by: junxiang Mu <1948535941@qq.com>
This commit is contained in:
junxiang Mu
2025-04-24 08:12:57 +00:00
parent 08d0021cd6
commit b2da8148bd
4 changed files with 52 additions and 39 deletions
+13 -6
View File
@@ -6,6 +6,7 @@ use crate::{
}; };
use blake2::Blake2b512; use blake2::Blake2b512;
use blake2::Digest as _; use blake2::Digest as _;
use bytes::Bytes;
use common::error::{Error, Result}; use common::error::{Error, Result};
use highway::{HighwayHash, HighwayHasher, Key}; use highway::{HighwayHash, HighwayHasher, Key};
use lazy_static::lazy_static; use lazy_static::lazy_static;
@@ -533,20 +534,26 @@ impl Writer for BitrotFileWriter {
self self
} }
async fn write(&mut self, buf: &[u8]) -> Result<()> { async fn write(&mut self, buf: Bytes) -> Result<()> {
if buf.is_empty() { if buf.is_empty() {
return Ok(()); return Ok(());
} }
self.hasher.reset(); let mut hasher = self.hasher.clone();
self.hasher.update(buf); let h_buf = buf.clone();
let hash_bytes = self.hasher.clone().finalize(); 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() { if let Some(f) = self.inner.as_mut() {
f.write_all(&hash_bytes).await?; f.write_all(&hash_bytes).await?;
f.write_all(buf).await?; f.write_all(&buf).await?;
} else { } else {
self.inline_data.extend_from_slice(&hash_bytes); self.inline_data.extend_from_slice(&hash_bytes);
self.inline_data.extend_from_slice(buf); self.inline_data.extend_from_slice(&buf);
} }
Ok(()) Ok(())
+16 -15
View File
@@ -86,7 +86,6 @@ impl Erasure {
} }
self.buf.resize(new_len, 0u8); self.buf.resize(new_len, 0u8);
match reader.read_exact(&mut self.buf).await { match reader.read_exact(&mut self.buf).await {
Ok(res) => res, Ok(res) => res,
Err(e) => { Err(e) => {
@@ -101,20 +100,22 @@ impl Erasure {
} }
self.encode_data(&self.buf, &mut blocks)?; self.encode_data(&self.buf, &mut blocks)?;
let mut errs = Vec::new();
// TODO: 并发写入 let write_futures = writers.iter_mut().enumerate().map(|(i, w_op)| {
for (i, w_op) in writers.iter_mut().enumerate() { let i_inner = i.clone();
if let Some(w) = w_op { let blocks_inner = blocks.clone();
match w.write(blocks[i].as_ref()).await { async move {
Ok(_) => errs.push(None), if let Some(w) = w_op {
Err(e) => errs.push(Some(e)), 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(); let none_count = errs.iter().filter(|&x| x.is_none()).count();
if none_count >= write_quorum { if none_count >= write_quorum {
if total_size == 0 { if total_size == 0 {
@@ -472,7 +473,7 @@ impl Erasure {
self.encoder.as_ref().unwrap().reconstruct(&mut bufs)?; self.encoder.as_ref().unwrap().reconstruct(&mut bufs)?;
} }
let shards: Vec<Vec<u8>> = bufs.into_iter().flatten().collect::<Vec<_>>(); let shards = bufs.into_iter().flatten().collect::<Vec<_>>();
if shards.len() != self.parity_shards + self.data_shards { if shards.len() != self.parity_shards + self.data_shards {
return Err(Error::from_string("can not reconstruct data")); return Err(Error::from_string("can not reconstruct data"));
} }
@@ -481,7 +482,7 @@ impl Erasure {
if w.is_none() { if w.is_none() {
continue; continue;
} }
match w.as_mut().unwrap().write(shards[i].as_ref()).await { match w.as_mut().unwrap().write(shards[i].clone().into()).await {
Ok(_) => {} Ok(_) => {}
Err(e) => { Err(e) => {
info!("write failed, err: {:?}", e); info!("write failed, err: {:?}", e);
@@ -501,7 +502,7 @@ impl Erasure {
#[async_trait::async_trait] #[async_trait::async_trait]
pub trait Writer { pub trait Writer {
fn as_any(&self) -> &dyn Any; 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<()> { async fn close(&mut self) -> Result<()> {
Ok(()) Ok(())
} }
+1 -1
View File
@@ -2,13 +2,13 @@ use bytes::Bytes;
use futures::TryStreamExt; use futures::TryStreamExt;
use md5::Digest; use md5::Digest;
use md5::Md5; use md5::Md5;
use tokio::sync::mpsc;
use std::pin::Pin; use std::pin::Pin;
use std::task::Context; use std::task::Context;
use std::task::Poll; use std::task::Poll;
use tokio::io::AsyncRead; use tokio::io::AsyncRead;
use tokio::io::AsyncWrite; use tokio::io::AsyncWrite;
use tokio::io::ReadBuf; use tokio::io::ReadBuf;
use tokio::sync::mpsc;
use tokio::sync::oneshot; use tokio::sync::oneshot;
use tokio_util::io::ReaderStream; use tokio_util::io::ReaderStream;
use tokio_util::io::StreamReader; use tokio_util::io::StreamReader;
+22 -17
View File
@@ -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) { if let Some(err) = reduce_write_quorum_errs(&errs, object_op_ignored_errs().as_ref(), write_quorum) {
// TODO: 并发 // TODO: 并发
for (i, err) in errs.iter().enumerate() { for (i, err) in errs.iter().enumerate() {
@@ -334,26 +335,30 @@ impl SetDisks {
if let Some(disk) = disks[i].as_ref() { if let Some(disk) = disks[i].as_ref() {
let fi = file_infos[i].clone(); let fi = file_infos[i].clone();
let _ = disk let old_data_dir = data_dirs[i];
.delete_version( futures.push(async move {
src_bucket, let _ = disk
src_object, .delete_version(
fi, src_bucket,
false, src_object,
DeleteOptions { fi,
undo_write: true, false,
old_data_dir: data_dirs[i], DeleteOptions {
..Default::default() undo_write: true,
}, old_data_dir: old_data_dir,
) ..Default::default()
.await },
.map_err(|e| { )
debug!("rename_data delete_version err {:?}", e); .await
e .map_err(|e| {
}); debug!("rename_data delete_version err {:?}", e);
e
});
});
} }
} }
let _ = join_all(futures).await;
return Err(err); return Err(err);
} }