mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-24 21:26:28 +00:00
perf(ecstore): reduce batch read identity cloning (#6441)
Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
@@ -140,6 +140,21 @@ struct CoalescedReadVersionRequest {
|
|||||||
tx: oneshot::Sender<disk::error::Result<FileInfo>>,
|
tx: oneshot::Sender<disk::error::Result<FileInfo>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Clone, Debug, Eq, PartialEq)]
|
||||||
|
struct ExpectedBatchReadVersionItem {
|
||||||
|
path: String,
|
||||||
|
version_id: String,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl From<&BatchReadVersionItem> for ExpectedBatchReadVersionItem {
|
||||||
|
fn from(item: &BatchReadVersionItem) -> Self {
|
||||||
|
Self {
|
||||||
|
path: item.path.clone(),
|
||||||
|
version_id: item.version_id.clone(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
|
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
|
||||||
struct ReadVersionCoalescerKey {
|
struct ReadVersionCoalescerKey {
|
||||||
disk: usize,
|
disk: usize,
|
||||||
@@ -266,7 +281,7 @@ async fn flush_read_version_coalescer_pending(
|
|||||||
items.push(request.item);
|
items.push(request.item);
|
||||||
}
|
}
|
||||||
|
|
||||||
let expected_items = items.clone();
|
let expected_items = items.iter().map(ExpectedBatchReadVersionItem::from).collect::<Vec<_>>();
|
||||||
record_read_version_coalescer_event("attempted_batch", items.len());
|
record_read_version_coalescer_event("attempted_batch", items.len());
|
||||||
let result =
|
let result =
|
||||||
match tokio::time::timeout(get_drive_metadata_timeout(), disk.batch_read_version(BatchReadVersionReq { items, opts }))
|
match tokio::time::timeout(get_drive_metadata_timeout(), disk.batch_read_version(BatchReadVersionReq { items, opts }))
|
||||||
@@ -292,7 +307,7 @@ async fn flush_read_version_coalescer_pending(
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn map_batch_read_version_responses(
|
fn map_batch_read_version_responses(
|
||||||
expected_items: &[BatchReadVersionItem],
|
expected_items: &[ExpectedBatchReadVersionItem],
|
||||||
responses: Vec<BatchReadVersionResp>,
|
responses: Vec<BatchReadVersionResp>,
|
||||||
) -> Vec<crate::disk::error::Result<FileInfo>> {
|
) -> Vec<crate::disk::error::Result<FileInfo>> {
|
||||||
let mut results = (0..expected_items.len())
|
let mut results = (0..expected_items.len())
|
||||||
@@ -6878,6 +6893,7 @@ mod tests {
|
|||||||
},
|
},
|
||||||
];
|
];
|
||||||
|
|
||||||
|
let expected_items = expected_batch_read_version_items(&expected_items);
|
||||||
let mut results = map_batch_read_version_responses(&expected_items, responses).into_iter();
|
let mut results = map_batch_read_version_responses(&expected_items, responses).into_iter();
|
||||||
let first = results
|
let first = results
|
||||||
.next()
|
.next()
|
||||||
@@ -6918,6 +6934,7 @@ mod tests {
|
|||||||
version_id: "v-b".to_string(),
|
version_id: "v-b".to_string(),
|
||||||
},
|
},
|
||||||
];
|
];
|
||||||
|
let expected_items = expected_batch_read_version_items(&expected_items);
|
||||||
let results = map_batch_read_version_responses(
|
let results = map_batch_read_version_responses(
|
||||||
&expected_items,
|
&expected_items,
|
||||||
vec![
|
vec![
|
||||||
@@ -6957,6 +6974,7 @@ mod tests {
|
|||||||
path: "object-a".to_string(),
|
path: "object-a".to_string(),
|
||||||
version_id: "v-a".to_string(),
|
version_id: "v-a".to_string(),
|
||||||
}];
|
}];
|
||||||
|
let expected_items = expected_batch_read_version_items(&expected_items);
|
||||||
let mismatched = map_batch_read_version_responses(
|
let mismatched = map_batch_read_version_responses(
|
||||||
&expected_items,
|
&expected_items,
|
||||||
vec![BatchReadVersionResp {
|
vec![BatchReadVersionResp {
|
||||||
@@ -7018,6 +7036,10 @@ mod tests {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn expected_batch_read_version_items(items: &[BatchReadVersionItem]) -> Vec<ExpectedBatchReadVersionItem> {
|
||||||
|
items.iter().map(ExpectedBatchReadVersionItem::from).collect()
|
||||||
|
}
|
||||||
|
|
||||||
/// Isolation guard: unobserved objects record nothing (so parallel tests do
|
/// Isolation guard: unobserved objects record nothing (so parallel tests do
|
||||||
/// not inflate one another), and a scope clears its own counts on drop.
|
/// not inflate one another), and a scope clears its own counts on drop.
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
|
|||||||
@@ -1203,7 +1203,8 @@ fn inject_batch_delete_pool_errors(
|
|||||||
bucket: &str,
|
bucket: &str,
|
||||||
pool_idx: usize,
|
pool_idx: usize,
|
||||||
object_names: &[String],
|
object_names: &[String],
|
||||||
result: &mut (Vec<DeletedObject>, Vec<Option<Error>>),
|
deleted: &[DeletedObject],
|
||||||
|
errors: &mut [Option<Error>],
|
||||||
) {
|
) {
|
||||||
let state = BATCH_DELETE_POOL_ERROR_INJECTION
|
let state = BATCH_DELETE_POOL_ERROR_INJECTION
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
.get_or_init(|| std::sync::Mutex::new(None))
|
||||||
@@ -1220,8 +1221,8 @@ fn inject_batch_delete_pool_errors(
|
|||||||
let Some(error) = state.errors.get(object_name) else {
|
let Some(error) = state.errors.get(object_name) else {
|
||||||
continue;
|
continue;
|
||||||
};
|
};
|
||||||
if result.1[idx].is_none() && result.0[idx].found {
|
if errors[idx].is_none() && deleted[idx].found {
|
||||||
result.1[idx] = Some(error.clone());
|
errors[idx] = Some(error.clone());
|
||||||
state.observed.fetch_add(1, Ordering::AcqRel);
|
state.observed.fetch_add(1, Ordering::AcqRel);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -3194,7 +3195,7 @@ impl ECStore {
|
|||||||
|
|
||||||
// Default return value
|
// Default return value
|
||||||
let mut del_objects = vec![DeletedObject::default(); objects.len()];
|
let mut del_objects = vec![DeletedObject::default(); objects.len()];
|
||||||
let accounting = vec![None; objects.len()];
|
let mut accounting = vec![None; objects.len()];
|
||||||
|
|
||||||
let mut del_errs = Vec::with_capacity(objects.len());
|
let mut del_errs = Vec::with_capacity(objects.len());
|
||||||
for _ in 0..objects.len() {
|
for _ in 0..objects.len() {
|
||||||
@@ -3333,11 +3334,12 @@ impl ECStore {
|
|||||||
.iter()
|
.iter()
|
||||||
.map(|object| object.object_name.clone())
|
.map(|object| object.object_name.clone())
|
||||||
.collect::<Vec<_>>();
|
.collect::<Vec<_>>();
|
||||||
let result = pool.delete_objects(bucket, pool_objects, pool_opts).await;
|
let result = pool.delete_objects_with_accounting(bucket, pool_objects, pool_opts).await;
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
let result = {
|
let result = {
|
||||||
let mut result = result;
|
let mut result = result;
|
||||||
inject_batch_delete_pool_errors(bucket, pool.pool_idx, &pool_object_names, &mut result);
|
let (deleted, errors, _) = &mut result;
|
||||||
|
inject_batch_delete_pool_errors(bucket, pool.pool_idx, &pool_object_names, deleted, errors);
|
||||||
result
|
result
|
||||||
};
|
};
|
||||||
(object_indices, result)
|
(object_indices, result)
|
||||||
@@ -3347,7 +3349,7 @@ impl ECStore {
|
|||||||
let results = join_all(futures).await;
|
let results = join_all(futures).await;
|
||||||
|
|
||||||
for idx in 0..del_objects.len() {
|
for idx in 0..del_objects.len() {
|
||||||
let pool_results = results.iter().filter_map(|(object_indices, (dels, errs))| {
|
let pool_results = results.iter().filter_map(|(object_indices, (dels, errs, _))| {
|
||||||
let pool_object_idx = object_indices.binary_search(&idx).ok()?;
|
let pool_object_idx = object_indices.binary_search(&idx).ok()?;
|
||||||
Some((&dels[pool_object_idx], &errs[pool_object_idx]))
|
Some((&dels[pool_object_idx], &errs[pool_object_idx]))
|
||||||
});
|
});
|
||||||
@@ -3367,6 +3369,12 @@ impl ECStore {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
for (object_indices, (_, _, pool_accounting)) in &results {
|
||||||
|
for (pool_object_idx, object_idx) in object_indices.iter().enumerate() {
|
||||||
|
accounting[*object_idx] = pool_accounting.get(pool_object_idx).cloned().flatten();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
for (idx, object) in objects.iter().enumerate() {
|
for (idx, object) in objects.iter().enumerate() {
|
||||||
if del_errs[idx].is_none() && del_objects[idx].delete_marker {
|
if del_errs[idx].is_none() && del_objects[idx].delete_marker {
|
||||||
|
|||||||
Reference in New Issue
Block a user