comphead opened a new issue, #6124:
URL: https://github.com/apache/datafusion-comet/issues/6124
### Describe the bug
A native Iceberg scan fails with a hard-wired 10 second per-IO timeout that
nothing in Comet can configure:
```
org.apache.comet.CometNativeException: Iceberg scan error: Unexpected =>
Failed to read a Parquet file,
source: External: Unexpected => Failure in doing io operation,
source: Unexpected (persistent) at read, context: { timeout: 10 } => io
operation timeout reached
at org.apache.comet.Native.executePlan(Native Method)
at
org.apache.comet.CometExecIterator.$anonfun$getNextBatch$2(CometExecIterator.scala:234)
...
at
org.apache.spark.sql.comet.execution.shuffle.CometNativeShuffleWriter.write(CometNativeShuffleWriter.scala:95)
```
**Where the 10 comes from.** `iceberg-storage-opendal` wraps every `FileIO`
operator in `TimeoutLayer::new()`, whose default `io_timeout` is 10 seconds:
https://github.com/apache/iceberg-rust/blob/665c64e48e8d33797ecb1a421f327edd9b024879/crates/storage/opendal/src/lib.rs#L394
That budget bounds a single `ReadStream::read()` call, meaning one chunk of
the response body, not the whole file. So this is not a large-file problem.
**Why it surfaces as persistent rather than flaky.** `RetryLayer::new()`
sits outside the timeout layer with `backon`'s defaults (`max_times: 3`,
`min_delay: 1s`, `factor: 2.0`). Four attempts with 1s, 2s and 4s backoff gives
roughly a 47 second budget per stalled byte range, after which the error is
re-marked persistent. Every attempt re-sends the same request, so a range that
cannot fit in 10 seconds fails identically each time.
**Why Comet cannot fix it locally.** Every hook is `pub(crate)` in the
pinned dependency:
- `OpenDalStorage::create_operator`, where the layer is applied, is
`pub(crate)`.
- `s3_config_parse` / `s3_config_build` are `pub(crate)`, so a Comet-owned
`StorageFactory` cannot reuse the config parsing.
- `iceberg::io::Storage` is opaque (`read` / `reader` / `write` / ...), so
`BlobHostPromotingS3StorageFactory` in
`native/core/src/parquet/objectstore/s3_blob_fs_support.rs` receives an
already-layered `Arc<dyn Storage>` and cannot re-layer. Layers stack rather
than replace, so wrapping a second, longer `TimeoutLayer` on the outside does
not help either.
### Steps to reproduce
Run a native Iceberg scan over object storage where one byte range takes
longer than 10 seconds to deliver a chunk. In the report above, the scan was
feeding a native shuffle write (`CometNativeShuffleWriter.drainAndClose`).
### Expected behavior
The per-IO timeout should be tunable, so a deployment on a slow or contended
link can raise it instead of failing the job.
### Additional context
**Upstream.** Tracked as
[apache/iceberg-rust#2977](https://github.com/apache/iceberg-rust/issues/2977).
[apache/iceberg-rust#3179](https://github.com/apache/iceberg-rust/pull/3179) is
open against that issue but only bounds S3 *write* request size, so it does not
cover the read path above.
**Fix is proposed upstream:**
https://github.com/apache/iceberg-rust/pull/3263
It adds `client.io-timeout-ms` in the existing `client.*` namespace (next to
`client.region`) and hands it to `TimeoutLayer::with_io_timeout`. One property
covers every backend, and unset keeps OpenDAL's 10s so the default is unchanged.
**What lands in Comet once that merges.** Very little, because `client.` is
already forwarded:
- Bump the pinned rev at `native/Cargo.toml:67-68`.
- No native change needed. `STORAGE_PROPERTY_PREFIXES` at
`native/core/src/execution/operators/iceberg_common.rs:39` already includes
`"client."`, and `load_file_io` forwards matching keys to
`FileIOBuilder::with_prop`.
- Add one `CometConf` entry (`.timeConf(TimeUnit.MILLISECONDS)`, mirroring
`COMET_SHUFFLE_JVM_MEMORY_WAIT_TIMEOUT`) and inject it into the property bag at
the two assembly sites: `CometScanRule.scala:590` (scan) and
`CometIcebergNativeWrite.scala:687` (write).
Both sites read from `SQLConf` at planning time, so unlike the `fs.comet.*`
Hadoop keys this config would honor a runtime `spark.conf.set`.
**One thing worth ruling out first.** `TimeoutLayer` uses wall-clock
`tokio::time::timeout`, so it fires whether the stall is in the network or in
the scheduler. Comet runs one shared multi-threaded tokio runtime per executor
sized to `spark.executor.cores`
(`native/core/src/execution/jni_api.rs:288-296`). With CPU-bound shuffle
partitioning and compression competing for the same workers, an IO future on a
healthy connection can go unpolled past 10 seconds. I have not verified which
applies to the report above. Raising `COMET_WORKER_THREADS` is the cheaper
experiment and needs no code change.
--
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]