Skip to content

Commit 009d461

Browse files
committed
feat(upload): upload multipart parts concurrently
Reuse the existing "batch transfer concurrency" setting to upload several parts of a large file at once instead of one at a time. Dispatch happens in waves sized to that concurrency: every part in a wave is allowed to finish (success or failure) before deciding whether to start the next wave, so a failing part doesn't discard sibling parts that were already in flight, and completed parts still land in the resumable session for a later retry.
1 parent e499103 commit 009d461

1 file changed

Lines changed: 130 additions & 22 deletions

File tree

crates/app-core/src/upload.rs

Lines changed: 130 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -7,11 +7,16 @@
77
//!
88
//! 只对超过一个分片大小的文件启用;小文件、以及 `resumable_upload` 能力为 false 的 provider
99
//! 回退到整体上传。出错时**不**放弃服务端已上传分片(不 abort),以便下次续传。
10+
//!
11+
//! 分片之间**并发**上传(并发度复用 [`Settings::concurrency`](crate::Settings) 这一个
12+
//! 用户已可调的旋钮,不单独引入新设置),每个分片各开一个文件句柄独立 seek+读,避免共享游标。
1013
1114
use std::collections::HashSet;
12-
use std::sync::atomic::{AtomicBool, Ordering};
15+
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
16+
use std::sync::Mutex;
1317

1418
use bytes::{Bytes, BytesMut};
19+
use futures::stream::{self, StreamExt};
1520
use tokio::io::{AsyncReadExt, AsyncSeekExt};
1621

1722
use crate::content_type::guess_content_type;
@@ -81,7 +86,7 @@ impl App {
8186
let key = session_key(account, remote_path, local_path);
8287

8388
// 恢复匹配的会话,否则新开一个。
84-
let (upload_id, mut done) = match self.load_session(&key, size, mtime, part_size) {
89+
let (upload_id, done) = match self.load_session(&key, size, mtime, part_size) {
8590
Some(resumed) => resumed,
8691
None => {
8792
let uid = provider.begin_multipart(remote_path, content_type).await?;
@@ -93,37 +98,79 @@ impl App {
9398

9499
let num_parts = size.div_ceil(part_size) as u32;
95100
let done_nums: HashSet<u32> = done.iter().map(|(n, _)| *n).collect();
96-
let mut uploaded: u64 = done_nums
101+
let uploaded_init: u64 = done_nums
97102
.iter()
98103
.map(|n| part_len(*n, num_parts, size, part_size))
99104
.sum();
100-
progress(uploaded, size);
101-
102-
let mut file = tokio::fs::File::open(local_path).await?;
103-
for n in 1..=num_parts {
104-
if done_nums.contains(&n) {
105-
continue;
106-
}
107-
// 取消:中止但保留会话与已传分片,重发即续传。
108-
if cancel.load(Ordering::Relaxed) {
109-
return Err(AppError::Cancelled);
110-
}
105+
progress(uploaded_init, size);
106+
107+
let pending: Vec<u32> = (1..=num_parts).filter(|n| !done_nums.contains(n)).collect();
108+
// 并发度复用「批量传输并发数」这一个设置,1..=10(前端已限定范围),这里再兜底一次。
109+
let concurrency = (self.settings().concurrency as usize).clamp(1, 10);
110+
let done = Mutex::new(done);
111+
let uploaded = AtomicU64::new(uploaded_init);
112+
// 提前把非 Copy 的捕获量重绑成引用:`async move` 按值捕获,若不这样做,
113+
// 外层 `map` 闭包在第二次调用时会尝试重新移动已经移走的值,导致只能实现 FnOnce。
114+
let provider = &provider;
115+
let key = &key;
116+
let upload_id = &upload_id;
117+
let done_ref = &done;
118+
let uploaded_ref = &uploaded;
119+
120+
let upload_one = move |n: u32| async move {
111121
let len = part_len(n, num_parts, size, part_size);
122+
// 各分片独立开文件句柄各自 seek+读,避免共享一个游标在并发下互相踩。
123+
let mut file = tokio::fs::File::open(local_path).await?;
112124
let bytes = read_part(&mut file, (n as u64 - 1) * part_size, len).await?;
113125
let etag = provider
114-
.upload_part(remote_path, &upload_id, n, bytes)
126+
.upload_part(remote_path, upload_id, n, bytes)
115127
.await?;
116-
done.push((n, etag));
117-
self.save_session(&key, &upload_id, size, mtime, part_size, &done);
118-
uploaded += len;
119-
progress(uploaded, size);
128+
{
129+
let mut d = done_ref.lock().unwrap();
130+
d.push((n, etag));
131+
self.save_session(key, upload_id, size, mtime, part_size, &d);
132+
}
133+
let total = uploaded_ref.fetch_add(len, Ordering::SeqCst) + len;
134+
progress(total, size);
135+
Ok(())
136+
};
137+
138+
// 按并发度分波:每一波内的分片并发上传、**全部落定**(不管成败)后才决定要不要下一波。
139+
// 这样一波内和分片 X 同批调度的其它分片不会因为 X 失败而被半路丢弃、白白浪费已发出的请求;
140+
// 但失败的那一波过后就不再开新的一波——已完成的进度全部存进会话,供下次续传。
141+
let mut first_err: Option<AppError> = None;
142+
for chunk in pending.chunks(concurrency) {
143+
if cancel.load(Ordering::Relaxed) {
144+
first_err = Some(AppError::Cancelled);
145+
break;
146+
}
147+
let results: Vec<Result<()>> = stream::iter(chunk.iter().copied().map(upload_one))
148+
.buffer_unordered(chunk.len())
149+
.collect()
150+
.await;
151+
let mut wave_failed = false;
152+
for r in results {
153+
if let Err(e) = r {
154+
wave_failed = true;
155+
if first_err.is_none() {
156+
first_err = Some(e);
157+
}
158+
}
159+
}
160+
if wave_failed {
161+
break;
162+
}
163+
}
164+
if let Some(e) = first_err {
165+
return Err(e);
120166
}
121167

168+
let mut done = done.into_inner().unwrap();
122169
done.sort_by_key(|(n, _)| *n);
123170
provider
124-
.complete_multipart(remote_path, &upload_id, &done)
171+
.complete_multipart(remote_path, upload_id, &done)
125172
.await?;
126-
self.delete_session(&key);
173+
self.delete_session(key);
127174
Ok(())
128175
}
129176

@@ -318,6 +365,10 @@ mod tests {
318365
async fn resumable_upload_sends_all_parts_in_order() {
319366
let rec = Arc::new(ResumableRec::new(None));
320367
let (app, db) = app_with_store("all", rec.clone());
368+
// 并发度钉在 1:这个用例要测的是"分片按顺序传完",并发下完成顺序不再保证。
369+
let mut s = app.settings();
370+
s.concurrency = 1;
371+
app.save_settings(&s).unwrap();
321372
let file = temp_file("all", &[7u8; 10]); // 10 字节,分片 4 → (4,4,2)
322373

323374
app.upload_resumable_parted(
@@ -348,9 +399,13 @@ mod tests {
348399

349400
#[tokio::test]
350401
async fn resume_after_interruption_skips_completed_parts() {
351-
// 让分片 2 首次失败。
402+
// 并发度钉在 1:这个用例要测的是"续传跳过已完成分片",不是并发本身——
403+
// 并发>1 时分片 2 失败不妨碍同批次的分片 3 照样成功,断言会变得依赖调度顺序。
352404
let rec = Arc::new(ResumableRec::new(Some(2)));
353405
let (app, db) = app_with_store("resume", rec.clone());
406+
let mut s = app.settings();
407+
s.concurrency = 1;
408+
app.save_settings(&s).unwrap();
354409
let file = temp_file("resume", &[9u8; 10]); // (4,4,2)
355410
let path = file.to_str().unwrap();
356411

@@ -389,10 +444,63 @@ mod tests {
389444
let _ = std::fs::remove_file(&db);
390445
}
391446

447+
#[tokio::test]
448+
async fn concurrent_parts_sharing_a_failure_still_persist_their_progress() {
449+
// 并发度 3、3 个分片:分片 2 失败不妨碍同批次并发调度的分片 1/3 成功并入会话,
450+
// 续传时只需再补分片 2——这正是"并发"相对旧的严格顺序循环带来的行为变化。
451+
let rec = Arc::new(ResumableRec::new(Some(2)));
452+
let (app, db) = app_with_store("concurrent", rec.clone());
453+
let mut s = app.settings();
454+
s.concurrency = 3;
455+
app.save_settings(&s).unwrap();
456+
let file = temp_file("concurrent", &[9u8; 10]); // (4,4,2)
457+
let path = file.to_str().unwrap();
458+
459+
assert!(app
460+
.upload_resumable_parted(
461+
"rec",
462+
"b/k",
463+
path,
464+
None,
465+
4,
466+
&AtomicBool::new(false),
467+
&|_, _| {}
468+
)
469+
.await
470+
.is_err());
471+
let mut got = rec.parts.lock().unwrap().clone();
472+
got.sort();
473+
assert_eq!(got, vec![(1, 4), (3, 2)], "分片 1、3 应已并发完成并持久化");
474+
assert!(rec.completed.lock().unwrap().is_none());
475+
476+
app.upload_resumable_parted(
477+
"rec",
478+
"b/k",
479+
path,
480+
None,
481+
4,
482+
&AtomicBool::new(false),
483+
&|_, _| {},
484+
)
485+
.await
486+
.unwrap();
487+
let mut got = rec.parts.lock().unwrap().clone();
488+
got.sort();
489+
assert_eq!(got, vec![(1, 4), (2, 4), (3, 2)], "续传只应再补分片 2");
490+
assert!(rec.completed.lock().unwrap().is_some());
491+
492+
let _ = std::fs::remove_file(&file);
493+
let _ = std::fs::remove_file(&db);
494+
}
495+
392496
#[tokio::test]
393497
async fn cancelled_upload_stops_and_resumes_later() {
394498
let rec = Arc::new(ResumableRec::new(None));
395499
let (app, db) = app_with_store("cancel", rec.clone());
500+
// 并发度钉在 1,让续传后的分片顺序断言保持确定性。
501+
let mut s = app.settings();
502+
s.concurrency = 1;
503+
app.save_settings(&s).unwrap();
396504
let file = temp_file("cancel", &[5u8; 10]); // (4,4,2)
397505
let path = file.to_str().unwrap();
398506

0 commit comments

Comments
 (0)