Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
24 changes: 4 additions & 20 deletions native/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

4 changes: 2 additions & 2 deletions native/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down
16 changes: 14 additions & 2 deletions native/core/src/execution/operators/iceberg_common.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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);
Expand Down
16 changes: 9 additions & 7 deletions native/core/src/execution/operators/iceberg_location_scoped.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<T>(
inner: &Inner,
bucket: &BucketLocations,
generation: u64,
result: Result<()>,
) -> Result<()> {
result: Result<T>,
) -> Result<T> {
match result {
Err(e) if may_mean_stale_locations(&e) => {
match inner.shared.refresh(bucket, generation).await {
Expand All @@ -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<FileMetadata> {
let generation = self.generation();
let result = self.writer.close().await;
refresh_on_failure(&self.inner, &self.route.bucket, generation, result).await
Expand Down Expand Up @@ -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<FileMetadata> {
self.0
.check("close", &self.1)
.map(|_| FileMetadata { size: 0 })
}
}

Expand Down Expand Up @@ -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();
Expand Down
31 changes: 13 additions & 18 deletions native/core/src/execution/planner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7298,35 +7298,30 @@ 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");

// 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");
Expand Down
9 changes: 9 additions & 0 deletions spark/src/main/scala/org/apache/comet/CometConf.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
*
Expand Down Expand Up @@ -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)
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Loading