Repository navigation
feat(storage-opendal)!: make the per-IO-operation timeout configurable - #3263
Conversation
43cc854 to
b1751f9
Compare
|
@mbutrovich FYI |
There was a problem hiding this comment.
Warning
Copilot couldn't run its full agentic review because it didn't start before the timeout. Make sure your repository has a runner available, or add a copilot-code-review.yml file specifying one with the runs-on attribute. See the docs for more details.
Copilot review overview
Review effort: Lite
Findings: 3
Open (4)
The wildcard arm (_ => default_io_timeout_ms()) plus#[allow(unreachable_patterns)]makes this… · New Hardcoding10_000as the OpenDALTimeoutLayerdefault is brittle: if OpenDAL changes its… · New For values like the empty string, the error renders asInvalid client.io-timeout-ms: , ..., which… · New The tests duplicate the default value (10_000). Usingdefault_io_timeout_ms()for the expected… · New
What changed in this PR
Adds a new storage configuration property to make OpenDAL’s per-IO-operation timeout configurable, and wires it through the OpenDAL storage factory/resolver into TimeoutLayer.
Changes:
- Introduce
client.io-timeout-ms(CLIENT_IO_TIMEOUT_MS) as a shared storage property iniceberg. - Parse/validate the timeout in
iceberg-storage-opendaland propagate it through factory + resolving storage variants. - Apply the configured timeout to
TimeoutLayer::with_io_timeout, plus add unit tests for parsing and propagation.
| File | Description |
|---|---|
| crates/storage/opendal/src/utils.rs | Adds default + parsing/validation for client.io-timeout-ms and unit tests. |
| crates/storage/opendal/src/resolving.rs | Parses timeout once per build and stores it in each resolved OpenDalStorage variant; adds propagation test. |
| crates/storage/opendal/src/lib.rs | Extends OpenDalStorage variants to carry io_timeout_ms, parses it in StorageFactory::build, and applies it to TimeoutLayer. |
| crates/storage/opendal/public-api.txt | Updates exported API surface to reflect enum variant shape/fields changes. |
| crates/iceberg/src/io/storage/config/mod.rs | Introduces the CLIENT_IO_TIMEOUT_MS property constant with docs. |
| crates/iceberg/public-api.txt | Public API snapshot updated to include CLIENT_IO_TIMEOUT_MS. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
xanderbailey
left a comment
There was a problem hiding this comment.
Thanks for working on this! I do think having a more generic OpenDal config might work better here and we can parse it with the new properties macro. It'll make this more expandable in the future I think. WDYT?
mbutrovich
left a comment
There was a problem hiding this comment.
Thanks @comphead for picking this up, Comet will be glad to have it. My comments are about how the new setting fits the repo's config conventions: how it gets parsed and how it's carried on the OpenDalStorage variants.
0c255ee to
1ef4025
Compare
`iceberg-storage-opendal` wraps every FileIO operator in `TimeoutLayer::new()`, whose 10s `io_timeout` bounds each `read`/`write` and each method call on a returned reader or writer. Nothing in the property or builder surface can override it, so an operation that legitimately needs longer fails deterministically: the `RetryLayer` above re-sends the same request, which cannot fit in the budget either. Add `client.io-timeout-ms`, parsed into an `OpenDalClientConfig` built with `#[derive(Properties)]` and handed to `TimeoutLayer::with_io_timeout`. The config's fields are private, so later client settings such as retry are additive rather than breaking. The default is unchanged. BREAKING CHANGE: `OpenDalStorage` variants carry a `client` field, and `OpenDalStorage::Memory` is now a struct variant.
Quote the rejected value in the parse error, so an empty input shows as `value: ""`. `TimeoutLayer` exposes no getters, so a new test compares its `Debug` output to pin `DEFAULT_IO_TIMEOUT_MS` to OpenDAL's default. The parsing test now uses the constant instead of repeating the literal. Add a serde round-trip test for the `client` field, including a payload that omits it. Drop the `DEFAULT` associated const, which only fed a fallback arm that cannot run. Update the upstream `test_writer_close_returns_stored_size` for the struct variant, and reuse `StorageConfig::with_prop` and `empty_resolving_storage` in the tests.
1ef4025 to
a075964
Compare
…meout-ms Only `iceberg-storage-opendal` honors the key, so namespace it under `opendal.` and define it as `OPENDAL_IO_TIMEOUT_MS` next to `OpenDalClientConfig`, rather than in `iceberg::io`. This matches how `ADLS_SAS_TOKEN` and `GCS_ALLOW_ANONYMOUS` sit next to the configs that parse them. With the key moved, the `iceberg` crate needs no change.
|
Depends on #3263 |
|
@kevinjqliu cc appreciate if you can take a look, in Comet we got file scan tasks restarted because of uncontrollable timeout which makes the entire execution time longer. the PR is to expose timeout value to external users |
There was a problem hiding this comment.
Warning
Copilot couldn't run its full agentic review because it didn't start before the timeout. Make sure your repository has a runner available, or add a copilot-code-review.yml file specifying one with the runs-on attribute. See the docs for more details.
Copilot review overview
Review effort: Lite
Findings: 1
Open (3)
Resolved since last review (4)
For values like the empty string, the error renders asInvalid client.io-timeout-ms: , ..., which… Hardcoding10_000as the OpenDALTimeoutLayerdefault is brittle: if OpenDAL changes its… The wildcard arm (_ => default_io_timeout_ms()) plus#[allow(unreachable_patterns)]makes this… The tests duplicate the default value (10_000). Usingdefault_io_timeout_ms()for the expected…
| // `TimeoutLayer` has no getters, so compare through `Debug`. An OpenDAL upgrade that | ||
| // changes its default fails here instead of silently diverging from it. | ||
| assert_eq!( | ||
| format!("{:?}", TimeoutLayer::new()), | ||
| format!( | ||
| "{:?}", | ||
| TimeoutLayer::new().with_io_timeout(Duration::from_millis(DEFAULT_IO_TIMEOUT_MS)) | ||
| ), | ||
| ); |
There was a problem hiding this comment.
Kept it, since TimeoutLayer has no getters and both sides go through the same derived Debug, so a format change alone cannot make them differ. The assert_ne! from 5e9641c catches a Debug that stops printing io_timeout.
|
|
||
| /// Deadline in milliseconds for one IO operation, and for every method call on a returned | ||
| /// reader, writer, lister or deleter. Honored by every [`OpenDalStorage`] backend, where it | ||
| /// defaults to 10000 to match OpenDAL's `TimeoutLayer`. |
There was a problem hiding this comment.
Fixed. The doc now links to the new OPENDAL_IO_TIMEOUT_MS_DEFAULT instead of repeating the number.
`client` reads like an HTTP client handle. Rename the field on every `OpenDalStorage` variant, and the `OpenDalStorage::client()` accessor, to `client_config` to match `OpenDalClientConfig`. `opendal_config` would be ambiguous next to the backend `config` field, which holds an OpenDAL config too. `test_default_io_timeout_matches_opendal` compares `Debug` output. Both sides go through the same derived impl, so a format change cannot make them differ. A `Debug` that stopped printing `io_timeout` would make the check pass vacuously, so also assert that a different `io_timeout` produces a different string.
`OpenDalClientConfig::io_timeout_ms()` returned `&NonZeroU64`, because the `Properties` derive only returns primitives by value. Replace the generated getter with `io_timeout()`, which returns a `Duration`, and keep the millisecond field private. Expose the default as `OPENDAL_IO_TIMEOUT_MS_DEFAULT`, following the `*_DEFAULT` consts in `iceberg`, so the key's doc links to it instead of repeating the number. Keep the `ParseIntError` as the source of a rejected value, as `parse_pool_property` does, so the message says why it failed. Document that the serialized form of `OpenDalStorage` is not stable across crate versions, as `FileIO::serialize_all` already does for its own. Test that the configured timeout reaches `TimeoutLayer`. A memory operator behind `ConcurrentLimitLayer::new(0)` stalls every IO call, and paused tokio time skips the wait and the retry backoff. Check propagation through `OpenDalResolvingStorage` for every enabled scheme, not only S3.
|
@laskoviymishka The description now lists every variant that gained @xanderbailey Agreed. |
…erive Build `OpenDalClientConfig::default()` from `from_properties` with no properties, so each default is declared once in `#[property]` and serde and `from_properties` cannot disagree as settings are added. Say in the `parse_io_timeout_ms` doc that it exists to quote the value and name the unit, and that `NonZeroU64` is what rejects zero. Fold the valid-value asserts in the parsing test into one loop.
|
Thanks @laskoviymishka and @xanderbailey for the review, PTAL |
|
taking a look! sorry github ui is glitching for me |
laskoviymishka
left a comment
There was a problem hiding this comment.
Thanks for the quick turnaround. The round-1 blocker is resolved now: the serialized-form break is called out separately, the new client_config fields are listed, and the timeout getter / parse error / new tests cover the gaps from last round.
One thing I’d still fix before merge: the title says datafusion, but nothing in the diff touches DataFusion. I’d drop that scope and use feat(storage-opendal)!: … instead.
The rest is minor and inline: the .expect() in Default is a future panic if a field ever has no default =, the control-op budget note could mention the 60s value, and there are a couple of small test nits.
Fix the title and I’m good with this.
| /// defaults to [`OPENDAL_IO_TIMEOUT_MS_DEFAULT`]. | ||
| /// | ||
| /// Each retry attempt is bounded separately, so it is a per-attempt budget, not a total one. | ||
| /// Control operations such as `stat` and `rename` are bounded by a separate, fixed budget. |
There was a problem hiding this comment.
While we're documenting this, I'd give the "separate, fixed budget" line a concrete number — it's OpenDAL's 60s control-op default, and it also governs delete and the lister/deleter open, not just stat/rename. Otherwise someone who sets io-timeout-ms=5000 to bound all their I/O gets surprised when metadata fetches and snapshot-expiry deletes still hang for up to a minute.
There was a problem hiding this comment.
Added the 60s, pinned to OpenDAL's default by test_default_timeouts_match_opendal, and the doc now names exists and metadata as the calls that keep it. Deletes and listing are not control operations in OpenDAL 0.58, which wraps deleters and listers in the IO timeout, and a stalled delete with a 45000 ms setting fails with timeout: 45.
| impl Default for OpenDalClientConfig { | ||
| fn default() -> Self { | ||
| // Reuse the `#[property]` defaults, so serde and `from_properties` cannot disagree. | ||
| Self::from_properties(&HashMap::new()).expect("every client setting has a default") |
There was a problem hiding this comment.
I'd build this default directly (Self { io_timeout_ms: DEFAULT_IO_TIMEOUT_MS }) rather than route through from_properties().expect(). It's safe today, but the moment someone adds a #[property] field without a default =, this becomes a production panic on the deserialization path — serde calls default() for every absent client_config, so a stored payload missing that field would panic instead of erroring. If you want to keep the single-source guarantee, a #[test] asserting default() == from_properties(&empty) turns the drift into a CI failure instead.
There was a problem hiding this comment.
Agreed, and the exposure is wider than absent payloads, since the container #[serde(default)] makes serde call default() on every client_config it deserializes. Reverted to the struct literal and added test_default_matches_property_defaults, which compares default() with from_properties on an empty map.
|
|
||
| #[cfg(feature = "opendal-s3")] | ||
| #[test] | ||
| fn test_client_config_serde_round_trip() { |
There was a problem hiding this comment.
The "Breaking (serialized form)" section covers the wire break well and I'm happy with the docs-only path we settled on last round — not reopening that. One cheap addition while you're in here: a test asserting an old bare-string "LocalFs" / {"Memory": null} payload now is_err() would pin the break as intended behavior rather than an accident, sitting right next to this round-trip test. Not blocking.
There was a problem hiding this comment.
Added test_old_unit_variant_forms_are_rejected next to the round trip. It checks that "LocalFs", "Memory" and {"Memory": null} all fail to deserialize.
| .unwrap_err() | ||
| .to_string(); | ||
| assert!(err.contains("io operation timeout reached"), "{err}"); | ||
| assert!(err.contains("timeout: 45"), "{err}"); |
There was a problem hiding this comment.
"timeout: 45" also matches "timeout: 450", so this passes if the layer ever reports 450s — "timeout: 45s" with the unit pins it to the value you actually mean.
There was a problem hiding this comment.
OpenDAL renders the budget as plain seconds (timeout.as_secs_f64().to_string()), so the error reads context: { timeout: 45 } and "timeout: 45s" would never match. The assertion now matches { timeout: 45 }, whose closing brace rules out 450.
|
Thanks @laskoviymishka for the review, addressing |
`Default` delegated to `from_properties` with no properties and called
`expect`. A later setting without a `#[property]` default would make it
panic, and because of the container-level `#[serde(default)]`, serde
calls `default()` on every `client_config` it deserializes. Build the
struct directly, and keep the two defaults equal with
`test_default_matches_property_defaults`.
Name OpenDAL's 60-second control-operation budget in the
`OPENDAL_IO_TIMEOUT_MS` doc, along with the calls it covers here,
`exists` and `metadata`. `test_default_timeouts_match_opendal` now pins
both defaults that the doc states.
Pin the serialized-form break with a test that rejects the old
`"LocalFs"`, `"Memory"` and `{"Memory": null}` forms. Match the stall
test's error on `{ timeout: 45 }`, so a 450-second budget cannot pass.
|
@laskoviymishka Thanks. The title now drops the |
|
Will wait for some time before merging (looks like @kevinjqliu still looking) |
…vate `OpenDalStorage::client_config()`, `OpenDalClientConfig::io_timeout()` and `OPENDAL_IO_TIMEOUT_MS_DEFAULT` have no callers outside the crate. Keep them out of the public API so they can be exposed later without a breaking change. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
Fold the `u64` default into the `NonZeroU64` const, drop `test_default_matches_property_defaults` (one field, one shared const), shorten docs and comments that restated the code, and keep the original `op` binding in the `Memory` arm of `create_operator`. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
Name OpenDAL's `io_timeout` and `timeout` settings so it is clear the property bounds each IO call per retry attempt and leaves control operations such as `stat` at 60 seconds. Align the `OpenDalStorage` serialization note with the wording on `FileIO::serialize_all`. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
…timeout Every `OpenDalStorageFactory::build` arm could fall back to the default `OpenDalClientConfig` without any test failing. Build each enabled backend with a 45000 ms timeout and read it back from the serialized storage. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
kevinjqliu
left a comment
There was a problem hiding this comment.
LGTM! Thanks for adding this (and for your patience waiting for reviews).
I only have a few nit comments. Trying out a new workflow where i add some of the nits and comments as commits on top of the PR. Hope you don't mind, and please lmk if you dont like this flow, I can add regular comments to the PR.
I added a few changes:
- 25612d5 makes
client_config(),io_timeout()andOPENDAL_IO_TIMEOUT_MS_DEFAULTcrate-private, since nothing in the repo uses them. Happy to revert and expose them again if Comet or other downstream users need them. - 3b169cd trims redundant consts, comments and a test to keep the diff tight.
- 38af8ab clarifies
io_timeoutvstimeoutin the docs, and aligns the serialization note withFileIO::serialize_all. - b56cbb4 adds a test that
OpenDalStorageFactorypropagates the timeout. Nothing covered it before.
The description still lists those three items as public, so it'll need a small update.
PTAL! And feel free to revert any you disagree with.
|
Thanks @kevinjqliu appreciate if we can merge the PR |
|
@comphead if you’re happy with the changes, lmk and I can press merge! |
|
Thanks for the pr! I’m excited to see more opendal integration. And thanks everyone for the reviews! 🚢 |
Bump iceberg-rust to af1da4c (apache/iceberg-rust#3263), which adds the `opendal.io-timeout-ms` FileIO property, and set it from the new `spark.comet.iceberg.ioTimeout` config (default 30s) for native Iceberg scans and writes. Forward `opendal.*` keys to the native FileIO. The bump also needs two API updates: `FileWrite::close` now returns `FileMetadata`, and `UnboundPartitionField` is built with its builder. Closes apache#6124.
Resolve conflicts with the configurable OpenDAL IO timeout (apache#3263): storage variants carry both the credential provider and the client config.


Which issue does this PR close?
What changes are included in this PR?
iceberg-storage-opendalwraps every FileIO operator inTimeoutLayer::new(). Its 10sio_timeoutbounds eachread/writeand every method call on a returned reader, writer, lister or deleter, and nothing in the property or builder surface can override it:https://github.com/apache/iceberg-rust/blob/bb1e4a4/crates/storage/opendal/src/lib.rs#L394
An operation that legitimately needs longer then fails deterministically, not flakily:
RetryLayerre-sends the same request, which cannot fit in the budget either, so every attempt dies at the same place and the error surfaces as persistent.Add
opendal.io-timeout-msand hand it toTimeoutLayer::with_io_timeout. Only this crate honors it, so the key is namespaced underopendal.and defined here asOPENDAL_IO_TIMEOUT_MS, not iniceberg::io. One property covers every backend rather than one per service. Unset keeps OpenDAL's 10s, exposed asOPENDAL_IO_TIMEOUT_MS_DEFAULT. Zero and non-numeric values are rejected rather than silently ignored, and the error keeps the parse failure as its source, so it says why.OpenDalStorageFactoryrejects them atbuildtime.OpenDalResolvingStoragerejects them when it first resolves a scheme, which is also when it parses the backend properties. A serialized storage with a zero timeout fails to deserialize.TimeoutLayerstays insideRetryLayer, so each attempt is still independently bounded.The property is parsed into a new
OpenDalClientConfig, built with#[derive(Properties)]. EveryOpenDalStoragevariant carries it as aclient_configfield, andOpenDalStorage::client_config()returns it. Its field is private, so later client settings such as retry can be added without changing the enum again.OpenDalClientConfig::io_timeout()returns the timeout as aDuration. The field itself is aNonZeroU64in milliseconds, so parsing and deserializing share the zero check.Breaking (API): every
OpenDalStoragevariant gains aclient_config: OpenDalClientConfigfield:Memory,LocalFs,S3,Gcs,Oss,Azdls, andHf.Memory(Operator)becomesMemory { operator, client_config }, andLocalFsbecomesLocalFs { client_config }. Code that builds any variant, or matches one without.., needs updating. PassingOpenDalClientConfig::default()keeps the current behavior.Breaking (serialized form): this affects code that serializes an
OpenDalStorageitself, for example as adyn Storagethroughtypetag.FileIO::serialize_allis not affected, because it serializes the factory and properties, and neither factory changes.LocalFsandMemoryused to serialize as the bare strings"LocalFs"and"Memory", and now serialize as{"LocalFs": {"client_config": {...}}}and{"Memory": {"client_config": {...}}}. Neither version can read the other's form, so aLocalFsorMemorystorage serialized by an earlier version must be rebuilt rather than deserialized. In JSON,S3,Gcs,Oss,Azdls, andHfstay compatible in both directions, because a missingclient_configfalls back to the default and older versions ignore it as an unknown field. TheOpenDalStoragedocs now say that its serialized form is not stable across crate versions, asFileIO::serialize_allalready says of its own.#3179 bounds the S3 write request size, which removes the oversized-part case on the write path. This covers the general one, including reads. Scaling the deadline with payload size (option 2 in #2977) is not available: OpenDAL dropped
TimeoutLayer::with_speedin apache/opendal#6793.The 60s
timeoutfor control operations is left alone, as it has not been reported as a problem. The property docs now name it, and that here it coversexistsandmetadata, while reads, writes, listing and deletes use the new property.Are these changes tested?
Unit tests:
test_io_timeout_parsing: default when unset, override,1andu64::MAXas the smallest and largest accepted values, and rejection of0,-1,12.5,abc, the empty string, andu64::MAX + 1. The error names the property, quotes the rejected value, and carries the parse error that says why it was rejected.test_default_matches_property_defaults:OpenDalClientConfig::default(), which serde uses, matches the#[property]defaults thatfrom_propertiesuses.test_default_timeouts_match_opendal:OPENDAL_IO_TIMEOUT_MS_DEFAULTand the 60s control budget named in the docs stay equal to OpenDAL'sTimeoutLayerdefaults, so an OpenDAL upgrade that changes either fails here.test_factory_rejects_invalid_io_timeout:OpenDalStorageFactory::buildsurfaces a parse failure instead of falling back.test_io_timeout_reaches_timeout_layer: the configured value reachesTimeoutLayer::with_io_timeout. A memory operator behindConcurrentLimitLayer::new(0)stalls every IO call, and a read with a 45000 ms setting fails withtimeout: 45. Paused tokio time, from the dev-onlytest-utilfeature, skips the wait and the retry backoff. With a bareTimeoutLayer::new()the test fails withtimeout: 10.test_resolve_propagates_io_timeout:OpenDalResolvingStorage::resolvepropagates the property for every enabled scheme, including in a memory-only build.test_client_config_serde_round_trip:client_configsurvives a serde round trip, a payload with a zero timeout fails to deserialize, and a payload withoutclient_configfalls back to the default.test_old_unit_variant_forms_are_rejected: the old"LocalFs","Memory"and{"Memory": null}forms fail to deserialize, which pins the serialized-form break as intended.AI Disclosure