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]

Reply via email to