Skip to content

Native Iceberg scan and write build a new FileIO, and a new storage client, for every task #6105

Description

@mixermt

What is the problem the feature request solves?

The native Iceberg operators build a fresh FileIO for every Spark task.

IcebergScanExec::execute_with_tasks calls load_file_io on each partition (native/core/src/execution/operators/iceberg_scan.rs:181), and the native write path does the same per task (iceberg_write.rs:484). Each call runs storage_factory_for, parses the storage configuration, and for S3 with a custom access provider constructs a CometS3CredentialBridge, which calls ensureInitialized over JNI and interns global refs. The instance is dropped when the task finishes, so nothing carries over to the next task on the same executor. On a production workload that was 40,317 constructions for one query.

What a shared FileIO would additionally reuse depends on the backend, because clones share the lazily created Storage behind Arc<OnceLock<Arc<dyn Storage>>>:

Spark's own Hadoop layer shares one FileSystem per (scheme, authority, user) for the lifetime of the JVM.

This is independent of #6091, though both were found on the same workload.

Describe the potential solution

Cache FileIO instances per executor in iceberg_common.rs and hand tasks clones. Correctness rests on the cache key, because a FileIO carries the access configuration it was built with. Key on every input that shapes the client:

  • Access mode. Read and write can be granted different access by the S3 bridge.
  • Catalog name. The dispatch key used when vending access.
  • The full reference path. The S3 access bridge is constructed with the bucket and url.path() of the reference location and the JVM provider is called with exactly that pair, so a bridge built for one table must not serve another table in the same bucket. All tasks of one scan or write share the path.
  • The full catalog property bag. A REST catalog vends short-lived access that rotates, and endpoints can be reconfigured, so a changed property produces a new entry rather than reusing a client built with the previous one.

Constraints:

  • Never cache memory:///. The write path assembles manifest bytes in that in-process namespace (iceberg_write.rs:1156); sharing it would let concurrent tasks observe each other's state.
  • Bound the cache with LRU eviction, one entry at a time, and drop the evicted FileIO outside the lock: the last clone of a bridge releases JNI global refs. Per-query rotation makes the key space unbounded, and clearing everything would also discard the entry a running query is using.

Clear the cache from release_runtime so the clients are released together with the Tokio runtime.

Additional context

Implementation in #6106. The S3-family follow-up, an operator cache inside OpenDalStorage::S3 in iceberg-storage-opendal, is tracked separately in #6109.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions