From b629d876cc661320da4b3977b7ad7fffa9f941be Mon Sep 17 00:00:00 2001 From: Mehul Arora Date: Mon, 21 Sep 2026 14:19:06 -0400 Subject: [PATCH] perf(lite): allow limiting multipart upload concurrency --- README.md | 15 +++ cli/src/cli.rs | 10 ++ lite/src/lib.rs | 1 + lite/src/multipart_limit.rs | 213 ++++++++++++++++++++++++++++++++++++ lite/src/server.rs | 37 +++++++ 5 files changed, 276 insertions(+) create mode 100644 lite/src/multipart_limit.rs diff --git a/README.md b/README.md index 9c3eb14e..d5ea18fd 100644 --- a/README.md +++ b/README.md @@ -213,6 +213,21 @@ Use `SL8_` prefixed environment variables, e.g.: SL8_FLUSH_INTERVAL=10ms ``` +For object storage sharing a disk with latency-sensitive WAL writes, Lite can +cap the active parts within each multipart upload: + +```bash +s2 lite --bucket my-bucket --multipart-upload-concurrency 1 +``` + +The same cap can be set with `S2LITE_MULTIPART_UPLOAD_CONCURRENCY`. It applies to +both memtable-flush and compaction uploads; multiple uploads can still run +concurrently. Ordinary PUTs, including WAL writes, are unaffected by the cap. +Lower concurrency trades bulk upload throughput for less contention, so measure +append latency and compaction backlog under your workload. Omitting the option +preserves the writer's default concurrency. The cap must be positive and does +not increase the writer's concurrency when set above its existing limit. + #### Design [Concepts](https://s2.dev/docs/concepts) diff --git a/cli/src/cli.rs b/cli/src/cli.rs index fe5d148a..7c0584a3 100644 --- a/cli/src/cli.rs +++ b/cli/src/cli.rs @@ -980,6 +980,16 @@ mod tests { use super::{Cli, Command, DiffArgs, DiffOutput, DiffResourceKind, IssueAccessTokenArgs}; + #[test] + fn lite_rejects_invalid_multipart_concurrency() { + for value in ["0".to_owned(), usize::MAX.to_string()] { + let error = + Cli::try_parse_from(["s2", "lite", "--multipart-upload-concurrency", &value]) + .unwrap_err(); + assert_eq!(error.kind(), clap::error::ErrorKind::ValueValidation); + } + } + fn issue_access_token_args_from(args: I) -> IssueAccessTokenArgs where I: IntoIterator, diff --git a/lite/src/lib.rs b/lite/src/lib.rs index 8d57a5cb..64924a2f 100644 --- a/lite/src/lib.rs +++ b/lite/src/lib.rs @@ -4,5 +4,6 @@ pub mod backend; pub mod handlers; pub mod init; pub mod metrics; +mod multipart_limit; pub mod server; pub mod stream_id; diff --git a/lite/src/multipart_limit.rs b/lite/src/multipart_limit.rs new file mode 100644 index 00000000..c9e69438 --- /dev/null +++ b/lite/src/multipart_limit.rs @@ -0,0 +1,213 @@ +//! Optional per-upload concurrency cap. Ordinary PUTs remain independent. + +use std::{fmt, num::NonZeroUsize, ops::Range, sync::Arc}; + +use async_trait::async_trait; +use bytes::Bytes; +use futures::stream::BoxStream; +use slatedb::object_store::{ + CopyOptions, GetOptions, GetResult, ListResult, MultipartUpload, ObjectMeta, ObjectStore, + PutMultipartOptions, PutOptions, PutPayload, PutResult, RenameOptions, Result, + limit::LimitUpload, path::Path, +}; + +#[derive(Debug)] +pub(crate) struct MultipartLimitStore { + inner: Arc, + concurrency: NonZeroUsize, +} + +impl MultipartLimitStore { + pub(crate) fn new(inner: Arc, concurrency: NonZeroUsize) -> Self { + Self { inner, concurrency } + } +} + +impl fmt::Display for MultipartLimitStore { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!( + f, + "MultipartLimitStore({}, {})", + self.concurrency, self.inner + ) + } +} + +#[async_trait] +impl ObjectStore for MultipartLimitStore { + async fn put_opts(&self, path: &Path, data: PutPayload, opts: PutOptions) -> Result { + self.inner.put_opts(path, data, opts).await + } + + async fn put_multipart_opts( + &self, + path: &Path, + opts: PutMultipartOptions, + ) -> Result> { + let upload = self.inner.put_multipart_opts(path, opts).await?; + Ok(Box::new(LimitUpload::new(upload, self.concurrency.get()))) + } + + async fn get_opts(&self, path: &Path, opts: GetOptions) -> Result { + self.inner.get_opts(path, opts).await + } + + async fn get_ranges(&self, path: &Path, ranges: &[Range]) -> Result> { + self.inner.get_ranges(path, ranges).await + } + + fn delete_stream( + &self, + locations: BoxStream<'static, Result>, + ) -> BoxStream<'static, Result> { + self.inner.delete_stream(locations) + } + + fn list(&self, prefix: Option<&Path>) -> BoxStream<'static, Result> { + self.inner.list(prefix) + } + + fn list_with_offset( + &self, + prefix: Option<&Path>, + offset: &Path, + ) -> BoxStream<'static, Result> { + self.inner.list_with_offset(prefix, offset) + } + + async fn list_with_delimiter(&self, prefix: Option<&Path>) -> Result { + self.inner.list_with_delimiter(prefix).await + } + + async fn copy_opts(&self, from: &Path, to: &Path, opts: CopyOptions) -> Result<()> { + self.inner.copy_opts(from, to, opts).await + } + + async fn rename_opts(&self, from: &Path, to: &Path, opts: RenameOptions) -> Result<()> { + self.inner.rename_opts(from, to, opts).await + } +} + +#[cfg(test)] +mod tests { + use std::time::Duration; + + use futures::future::try_join_all; + use slatedb::object_store::{ + Error, ObjectStoreExt, PutMode, + memory::InMemory, + throttle::{ThrottleConfig, ThrottledStore}, + }; + use tokio::time::Instant; + + use super::*; + + #[tokio::test(start_paused = true)] + async fn separate_uploads_have_independent_limits() { + let inner = ThrottledStore::new( + InMemory::new(), + ThrottleConfig { + wait_put_per_call: Duration::from_millis(100), + ..Default::default() + }, + ); + let store = MultipartLimitStore::new(Arc::new(inner), NonZeroUsize::new(1).unwrap()); + let mut first = store.put_multipart(&Path::from("first")).await.unwrap(); + let mut second = store.put_multipart(&Path::from("second")).await.unwrap(); + let parts = [&mut first, &mut second] + .into_iter() + .flat_map(|upload| { + ["a", "b"].map(|part| upload.put_part(Bytes::from_static(part.as_bytes()).into())) + }) + .collect::>(); + let started = Instant::now(); + try_join_all(parts).await.unwrap(); + assert_eq!(started.elapsed(), Duration::from_millis(200)); + first.complete().await.unwrap(); + second.complete().await.unwrap(); + for name in ["first", "second"] { + assert_eq!( + store + .get(&Path::from(name)) + .await + .unwrap() + .bytes() + .await + .unwrap(), + "ab" + ); + } + } + + #[tokio::test(start_paused = true)] + async fn limits_parts_without_blocking_ordinary_puts() { + let inner = ThrottledStore::new( + InMemory::new(), + ThrottleConfig { + wait_put_per_call: Duration::from_millis(100), + ..Default::default() + }, + ); + let store = MultipartLimitStore::new(Arc::new(inner), NonZeroUsize::new(1).unwrap()); + let path = Path::from("bulk"); + let ordinary_path = Path::from("wal"); + let mut upload = store.put_multipart(&path).await.unwrap(); + let parts = ["first", "second", "third"] + .map(|s| upload.put_part(Bytes::from_static(s.as_bytes()).into())); + let started = Instant::now(); + let (uploaded, ordinary_elapsed) = tokio::join!(try_join_all(parts), async { + store + .put(&ordinary_path, Bytes::from_static(b"durable").into()) + .await + .unwrap(); + started.elapsed() + }); + uploaded.unwrap(); + assert_eq!(ordinary_elapsed, Duration::from_millis(100)); + assert_eq!(started.elapsed(), Duration::from_millis(300)); + upload.complete().await.unwrap(); + assert_eq!( + store.get(&path).await.unwrap().bytes().await.unwrap(), + "firstsecondthird" + ); + assert_eq!( + store + .get(&ordinary_path) + .await + .unwrap() + .bytes() + .await + .unwrap(), + "durable" + ); + + let duplicate = store + .put_opts( + &ordinary_path, + Bytes::from_static(b"replacement").into(), + PutOptions { + mode: PutMode::Create, + ..Default::default() + }, + ) + .await; + assert!(matches!(duplicate, Err(Error::AlreadyExists { .. }))); + } + + #[tokio::test] + async fn abort_does_not_publish_partial_data() { + let store = + MultipartLimitStore::new(Arc::new(InMemory::new()), NonZeroUsize::new(1).unwrap()); + let path = Path::from("aborted"); + let mut upload = store.put_multipart(&path).await.unwrap(); + upload + .put_part(Bytes::from_static(b"partial").into()) + .await + .unwrap(); + upload.abort().await.unwrap(); + assert!(matches!( + store.get(&path).await, + Err(Error::NotFound { .. }) + )); + } +} diff --git a/lite/src/server.rs b/lite/src/server.rs index bb4ec354..19df269a 100644 --- a/lite/src/server.rs +++ b/lite/src/server.rs @@ -1,5 +1,6 @@ use std::{ net::SocketAddr, + num::NonZeroUsize, path::PathBuf, sync::Arc, time::{Duration, SystemTime}, @@ -83,6 +84,32 @@ pub struct LiteArgs { /// Maximum in-flight append metered bytes across all streams before admission blocks. #[arg(long, default_value = "128MiB")] pub append_inflight_bytes: ByteSize, + + /// Cap the number of active parts in each multipart object upload. + /// + /// Lower values can reduce interference with WAL writes on shared storage, + /// at the cost of bulk upload throughput. Applies to memtable flushes and + /// compaction output. This does not limit the number of concurrent uploads + /// or increase the upload concurrency chosen by the underlying writer. + /// When omitted, the object store writer's concurrency is unchanged. + #[arg( + long, + env = "S2LITE_MULTIPART_UPLOAD_CONCURRENCY", + value_name = "PARTS", + value_parser = parse_multipart_upload_concurrency + )] + pub multipart_upload_concurrency: Option, +} + +fn parse_multipart_upload_concurrency(value: &str) -> Result { + let concurrency = value.parse::().map_err(|e| e.to_string())?; + if concurrency.get() > tokio::sync::Semaphore::MAX_PERMITS { + return Err(format!( + "must be at most {}", + tokio::sync::Semaphore::MAX_PERMITS + )); + } + Ok(concurrency) } #[derive(Debug, Clone)] @@ -174,6 +201,16 @@ pub async fn run(args: LiteArgs) -> eyre::Result<()> { }; let object_store = init_object_store(&store_type).await?; + let object_store = match args.multipart_upload_concurrency { + Some(concurrency) => { + info!(%concurrency, "capping concurrent parts per multipart upload"); + Arc::new(crate::multipart_limit::MultipartLimitStore::new( + object_store, + concurrency, + )) as Arc + } + None => object_store, + }; let db_settings = slatedb::Settings::from_env_with_default( "SL8_",