diff --git a/docs/source/contributor-guide/s3-credential-provider-design.md b/docs/source/contributor-guide/s3-credential-provider-design.md index d2f00c0c41e..d4fa3e1baa9 100644 --- a/docs/source/contributor-guide/s3-credential-provider-design.md +++ b/docs/source/contributor-guide/s3-credential-provider-design.md @@ -142,7 +142,7 @@ The 5-minute fallback on the Iceberg path bounds how long reqsign reuses a crede ## Property-bag handling on the Iceberg path -The full unfiltered FileIO property bag crosses JNI as `catalog_properties`. The storage-prefix filter (`s3.`/`gcs.`/`adls.`/`client.`) is applied native-side in `iceberg_common.rs::build_file_io` immediately before `FileIOBuilder.with_prop`. This means the bridge sees `credentials.uri`, OAuth tokens, and any vendor-custom keys with no parallel field on the operator and no driver-side broadcast. Vendors set their own keys on the catalog config and read them back inside `initialize(Map)`. +The full unfiltered FileIO property bag crosses JNI as `catalog_properties`. The storage-prefix filter (`s3.`/`gcs.`/`adls.`/`client.`/`opendal.`) is applied native-side in `iceberg_common.rs::build_file_io` immediately before `FileIOBuilder.with_prop`. This means the bridge sees `credentials.uri`, OAuth tokens, and any vendor-custom keys with no parallel field on the operator and no driver-side broadcast. Vendors set their own keys on the catalog config and read them back inside `initialize(Map)`. `IcebergScanExec` derives a redacting `Debug`, and the `FileIO` cache key's `Debug` omits the property bag, so plan dumps and tracing do not leak it. diff --git a/native/Cargo.lock b/native/Cargo.lock index 07011be0d19..62edc47804b 100644 --- a/native/Cargo.lock +++ b/native/Cargo.lock @@ -2895,12 +2895,6 @@ dependencies = [ "syn 3.0.5", ] -[[package]] -name = "dissimilar" -version = "1.0.11" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "aeda16ab4059c5fd2a83f2b9c9e9c981327b18aa8e3b313f7e6563799d4f093e" - [[package]] name = "dlv-list" version = "0.5.2" @@ -3001,16 +2995,6 @@ dependencies = [ "pin-project-lite", ] -[[package]] -name = "expect-test" -version = "1.5.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "63af43ff4431e848fb47472a920f14fa71c24de13255a5692e93d4e90302acb0" -dependencies = [ - "dissimilar", - "once_cell", -] - [[package]] name = "fastnum" version = "0.7.5" @@ -3577,7 +3561,7 @@ dependencies = [ [[package]] name = "iceberg" version = "0.10.1" -source = "git+https://github.com/apache/iceberg-rust?rev=bb1e4a4861f02377489eff818b75138f414c4cb0#bb1e4a4861f02377489eff818b75138f414c4cb0" +source = "git+https://github.com/apache/iceberg-rust?rev=af1da4c5bc86178c38c1c6db0578bdd9fb7021e3#af1da4c5bc86178c38c1c6db0578bdd9fb7021e3" dependencies = [ "aes-gcm", "anyhow", @@ -3600,7 +3584,6 @@ dependencies = [ "chrono", "crc32fast", "derive_builder", - "expect-test", "fastnum", "flate2", "fnv", @@ -3636,7 +3619,7 @@ dependencies = [ [[package]] name = "iceberg-property-macro" version = "0.10.1" -source = "git+https://github.com/apache/iceberg-rust?rev=bb1e4a4861f02377489eff818b75138f414c4cb0#bb1e4a4861f02377489eff818b75138f414c4cb0" +source = "git+https://github.com/apache/iceberg-rust?rev=af1da4c5bc86178c38c1c6db0578bdd9fb7021e3#af1da4c5bc86178c38c1c6db0578bdd9fb7021e3" dependencies = [ "proc-macro2", "quote", @@ -3646,7 +3629,7 @@ dependencies = [ [[package]] name = "iceberg-storage-opendal" version = "0.10.1" -source = "git+https://github.com/apache/iceberg-rust?rev=bb1e4a4861f02377489eff818b75138f414c4cb0#bb1e4a4861f02377489eff818b75138f414c4cb0" +source = "git+https://github.com/apache/iceberg-rust?rev=af1da4c5bc86178c38c1c6db0578bdd9fb7021e3#af1da4c5bc86178c38c1c6db0578bdd9fb7021e3" dependencies = [ "anyhow", "async-trait", @@ -3654,6 +3637,7 @@ dependencies = [ "cfg-if", "futures", "iceberg", + "iceberg-property-macro", "opendal", "reqsign-aws-v4", "reqsign-core", diff --git a/native/Cargo.toml b/native/Cargo.toml index 3b144297657..fa6ea26a217 100644 --- a/native/Cargo.toml +++ b/native/Cargo.toml @@ -80,8 +80,8 @@ aws-sdk-sts = { version = "1.114.0", default-features = false } # Pinned to a revision rather than a release because no published iceberg-rust # version is on our arrow: 0.10.1 requires arrow/parquet ^58. Move to a plain # version requirement once a release ships on arrow 59. Bump policy: #5645. -iceberg = { git = "https://github.com/apache/iceberg-rust", rev = "bb1e4a4861f02377489eff818b75138f414c4cb0" } -iceberg-storage-opendal = { git = "https://github.com/apache/iceberg-rust", rev = "bb1e4a4861f02377489eff818b75138f414c4cb0", features = ["opendal-memory", "opendal-fs", "opendal-s3", "opendal-gcs", "opendal-oss", "opendal-azdls"] } +iceberg = { git = "https://github.com/apache/iceberg-rust", rev = "af1da4c5bc86178c38c1c6db0578bdd9fb7021e3" } +iceberg-storage-opendal = { git = "https://github.com/apache/iceberg-rust", rev = "af1da4c5bc86178c38c1c6db0578bdd9fb7021e3", features = ["opendal-memory", "opendal-fs", "opendal-s3", "opendal-gcs", "opendal-oss", "opendal-azdls"] } reqsign-core = "3" [profile.release] diff --git a/native/core/src/execution/operators/iceberg_common.rs b/native/core/src/execution/operators/iceberg_common.rs index 94df0bd70e5..f4ea3f7983b 100644 --- a/native/core/src/execution/operators/iceberg_common.rs +++ b/native/core/src/execution/operators/iceberg_common.rs @@ -44,8 +44,9 @@ const ICEBERG_PROVIDER_CLASS_PROPERTY: &str = "s3.comet.credential.provider.clas /// Key prefixes forwarded to iceberg-rust's `FileIO`. The full unfiltered catalog bag (catalog /// URI, OAuth tokens, credentials.uri, tenant-id, etc.) is kept upstream so -/// `CometS3CredentialBridge` can read whatever the vendor needs. -const STORAGE_PROPERTY_PREFIXES: &[&str] = &["s3.", "gcs.", "adls.", "client."]; +/// `CometS3CredentialBridge` can read whatever the vendor needs. `opendal.` carries +/// iceberg-storage-opendal's own settings, such as `opendal.io-timeout-ms`. +const STORAGE_PROPERTY_PREFIXES: &[&str] = &["s3.", "gcs.", "adls.", "client.", "opendal."]; /// Pick an OpenDAL storage backend from a URI's scheme. `file` (or no scheme) falls through to /// the local file system. `memory` is used by the write path to assemble manifest bytes that @@ -775,6 +776,17 @@ mod tests { .all(|k| scheme_of(&k.reference_path) != "memory")); } + /// The JVM sets this key from `spark.comet.iceberg.ioTimeout`. + #[test] + fn io_timeout_reaches_the_file_io() { + let key = iceberg_storage_opendal::OPENDAL_IO_TIMEOUT_MS; + assert_eq!(key, "opendal.io-timeout-ms"); + let props = HashMap::from([(key.to_string(), "30000".to_string())]); + let (file_io, _) = + build_file_io(&props, "file:///tmp/warehouse", "", AccessMode::Read).unwrap(); + assert_eq!(file_io.config().get(key).map(String::as_str), Some("30000")); + } + #[test] fn cache_evicts_the_least_recently_used_entry() { let mut cache = FileIoCache::new(2); diff --git a/native/core/src/execution/operators/iceberg_location_scoped.rs b/native/core/src/execution/operators/iceberg_location_scoped.rs index eb7fd937991..d1ad4744859 100644 --- a/native/core/src/execution/operators/iceberg_location_scoped.rs +++ b/native/core/src/execution/operators/iceberg_location_scoped.rs @@ -654,12 +654,12 @@ impl LocationScopedFileWrite { /// `result` of a write call made at the snapshot of `generation`, after fetching `bucket`'s /// locations again if it failed in a way that may mean they changed. A free function so the /// writer, which is not `Sync`, is not borrowed across the refresh. -async fn refresh_on_failure( +async fn refresh_on_failure( inner: &Inner, bucket: &BucketLocations, generation: u64, - result: Result<()>, -) -> Result<()> { + result: Result, +) -> Result { match result { Err(e) if may_mean_stale_locations(&e) => { match inner.shared.refresh(bucket, generation).await { @@ -679,7 +679,7 @@ impl FileWrite for LocationScopedFileWrite { refresh_on_failure(&self.inner, &self.route.bucket, generation, result).await } - async fn close(&mut self) -> Result<()> { + async fn close(&mut self) -> Result { let generation = self.generation(); let result = self.writer.close().await; refresh_on_failure(&self.inner, &self.route.bucket, generation, result).await @@ -797,8 +797,10 @@ mod tests { Ok(()) } - async fn close(&mut self) -> Result<()> { - self.0.check("close", &self.1) + async fn close(&mut self) -> Result { + self.0 + .check("close", &self.1) + .map(|_| FileMetadata { size: 0 }) } } @@ -1512,7 +1514,7 @@ mod tests { let output = storage.new_output("s3://b/t/new/f.parquet").unwrap(); let mut writer = output.writer().await.unwrap(); writer.write(Bytes::from("data")).await.unwrap(); - let err = writer.close().await.unwrap_err(); + let err = writer.close().await.err().expect("close should fail"); assert!(may_mean_stale_locations(&err), "{err}"); assert_eq!(fixture.fetches(), 1); let mut writer = output.writer().await.unwrap(); diff --git a/native/core/src/execution/planner.rs b/native/core/src/execution/planner.rs index aec252e2e8d..c8dc8fa46a0 100644 --- a/native/core/src/execution/planner.rs +++ b/native/core/src/execution/planner.rs @@ -7298,15 +7298,20 @@ mod tests { .expect("schema"); let schema_arc = Arc::new(iceberg_schema); + let identity_field = |source_id: i32, field_id: i32, name: &str| { + iceberg::spec::UnboundPartitionField::builder() + .source_ids(vec![source_id]) + .field_id(field_id) + .name(name) + .transform(iceberg::spec::Transform::Identity) + .build() + .expect("unbound field") + }; + // Spec 0 (older): field 1000 named "region_old". let spec0 = PartitionSpec::builder(Arc::clone(&schema_arc)) .with_spec_id(0) - .add_unbound_field(iceberg::spec::UnboundPartitionField { - source_id: 2, - field_id: Some(1000), - name: "region_old".to_string(), - transform: iceberg::spec::Transform::Identity, - }) + .add_unbound_field(identity_field(2, 1000, "region_old")) .expect("add field") .build() .expect("build spec0"); @@ -7314,19 +7319,9 @@ mod tests { // Spec 1 (newer): same field id 1000 renamed to "region_new", plus a new field 2000. let spec1 = PartitionSpec::builder(Arc::clone(&schema_arc)) .with_spec_id(1) - .add_unbound_field(iceberg::spec::UnboundPartitionField { - source_id: 2, - field_id: Some(1000), - name: "region_new".to_string(), - transform: iceberg::spec::Transform::Identity, - }) + .add_unbound_field(identity_field(2, 1000, "region_new")) .expect("add field") - .add_unbound_field(iceberg::spec::UnboundPartitionField { - source_id: 3, - field_id: Some(2000), - name: "category".to_string(), - transform: iceberg::spec::Transform::Identity, - }) + .add_unbound_field(identity_field(3, 2000, "category")) .expect("add field") .build() .expect("build spec1"); diff --git a/spark/src/main/scala/org/apache/comet/CometConf.scala b/spark/src/main/scala/org/apache/comet/CometConf.scala index 5c588017edc..5465de1f64f 100644 --- a/spark/src/main/scala/org/apache/comet/CometConf.scala +++ b/spark/src/main/scala/org/apache/comet/CometConf.scala @@ -152,6 +152,15 @@ object CometConf extends ShimCometConf { .checkValue(v => v > 0, "Data file concurrency limit must be positive") .createWithDefault(1) + val COMET_ICEBERG_IO_TIMEOUT: ConfigEntry[Long] = + conf("spark.comet.iceberg.ioTimeout") + .category(CATEGORY_SCAN) + .doc("Timeout for each storage I/O call, such as one read or write, that Comet's native " + + "Iceberg scan and write make. Passed to iceberg-rust as `opendal.io-timeout-ms`.") + .timeConf(TimeUnit.MILLISECONDS) + .checkValue(_ > 0, "The Iceberg I/O timeout must be positive") + .createWithDefault(TimeUnit.SECONDS.toMillis(30)) + val COMET_CSV_V2_NATIVE_ENABLED: ConfigEntry[Boolean] = conf("spark.comet.scan.csv.v2.enabled") .category(CATEGORY_TESTING) diff --git a/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeScan.scala b/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeScan.scala index 689ed68389e..d7865833776 100644 --- a/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeScan.scala +++ b/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeScan.scala @@ -645,6 +645,10 @@ object CometIcebergNativeScan extends CometOperatorSerde[CometBatchScanExec] wit global.result() ++ promoted.result() } + /** iceberg-rust's per-I/O-call timeout, set from `spark.comet.iceberg.ioTimeout`. */ + def ioTimeoutProperty(): (String, String) = + "opendal.io-timeout-ms" -> CometConf.COMET_ICEBERG_IO_TIMEOUT.get().toString + /** * Converts an Iceberg residual Expression into an IcebergPredicate for the native scan. * @@ -1020,7 +1024,7 @@ object CometIcebergNativeScan extends CometOperatorSerde[CometBatchScanExec] wit commonBuilder.setDataFileConcurrencyLimit( CometConf.COMET_ICEBERG_DATA_FILE_CONCURRENCY_LIMIT.get()) metadata.catalogName.foreach(commonBuilder.setCatalogName) - metadata.catalogProperties.foreach { case (key, value) => + (metadata.catalogProperties + ioTimeoutProperty()).foreach { case (key, value) => commonBuilder.putCatalogProperties(key, value) } diff --git a/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala b/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala index 2c5d00f2b50..5d5510fb7a1 100644 --- a/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala +++ b/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala @@ -766,7 +766,8 @@ object CometIcebergNativeWrite extends CometOperatorSerde[IcebergWriteExec] { val hadoopDerivedProperties = CometIcebergNativeScan.hadoopToIcebergS3Properties( NativeConfig.extractObjectStoreOptions(writeHadoopConf, dataUri), dataBucket) - val catalogProperties = hadoopDerivedProperties ++ fileIOProperties + val catalogProperties = + hadoopDerivedProperties ++ fileIOProperties + CometIcebergNativeScan.ioTimeoutProperty() val common = IcebergWriteProtoTranslation.buildCommon( catalogProperties = catalogProperties,