Skip to content

feat: support decode-time Iceberg runtime selection - #758

Open
alexanderbianchi wants to merge 2 commits into
datafusion-contrib:mainfrom
alexanderbianchi:iceberg-codec-runtime-selection
Open

alexanderbianchi wants to merge 2 commits into
datafusion-contrib:mainfrom
alexanderbianchi:iceberg-codec-runtime-selection

Conversation

@alexanderbianchi

@alexanderbianchi alexanderbianchi commented Sep 29, 2026 •

Copy link
Copy Markdown
Collaborator

Summary

Allow a codec registered at startup to select a worker query's Iceberg runtime during plan decoding. This supports load/priority isolation when queries use different CPU runtimes while sharing a dedicated I/O runtime; it is not a data-correctness fix or a measured performance improvement.

  • Add IcebergCodec::new_with_runtime_resolver, taking a fallible closure over the decoding TaskContext.
  • Preserve new() and Default fixed-runtime behavior, encoding, and the wire format. Runtime handles remain local, and CPU/I/O policy stays with the application.
  • No worker-plumbing or fixture changes.

Validation

  • cargo fmt --all -- --check
  • cargo test --locked -p datafusion-distributed-iceberg --features integration — 88 unit/integration tests passed before the documentation-only cleanup.
  • cargo test --locked -p datafusion-distributed-iceberg --features integration --doc — the remaining doctest passed after removing the verbose example.
  • cargo clippy --locked -p datafusion-distributed-iceberg --all-targets --all-features -- -D warnings
  • RUSTDOCFLAGS='-D warnings' cargo doc --locked -p datafusion-distributed-iceberg --no-deps --all-features

New codec tests probe runtime IDs through the handles installed in decoded fixture scans: fixed/default compatibility, two query CPU runtimes sharing distinct I/O, and resolver-error propagation. Temporarily bypassing resolution made all four new tests fail; restoring it passed. These are focused codec tests, not an end-to-end scheduler or performance test.

@alexanderbianchi
alexanderbianchi marked this pull request as ready for review October 1, 2026 00:22
Comment thread iceberg/src/codec.rs
use crate::proto::generated::iceberg as pb;
use crate::{IcebergDataSource, IcebergWorkUnitFeed};

type RuntimeResolver = dyn Fn(&TaskContext) -> Result<iceberg::Runtime> + Send + Sync;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is threaded a bit weirdly. I don't understand why this would be necessary though.

an IcebergCodec is something that users are responsible for injecting per-session, not once at startup, as it goes in SessionConfig.extensions. This means that users should be able to just create one IcebergCodec per-query with the correct runtime already attached, rather than delegating to an opaque callback.

This branch has not been deployed

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants