comphead commented on code in PR #3263:
URL: https://github.com/apache/iceberg-rust/pull/3263#discussion_r4088766448
##########
crates/storage/opendal/src/lib.rs:
##########
@@ -163,35 +164,42 @@ where
impl StorageFactory for OpenDalStorageFactory {
#[allow(unused_variables)]
fn build(&self, config: &StorageConfig) -> Result<Arc<dyn Storage>> {
+ let io_timeout_ms = io_timeout_ms_parse(config.props())?;
Review Comment:
Thanks both. Done in 18a0407. `OpenDalClientConfig` uses
`#[derive(Properties)]`, and every `OpenDalStorage` variant carries it as
`client`. Its field is private behind the generated getter, and a075964 adds a
serde round-trip test for it.
One change from the sketch: the field is `io_timeout_ms: u64` rather than
`io_timeout: Duration`. The derive's getter returns by value only for types it
recognizes as `Copy`, which are mostly primitives, so a `Duration` field would
get a getter returning `&Duration`. A `u64` also serializes as a plain integer
in the key's own unit rather than as `{secs, nanos}`. Happy to switch to
`Duration` if you prefer it.
##########
crates/iceberg/src/io/storage/config/mod.rs:
##########
@@ -45,6 +45,12 @@ pub use oss::*;
pub use s3::*;
use serde::{Deserialize, Serialize};
+/// Deadline in milliseconds for one IO operation, and for every method call
on a returned
+/// reader, writer, lister or deleter. Applies to all backends. Defaults to
10000.
Review Comment:
Applied your wording in 18a0407. I also added a short paragraph saying the
budget is per attempt, and that control operations such as `stat` and `rename`
have a separate fixed budget.
##########
crates/storage/opendal/src/lib.rs:
##########
@@ -391,10 +426,35 @@ impl OpenDalStorage {
// Transient errors are common for object stores; we retry temporary
// failures with exponential backoff. The retry behavior also
// benefits non-object-store backends.
- let operator =
operator.layer(TimeoutLayer::new()).layer(RetryLayer::new());
+ let operator = operator
+ .layer(TimeoutLayer::new().with_io_timeout(self.io_timeout()))
+ .layer(RetryLayer::new());
Ok((operator, relative_path))
}
+ /// Per-IO-operation deadline, from
[`CLIENT_IO_TIMEOUT_MS`](iceberg::io::CLIENT_IO_TIMEOUT_MS).
+ #[allow(unreachable_patterns)]
+ fn io_timeout(&self) -> Duration {
+ let ms = match self {
+ #[cfg(feature = "opendal-memory")]
+ OpenDalStorage::Memory { io_timeout_ms, .. } => *io_timeout_ms,
+ #[cfg(feature = "opendal-fs")]
+ OpenDalStorage::LocalFs { io_timeout_ms } => *io_timeout_ms,
+ #[cfg(feature = "opendal-s3")]
+ OpenDalStorage::S3 { io_timeout_ms, .. } => *io_timeout_ms,
+ #[cfg(feature = "opendal-gcs")]
+ OpenDalStorage::Gcs { io_timeout_ms, .. } => *io_timeout_ms,
+ #[cfg(feature = "opendal-oss")]
+ OpenDalStorage::Oss { io_timeout_ms, .. } => *io_timeout_ms,
+ #[cfg(feature = "opendal-azdls")]
+ OpenDalStorage::Azdls { io_timeout_ms, .. } => *io_timeout_ms,
+ #[cfg(feature = "opendal-hf")]
+ OpenDalStorage::Hf { io_timeout_ms, .. } => *io_timeout_ms,
+ _ => default_io_timeout_ms(),
+ };
+ Duration::from_millis(ms)
Review Comment:
Done in 18a0407. The method is now `OpenDalStorage::client()`, which returns
`&OpenDalClientConfig`. Its `_` arm uses the same all-features-off gate as
`create_operator`, with `opendal-memory` added, and the `allow` is gone. The
enum has no variants in that build, so a075964 makes the arm `unreachable!()`.
##########
crates/storage/opendal/src/lib.rs:
##########
@@ -391,10 +429,35 @@ impl OpenDalStorage {
// Transient errors are common for object stores; we retry temporary
// failures with exponential backoff. The retry behavior also
// benefits non-object-store backends.
- let operator =
operator.layer(TimeoutLayer::new()).layer(RetryLayer::new());
+ let operator = operator
+ .layer(TimeoutLayer::new().with_io_timeout(self.io_timeout()))
+ .layer(RetryLayer::new());
Ok((operator, relative_path))
}
+ /// Per-IO-operation deadline, from
[`CLIENT_IO_TIMEOUT_MS`](iceberg::io::CLIENT_IO_TIMEOUT_MS).
+ #[allow(unreachable_patterns)]
+ fn io_timeout(&self) -> Duration {
+ let ms = match self {
+ #[cfg(feature = "opendal-memory")]
+ OpenDalStorage::Memory { io_timeout_ms, .. } => *io_timeout_ms,
+ #[cfg(feature = "opendal-fs")]
+ OpenDalStorage::LocalFs { io_timeout_ms } => *io_timeout_ms,
+ #[cfg(feature = "opendal-s3")]
+ OpenDalStorage::S3 { io_timeout_ms, .. } => *io_timeout_ms,
+ #[cfg(feature = "opendal-gcs")]
+ OpenDalStorage::Gcs { io_timeout_ms, .. } => *io_timeout_ms,
+ #[cfg(feature = "opendal-oss")]
+ OpenDalStorage::Oss { io_timeout_ms, .. } => *io_timeout_ms,
+ #[cfg(feature = "opendal-azdls")]
+ OpenDalStorage::Azdls { io_timeout_ms, .. } => *io_timeout_ms,
+ #[cfg(feature = "opendal-hf")]
+ OpenDalStorage::Hf { io_timeout_ms, .. } => *io_timeout_ms,
+ _ => default_io_timeout_ms(),
Review Comment:
Fixed in 18a0407. The match now lives in `OpenDalStorage::client()`, which
has no `#[allow(unreachable_patterns)]`. Its `_` arm only compiles when every
backend feature is off, so a new variant without an arm is a compile error.
##########
crates/storage/opendal/src/utils.rs:
##########
@@ -15,10 +15,34 @@
// specific language governing permissions and limitations
// under the License.
+use std::collections::HashMap;
+
+use iceberg::io::CLIENT_IO_TIMEOUT_MS;
+
pub(crate) fn is_truthy(value: &str) -> bool {
["true", "t", "1", "on"].contains(&value.to_lowercase().as_str())
}
+/// Matches the `opendal::layers::TimeoutLayer` default.
+pub(crate) fn default_io_timeout_ms() -> u64 {
+ 10_000
+}
Review Comment:
Kept a named `DEFAULT_IO_TIMEOUT_MS`, which is the shape suggested in
review. OpenDAL's `TimeoutLayer` has no getter for its default, so a075964 adds
`test_default_io_timeout_matches_opendal`. It compares the layer's `Debug`
output and fails if an OpenDAL upgrade changes the default.
##########
crates/storage/opendal/src/utils.rs:
##########
@@ -15,10 +15,34 @@
// specific language governing permissions and limitations
// under the License.
+use std::collections::HashMap;
+
+use iceberg::io::CLIENT_IO_TIMEOUT_MS;
+
pub(crate) fn is_truthy(value: &str) -> bool {
["true", "t", "1", "on"].contains(&value.to_lowercase().as_str())
}
+/// Matches the `opendal::layers::TimeoutLayer` default.
+pub(crate) fn default_io_timeout_ms() -> u64 {
+ 10_000
+}
+
+/// Parse iceberg props to the per-IO-operation timeout.
+pub(crate) fn io_timeout_ms_parse(m: &HashMap<String, String>) ->
iceberg::Result<u64> {
+ let Some(value) = m.get(CLIENT_IO_TIMEOUT_MS) else {
+ return Ok(default_io_timeout_ms());
+ };
+ // Zero would make every operation time out before it starts.
+ match value.parse::<u64>() {
+ Ok(ms) if ms > 0 => Ok(ms),
+ _ => Err(iceberg::Error::new(
+ iceberg::ErrorKind::DataInvalid,
+ format!("Invalid {CLIENT_IO_TIMEOUT_MS}: {value}, expected a
positive integer"),
+ )),
Review Comment:
Addressed in a075964. The rejected value now goes into the error context as
`{value:?}`, so an empty input renders as `value: ""`. The error also names the
property key.
##########
crates/storage/opendal/src/utils.rs:
##########
@@ -27,3 +51,23 @@ pub(crate) fn from_opendal_error(e: opendal::Error) ->
iceberg::Error {
)
.with_source(e)
}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+
+ fn props(value: &str) -> HashMap<String, String> {
+ HashMap::from([(CLIENT_IO_TIMEOUT_MS.to_string(), value.to_string())])
+ }
+
+ #[test]
+ fn test_io_timeout_ms_parse() {
+ assert_eq!(io_timeout_ms_parse(&HashMap::new()).unwrap(), 10_000);
Review Comment:
Fixed in a075964. The test now asserts against `DEFAULT_IO_TIMEOUT_MS`, and
`test_default_io_timeout_matches_opendal` pins that constant to OpenDAL's
actual default.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]