From 5e1445e03e6fda5ba7541a748b11468d4161d320 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sat, 8 Aug 2026 09:41:46 -0700 Subject: [PATCH 1/4] feat(storage): retry temporary OSS failures --- Cargo.lock | 35 ++++++ crates/paimon/Cargo.toml | 7 +- crates/paimon/src/io/storage.rs | 8 +- crates/paimon/src/io/storage_oss.rs | 159 ++++++++++++++++++++++++---- docs/src/getting-started.md | 3 + 5 files changed, 185 insertions(+), 27 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 87cdeca74..9f078c3b1 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -725,6 +725,17 @@ dependencies = [ "syn 2.0.119", ] +[[package]] +name = "backon" +version = "1.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cffb0e931875b666fc4fcb20fee52e9bbd1ef836fd9e9e04ec21555f9f85f7ef" +dependencies = [ + "fastrand", + "gloo-timers", + "tokio", +] + [[package]] name = "base64" version = "0.22.1" @@ -2813,6 +2824,18 @@ version = "0.3.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0cc23270f6e1808e30a928bdc84dea0b9b4136a8bc82338574f23baf47bbd280" +[[package]] +name = "gloo-timers" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bbb143cf96099802033e0d4f4963b19fd2e0b728bcf076cd9cf7f6634f092994" +dependencies = [ + "futures-channel", + "futures-core", + "js-sys", + "wasm-bindgen", +] + [[package]] name = "h2" version = "0.4.15" @@ -4254,6 +4277,17 @@ dependencies = [ "reqwest 0.13.4", ] +[[package]] +name = "opendal-layer-retry" +version = "0.58.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f2df70875ab7fd6f80720d4787c49c70883cef0d81bfae947ecba88b8d1cd62e" +dependencies = [ + "backon", + "log", + "opendal-core", +] + [[package]] name = "opendal-service-azdls" version = "0.58.0" @@ -4546,6 +4580,7 @@ dependencies = [ "md-5 0.10.6", "opendal-core", "opendal-http-transport-reqwest", + "opendal-layer-retry", "opendal-service-azdls", "opendal-service-cos", "opendal-service-fs", diff --git a/crates/paimon/Cargo.toml b/crates/paimon/Cargo.toml index 5922270da..ee5beacfd 100644 --- a/crates/paimon/Cargo.toml +++ b/crates/paimon/Cargo.toml @@ -47,7 +47,11 @@ vortex = ["dep:vortex"] storage-memory = ["opendal/services-memory"] storage-fs = ["dep:opendal-service-fs"] -storage-oss = ["dep:opendal-http-transport-reqwest", "dep:opendal-service-oss"] +storage-oss = [ + "dep:opendal-http-transport-reqwest", + "dep:opendal-layer-retry", + "dep:opendal-service-oss", +] storage-s3 = ["dep:opendal-http-transport-reqwest", "dep:opendal-service-s3"] storage-cos = ["dep:opendal-http-transport-reqwest", "dep:opendal-service-cos"] storage-azdls = ["dep:opendal-http-transport-reqwest", "dep:opendal-service-azdls"] @@ -71,6 +75,7 @@ snafu = "0.9.0" typed-builder = "^0.19" opendal = { package = "opendal-core", version = "0.58.0" } opendal-http-transport-reqwest = { version = "0.58.0", optional = true } +opendal-layer-retry = { version = "0.58.0", optional = true } opendal-service-azdls = { version = "0.58.0", optional = true } opendal-service-cos = { version = "0.58.0", optional = true } opendal-service-fs = { version = "0.58.0", optional = true } diff --git a/crates/paimon/src/io/storage.rs b/crates/paimon/src/io/storage.rs index eb37b2af9..5ba31e715 100644 --- a/crates/paimon/src/io/storage.rs +++ b/crates/paimon/src/io/storage.rs @@ -39,6 +39,8 @@ use std::sync::MutexGuard; #[cfg(feature = "storage-azdls")] use super::AzdlsStorageConfig; +#[cfg(feature = "storage-oss")] +use super::OssStorageConfig; use opendal::Operator; #[cfg(feature = "storage-cos")] use opendal_service_cos::CosConfig; @@ -48,8 +50,6 @@ use opendal_service_gcs::GcsConfig; use opendal_service_hdfs_native::HdfsNativeConfig; #[cfg(feature = "storage-obs")] use opendal_service_obs::ObsConfig; -#[cfg(feature = "storage-oss")] -use opendal_service_oss::OssConfig; #[cfg(feature = "storage-s3")] use opendal_service_s3::S3Config; #[cfg(any( @@ -77,7 +77,7 @@ pub enum Storage { LocalFs { op: Operator }, #[cfg(feature = "storage-oss")] Oss { - config: Box, + config: Box, operators: Mutex>, }, #[cfg(feature = "storage-s3")] @@ -412,7 +412,7 @@ impl Storage { #[cfg(feature = "storage-oss")] fn cached_oss_operator( - config: &OssConfig, + config: &OssStorageConfig, operators: &Mutex>, path: &str, bucket: &str, diff --git a/crates/paimon/src/io/storage_oss.rs b/crates/paimon/src/io/storage_oss.rs index 77b883a71..ee32d7d04 100644 --- a/crates/paimon/src/io/storage_oss.rs +++ b/crates/paimon/src/io/storage_oss.rs @@ -16,8 +16,10 @@ // under the License. use std::collections::HashMap; +use std::time::Duration; use opendal::{Configurator, Operator}; +use opendal_layer_retry::RetryLayer; use opendal_service_oss::OssConfig; use url::Url; @@ -45,14 +47,29 @@ pub(crate) const OSS_ACCESS_KEY_SECRET: &str = "fs.oss.accessKeySecret"; /// Required when using STS temporary credentials (e.g. from REST data tokens). pub(crate) const OSS_SECURITY_TOKEN: &str = "fs.oss.securityToken"; -/// Parse paimon catalog options into an [`OssConfig`]. +/// Number of retries after an OSS request fails. +pub(crate) const OSS_RETRY_COUNT: &str = "fs.oss.retry.count"; + +/// Initial exponential retry interval in milliseconds. +pub(crate) const OSS_RETRY_INTERVAL_MILLIS: &str = "fs.oss.retry.interval.millisecond"; + +const DEFAULT_OSS_RETRY_COUNT: usize = 5; +const DEFAULT_OSS_RETRY_INTERVAL_MILLIS: u64 = 500; + +#[derive(Debug)] +pub struct OssStorageConfig { + service: OssConfig, + retry_count: usize, + retry_interval: Duration, +} + +/// Parse paimon catalog options into an [`OssStorageConfig`]. /// /// Extracts OSS-related configuration keys (endpoint, access key, secret key, -/// and optional security token) from the provided properties map and maps them -/// to the corresponding [`OssConfig`] fields. +/// optional security token, and retry settings) from the provided properties. /// /// Returns an error if any required configuration key is missing. -pub(crate) fn oss_config_parse(mut props: HashMap) -> Result { +pub(crate) fn oss_config_parse(mut props: HashMap) -> Result { let mut cfg = OssConfig::default(); cfg.endpoint = Some( @@ -80,14 +97,36 @@ pub(crate) fn oss_config_parse(mut props: HashMap) -> Result(props: &mut HashMap, key: &str, default: T) -> Result +where + T: std::str::FromStr, +{ + match props.remove(key) { + Some(value) => value.parse().map_err(|_| Error::ConfigInvalid { + message: format!("Invalid OSS config {key}: {value}"), + }), + None => Ok(default), + } } /// Build an [`Operator`] for the given OSS path. /// /// Parses the bucket name from the `oss://bucket/key` URL and combines it -/// with the provided [`OssConfig`] to construct an OpenDAL operator. -pub(crate) fn oss_config_build(cfg: &OssConfig, path: &str) -> Result { +/// with the provided [`OssStorageConfig`] to construct an OpenDAL operator. +pub(crate) fn oss_config_build(cfg: &OssStorageConfig, path: &str) -> Result { let url = Url::parse(path).map_err(|_| Error::ConfigInvalid { message: format!("Invalid OSS url: {path}"), })?; @@ -96,31 +135,87 @@ pub(crate) fn oss_config_build(cfg: &OssConfig, path: &str) -> Result message: format!("Invalid OSS url: {path}, missing bucket"), })?; - let builder = cfg.clone().into_builder().bucket(bucket); - Ok(super::with_http_transport(Operator::new(builder)?)) + let builder = cfg.service.clone().into_builder().bucket(bucket); + let retry = RetryLayer::default() + .with_min_delay(cfg.retry_interval) + .with_max_times(cfg.retry_count) + .with_jitter(); + Ok(super::with_http_transport(Operator::new(builder)?).layer(retry)) } #[cfg(test)] mod tests { + use std::sync::atomic::{AtomicUsize, Ordering}; + use std::sync::Arc; + + use axum::body::Body; + use axum::extract::State; + use axum::http::{Response, StatusCode}; + use axum::routing::get; + use axum::Router; + use super::*; + fn storage_config(service: OssConfig) -> OssStorageConfig { + OssStorageConfig { + service, + retry_count: DEFAULT_OSS_RETRY_COUNT, + retry_interval: Duration::from_millis(1), + } + } + + fn required_props() -> HashMap { + HashMap::from([ + ( + OSS_ENDPOINT.to_string(), + "https://oss-cn-hangzhou.aliyuncs.com".to_string(), + ), + (OSS_ACCESS_KEY_ID.to_string(), "test-ak".to_string()), + (OSS_ACCESS_KEY_SECRET.to_string(), "test-sk".to_string()), + ]) + } + + async fn retry_once(State(attempts): State>) -> Response { + if attempts.fetch_add(1, Ordering::SeqCst) == 0 { + return Response::builder() + .status(StatusCode::SERVICE_UNAVAILABLE) + .body(Body::from( + "QpsLimitExceededretry", + )) + .unwrap(); + } + Response::new(Body::from("ok")) + } + #[test] fn test_oss_config_parse_with_all_keys() { - let mut props = HashMap::new(); - props.insert( - OSS_ENDPOINT.to_string(), - "https://oss-cn-hangzhou.aliyuncs.com".to_string(), - ); - props.insert(OSS_ACCESS_KEY_ID.to_string(), "test-ak".to_string()); - props.insert(OSS_ACCESS_KEY_SECRET.to_string(), "test-sk".to_string()); + let mut props = required_props(); + props.insert(OSS_RETRY_COUNT.to_string(), "7".to_string()); + props.insert(OSS_RETRY_INTERVAL_MILLIS.to_string(), "250".to_string()); let cfg = oss_config_parse(props).unwrap(); assert_eq!( - cfg.endpoint.as_deref(), + cfg.service.endpoint.as_deref(), Some("https://oss-cn-hangzhou.aliyuncs.com") ); - assert_eq!(cfg.access_key_id.as_deref(), Some("test-ak")); - assert_eq!(cfg.access_key_secret.as_deref(), Some("test-sk")); + assert_eq!(cfg.service.access_key_id.as_deref(), Some("test-ak")); + assert_eq!(cfg.service.access_key_secret.as_deref(), Some("test-sk")); + assert_eq!(cfg.retry_count, 7); + assert_eq!(cfg.retry_interval, Duration::from_millis(250)); + } + + #[test] + fn test_oss_retry_defaults_and_validation() { + let cfg = oss_config_parse(required_props()).unwrap(); + assert_eq!(cfg.retry_count, DEFAULT_OSS_RETRY_COUNT); + assert_eq!( + cfg.retry_interval, + Duration::from_millis(DEFAULT_OSS_RETRY_INTERVAL_MILLIS) + ); + + let mut props = required_props(); + props.insert(OSS_RETRY_COUNT.to_string(), "invalid".to_string()); + assert!(oss_config_parse(props).is_err()); } #[test] @@ -128,21 +223,41 @@ mod tests { let mut cfg = OssConfig::default(); cfg.endpoint = Some("https://oss-cn-hangzhou.aliyuncs.com".to_string()); - let op = oss_config_build(&cfg, "oss://my-bucket/some/path").unwrap(); + let op = oss_config_build(&storage_config(cfg), "oss://my-bucket/some/path").unwrap(); assert_eq!(op.info().name(), "my-bucket"); } #[test] fn test_oss_config_build_invalid_url() { - let cfg = OssConfig::default(); + let cfg = storage_config(OssConfig::default()); let result = oss_config_build(&cfg, "not-a-valid-url"); assert!(result.is_err()); } #[test] fn test_oss_config_build_missing_bucket() { - let cfg = OssConfig::default(); + let cfg = storage_config(OssConfig::default()); let result = oss_config_build(&cfg, "oss:///path/without/bucket"); assert!(result.is_err()); } + + #[tokio::test] + async fn test_oss_retries_temporary_failure() { + let attempts = Arc::new(AtomicUsize::new(0)); + let app = Router::new() + .fallback(get(retry_once)) + .with_state(attempts.clone()); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); + + let mut cfg = OssConfig::default(); + cfg.endpoint = Some(format!("http://{address}")); + cfg.addressing_style = Some("path".to_string()); + cfg.skip_signature = true; + + let op = oss_config_build(&storage_config(cfg), "oss://bucket/path").unwrap(); + assert_eq!(op.read("object").await.unwrap().to_bytes(), "ok"); + assert_eq!(attempts.load(Ordering::SeqCst), 2); + } } diff --git a/docs/src/getting-started.md b/docs/src/getting-started.md index 5e7ee20c2..ca86121a0 100644 --- a/docs/src/getting-started.md +++ b/docs/src/getting-started.md @@ -85,6 +85,9 @@ options.set(CatalogOptions::WAREHOUSE, "oss://bucket/warehouse"); options.set("fs.oss.accessKeyId", "your-access-key-id"); options.set("fs.oss.accessKeySecret", "your-access-key-secret"); options.set("fs.oss.endpoint", "oss-cn-hangzhou.aliyuncs.com"); +// Optional: configure retries for temporary OSS failures. +options.set("fs.oss.retry.count", "5"); +options.set("fs.oss.retry.interval.millisecond", "500"); let catalog = CatalogFactory::create(options).await?; // Tencent Cloud COS From a27dc85d72a695de4bb8895036099f232f4c542b Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sat, 8 Aug 2026 10:35:05 -0700 Subject: [PATCH 2/4] chore: update Rust dependency reports --- DEPENDENCIES.rust.tsv | 3 +++ benchmarks/tpcds/DEPENDENCIES.rust.tsv | 3 +++ bindings/c/DEPENDENCIES.rust.tsv | 3 +++ bindings/go/DEPENDENCIES.rust.tsv | 3 +++ bindings/python/DEPENDENCIES.rust.tsv | 3 +++ crates/integration_tests/DEPENDENCIES.rust.tsv | 3 +++ crates/integrations/datafusion/DEPENDENCIES.rust.tsv | 3 +++ crates/paimon-rest-server/DEPENDENCIES.rust.tsv | 3 +++ crates/paimon/DEPENDENCIES.rust.tsv | 3 +++ 9 files changed, 27 insertions(+) diff --git a/DEPENDENCIES.rust.tsv b/DEPENDENCIES.rust.tsv index 7b2d77729..f2985fdb5 100644 --- a/DEPENDENCIES.rust.tsv +++ b/DEPENDENCIES.rust.tsv @@ -58,6 +58,7 @@ aws-lc-sys@0.43.0 X X X X X axum@0.7.9 X axum-core@0.4.5 X axum-macros@0.4.2 X +backon@1.6.0 X base64@0.22.1 X X base64ct@1.8.3 X X better_io@0.2.0 X @@ -234,6 +235,7 @@ getrandom@0.2.17 X X getrandom@0.3.4 X X getrandom@0.4.3 X X glob@0.3.3 X X +gloo-timers@0.3.0 X X h2@0.4.15 X half@2.7.1 X X hashbrown@0.14.5 X X @@ -368,6 +370,7 @@ oneshot@0.1.13 X X oneshot@0.2.1 X X opendal-core@0.58.0 X opendal-http-transport-reqwest@0.58.0 X +opendal-layer-retry@0.58.0 X opendal-service-azdls@0.58.0 X opendal-service-azure-common@0.58.0 X opendal-service-cos@0.58.0 X diff --git a/benchmarks/tpcds/DEPENDENCIES.rust.tsv b/benchmarks/tpcds/DEPENDENCIES.rust.tsv index 4b76a75bd..e31ea76a2 100644 --- a/benchmarks/tpcds/DEPENDENCIES.rust.tsv +++ b/benchmarks/tpcds/DEPENDENCIES.rust.tsv @@ -40,6 +40,7 @@ atomic-waker@1.1.2 X X autocfg@1.5.1 X X aws-lc-rs@1.17.3 X X aws-lc-sys@0.43.0 X X X X X +backon@1.6.0 X base64@0.22.1 X X bigdecimal@0.4.10 X X bitflags@2.13.1 X X @@ -163,6 +164,7 @@ getrandom@0.2.17 X X getrandom@0.3.4 X X getrandom@0.4.3 X X glob@0.3.3 X X +gloo-timers@0.3.0 X X h2@0.4.15 X half@2.7.1 X X hashbrown@0.14.5 X X @@ -254,6 +256,7 @@ once_cell@1.21.4 X X once_cell_polyfill@1.70.2 X X opendal-core@0.58.0 X opendal-http-transport-reqwest@0.58.0 X +opendal-layer-retry@0.58.0 X opendal-service-fs@0.58.0 X opendal-service-oss@0.58.0 X openssl@0.10.81 X diff --git a/bindings/c/DEPENDENCIES.rust.tsv b/bindings/c/DEPENDENCIES.rust.tsv index 244fb1100..a01464e35 100644 --- a/bindings/c/DEPENDENCIES.rust.tsv +++ b/bindings/c/DEPENDENCIES.rust.tsv @@ -30,6 +30,7 @@ atomic-waker@1.1.2 X X autocfg@1.5.1 X X aws-lc-rs@1.17.3 X X aws-lc-sys@0.43.0 X X X X X +backon@1.6.0 X base64@0.22.1 X X bigdecimal@0.4.10 X X bitflags@2.13.1 X X @@ -106,6 +107,7 @@ generic-array@0.14.7 X getrandom@0.2.17 X X getrandom@0.3.4 X X getrandom@0.4.3 X X +gloo-timers@0.3.0 X X h2@0.4.15 X half@2.7.1 X X hashbrown@0.14.5 X X @@ -187,6 +189,7 @@ num-traits@0.2.19 X X once_cell@1.21.4 X X opendal-core@0.58.0 X opendal-http-transport-reqwest@0.58.0 X +opendal-layer-retry@0.58.0 X opendal-service-fs@0.58.0 X opendal-service-oss@0.58.0 X openssl@0.10.81 X diff --git a/bindings/go/DEPENDENCIES.rust.tsv b/bindings/go/DEPENDENCIES.rust.tsv index 244fb1100..a01464e35 100644 --- a/bindings/go/DEPENDENCIES.rust.tsv +++ b/bindings/go/DEPENDENCIES.rust.tsv @@ -30,6 +30,7 @@ atomic-waker@1.1.2 X X autocfg@1.5.1 X X aws-lc-rs@1.17.3 X X aws-lc-sys@0.43.0 X X X X X +backon@1.6.0 X base64@0.22.1 X X bigdecimal@0.4.10 X X bitflags@2.13.1 X X @@ -106,6 +107,7 @@ generic-array@0.14.7 X getrandom@0.2.17 X X getrandom@0.3.4 X X getrandom@0.4.3 X X +gloo-timers@0.3.0 X X h2@0.4.15 X half@2.7.1 X X hashbrown@0.14.5 X X @@ -187,6 +189,7 @@ num-traits@0.2.19 X X once_cell@1.21.4 X X opendal-core@0.58.0 X opendal-http-transport-reqwest@0.58.0 X +opendal-layer-retry@0.58.0 X opendal-service-fs@0.58.0 X opendal-service-oss@0.58.0 X openssl@0.10.81 X diff --git a/bindings/python/DEPENDENCIES.rust.tsv b/bindings/python/DEPENDENCIES.rust.tsv index f4d35ed91..785a3f063 100644 --- a/bindings/python/DEPENDENCIES.rust.tsv +++ b/bindings/python/DEPENDENCIES.rust.tsv @@ -40,6 +40,7 @@ atomic-waker@1.1.2 X X autocfg@1.5.1 X X aws-lc-rs@1.17.3 X X aws-lc-sys@0.43.0 X X X X X +backon@1.6.0 X base64@0.22.1 X X base64ct@1.8.3 X X bigdecimal@0.4.10 X X @@ -187,6 +188,7 @@ getrandom@0.2.17 X X getrandom@0.3.4 X X getrandom@0.4.3 X X glob@0.3.3 X X +gloo-timers@0.3.0 X X h2@0.4.15 X half@2.7.1 X X hashbrown@0.14.5 X X @@ -301,6 +303,7 @@ once_cell@1.21.4 X X oneshot@0.1.13 X X opendal-core@0.58.0 X opendal-http-transport-reqwest@0.58.0 X +opendal-layer-retry@0.58.0 X opendal-service-azdls@0.58.0 X opendal-service-azure-common@0.58.0 X opendal-service-cos@0.58.0 X diff --git a/crates/integration_tests/DEPENDENCIES.rust.tsv b/crates/integration_tests/DEPENDENCIES.rust.tsv index 07114e666..aa9bfc158 100644 --- a/crates/integration_tests/DEPENDENCIES.rust.tsv +++ b/crates/integration_tests/DEPENDENCIES.rust.tsv @@ -30,6 +30,7 @@ atomic-waker@1.1.2 X X autocfg@1.5.1 X X aws-lc-rs@1.17.3 X X aws-lc-sys@0.43.0 X X X X X +backon@1.6.0 X base64@0.22.1 X X bigdecimal@0.4.10 X X bitflags@2.13.1 X X @@ -106,6 +107,7 @@ generic-array@0.14.7 X getrandom@0.2.17 X X getrandom@0.3.4 X X getrandom@0.4.3 X X +gloo-timers@0.3.0 X X h2@0.4.15 X half@2.7.1 X X hashbrown@0.14.5 X X @@ -187,6 +189,7 @@ num-traits@0.2.19 X X once_cell@1.21.4 X X opendal-core@0.58.0 X opendal-http-transport-reqwest@0.58.0 X +opendal-layer-retry@0.58.0 X opendal-service-fs@0.58.0 X opendal-service-oss@0.58.0 X openssl@0.10.81 X diff --git a/crates/integrations/datafusion/DEPENDENCIES.rust.tsv b/crates/integrations/datafusion/DEPENDENCIES.rust.tsv index 8b242ee9c..5e4ac46d4 100644 --- a/crates/integrations/datafusion/DEPENDENCIES.rust.tsv +++ b/crates/integrations/datafusion/DEPENDENCIES.rust.tsv @@ -47,6 +47,7 @@ atomic-waker@1.1.2 X X autocfg@1.5.1 X X aws-lc-rs@1.17.3 X X aws-lc-sys@0.43.0 X X X X X +backon@1.6.0 X base64@0.22.1 X X better_io@0.2.0 X bigdecimal@0.4.10 X X @@ -200,6 +201,7 @@ getrandom@0.2.17 X X getrandom@0.3.4 X X getrandom@0.4.3 X X glob@0.3.3 X X +gloo-timers@0.3.0 X X h2@0.4.15 X half@2.7.1 X X hashbrown@0.14.5 X X @@ -327,6 +329,7 @@ oneshot@0.1.13 X X oneshot@0.2.1 X X opendal-core@0.58.0 X opendal-http-transport-reqwest@0.58.0 X +opendal-layer-retry@0.58.0 X opendal-service-fs@0.58.0 X opendal-service-oss@0.58.0 X openssl@0.10.81 X diff --git a/crates/paimon-rest-server/DEPENDENCIES.rust.tsv b/crates/paimon-rest-server/DEPENDENCIES.rust.tsv index e068e1d72..39dd29f64 100644 --- a/crates/paimon-rest-server/DEPENDENCIES.rust.tsv +++ b/crates/paimon-rest-server/DEPENDENCIES.rust.tsv @@ -33,6 +33,7 @@ aws-lc-sys@0.43.0 X X X X X axum@0.7.9 X axum-core@0.4.5 X axum-macros@0.4.2 X +backon@1.6.0 X base64@0.22.1 X X bigdecimal@0.4.10 X X bitflags@2.13.1 X X @@ -109,6 +110,7 @@ generic-array@0.14.7 X getrandom@0.2.17 X X getrandom@0.3.4 X X getrandom@0.4.3 X X +gloo-timers@0.3.0 X X h2@0.4.15 X half@2.7.1 X X hashbrown@0.14.5 X X @@ -191,6 +193,7 @@ num-traits@0.2.19 X X once_cell@1.21.4 X X opendal-core@0.58.0 X opendal-http-transport-reqwest@0.58.0 X +opendal-layer-retry@0.58.0 X opendal-service-fs@0.58.0 X opendal-service-oss@0.58.0 X openssl@0.10.81 X diff --git a/crates/paimon/DEPENDENCIES.rust.tsv b/crates/paimon/DEPENDENCIES.rust.tsv index 1cce09cf6..24855b1c1 100644 --- a/crates/paimon/DEPENDENCIES.rust.tsv +++ b/crates/paimon/DEPENDENCIES.rust.tsv @@ -44,6 +44,7 @@ atomic-waker@1.1.2 X X autocfg@1.5.1 X X aws-lc-rs@1.17.3 X X aws-lc-sys@0.43.0 X X X X X +backon@1.6.0 X base64@0.22.1 X X base64ct@1.8.3 X X better_io@0.2.0 X @@ -173,6 +174,7 @@ getrandom@0.2.17 X X getrandom@0.3.4 X X getrandom@0.4.3 X X glob@0.3.3 X X +gloo-timers@0.3.0 X X h2@0.4.15 X half@2.7.1 X X hashbrown@0.14.5 X X @@ -297,6 +299,7 @@ oneshot@0.1.13 X X oneshot@0.2.1 X X opendal-core@0.58.0 X opendal-http-transport-reqwest@0.58.0 X +opendal-layer-retry@0.58.0 X opendal-service-azdls@0.58.0 X opendal-service-azure-common@0.58.0 X opendal-service-cos@0.58.0 X From 9ab0b2fe0d3bf530b60aa689aa4e7780ed8883fb Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sat, 8 Aug 2026 19:26:45 -0700 Subject: [PATCH 3/4] fix(storage): increase default OSS retries --- crates/paimon/src/io/storage_oss.rs | 2 +- docs/src/getting-started.md | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/crates/paimon/src/io/storage_oss.rs b/crates/paimon/src/io/storage_oss.rs index ee32d7d04..90410acc7 100644 --- a/crates/paimon/src/io/storage_oss.rs +++ b/crates/paimon/src/io/storage_oss.rs @@ -53,7 +53,7 @@ pub(crate) const OSS_RETRY_COUNT: &str = "fs.oss.retry.count"; /// Initial exponential retry interval in milliseconds. pub(crate) const OSS_RETRY_INTERVAL_MILLIS: &str = "fs.oss.retry.interval.millisecond"; -const DEFAULT_OSS_RETRY_COUNT: usize = 5; +const DEFAULT_OSS_RETRY_COUNT: usize = 10; const DEFAULT_OSS_RETRY_INTERVAL_MILLIS: u64 = 500; #[derive(Debug)] diff --git a/docs/src/getting-started.md b/docs/src/getting-started.md index ca86121a0..8d90895e1 100644 --- a/docs/src/getting-started.md +++ b/docs/src/getting-started.md @@ -86,7 +86,7 @@ options.set("fs.oss.accessKeyId", "your-access-key-id"); options.set("fs.oss.accessKeySecret", "your-access-key-secret"); options.set("fs.oss.endpoint", "oss-cn-hangzhou.aliyuncs.com"); // Optional: configure retries for temporary OSS failures. -options.set("fs.oss.retry.count", "5"); +options.set("fs.oss.retry.count", "10"); options.set("fs.oss.retry.interval.millisecond", "500"); let catalog = CatalogFactory::create(options).await?; From fafcdb59c26be6346e01dc5d5ca3972d0110eac2 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sat, 8 Aug 2026 19:29:57 -0700 Subject: [PATCH 4/4] test(storage): verify OSS retry count semantics --- crates/paimon/src/io/storage_oss.rs | 44 +++++++++++++++++++++++++---- 1 file changed, 38 insertions(+), 6 deletions(-) diff --git a/crates/paimon/src/io/storage_oss.rs b/crates/paimon/src/io/storage_oss.rs index 90410acc7..83923d707 100644 --- a/crates/paimon/src/io/storage_oss.rs +++ b/crates/paimon/src/io/storage_oss.rs @@ -175,18 +175,27 @@ mod tests { ]) } + fn temporary_failure() -> Response { + Response::builder() + .status(StatusCode::SERVICE_UNAVAILABLE) + .body(Body::from( + "QpsLimitExceededretry", + )) + .unwrap() + } + async fn retry_once(State(attempts): State>) -> Response { if attempts.fetch_add(1, Ordering::SeqCst) == 0 { - return Response::builder() - .status(StatusCode::SERVICE_UNAVAILABLE) - .body(Body::from( - "QpsLimitExceededretry", - )) - .unwrap(); + return temporary_failure(); } Response::new(Body::from("ok")) } + async fn always_fail(State(attempts): State>) -> Response { + attempts.fetch_add(1, Ordering::SeqCst); + temporary_failure() + } + #[test] fn test_oss_config_parse_with_all_keys() { let mut props = required_props(); @@ -260,4 +269,27 @@ mod tests { assert_eq!(op.read("object").await.unwrap().to_bytes(), "ok"); assert_eq!(attempts.load(Ordering::SeqCst), 2); } + + #[tokio::test] + async fn test_oss_retry_count_is_additional_attempts() { + let attempts = Arc::new(AtomicUsize::new(0)); + let app = Router::new() + .fallback(get(always_fail)) + .with_state(attempts.clone()); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); + + let mut props = required_props(); + props.insert(OSS_ENDPOINT.to_string(), format!("http://{address}")); + props.insert(OSS_RETRY_COUNT.to_string(), "1".to_string()); + props.insert(OSS_RETRY_INTERVAL_MILLIS.to_string(), "1".to_string()); + let mut cfg = oss_config_parse(props).unwrap(); + cfg.service.addressing_style = Some("path".to_string()); + cfg.service.skip_signature = true; + + let op = oss_config_build(&cfg, "oss://bucket/path").unwrap(); + assert!(op.read("object").await.is_err()); + assert_eq!(attempts.load(Ordering::SeqCst), 2); + } }