From e7cf9cbe1da5dcea55e75f224e9fd2d5fca8be92 Mon Sep 17 00:00:00 2001 From: Alex Bianchi <75697973+alexanderbianchi@users.noreply.github.com> Date: Mon, 28 Sep 2026 22:02:56 -0400 Subject: [PATCH 1/2] feat: support decode-time Iceberg runtime selection --- iceberg/README.md | 16 ++++ iceberg/src/codec.rs | 177 +++++++++++++++++++++++++++++++++++++++++-- 2 files changed, 188 insertions(+), 5 deletions(-) diff --git a/iceberg/README.md b/iceberg/README.md index 404ec9c3c..6e927eead 100644 --- a/iceberg/README.md +++ b/iceberg/README.md @@ -29,6 +29,22 @@ The default storage factory resolves `file://`, S3 (`s3://`, `s3a://`, `s3n://`), and GCS (`gs://`, `gcs://`) URIs. Use `IcebergIntegrationOptions` to supply custom storage or an Iceberg runtime. +## Decode-time runtime selection + +`IcebergCodec::new` keeps a fixed runtime. For per-query selection, use +`IcebergCodec::new_with_runtime_resolver` (see its Rustdoc example). +Its closure receives the decoding `TaskContext` and can read worker-local session +extensions, for example an `iceberg::Runtime::new_with_split(&io, &query_cpu)`. +The caller controls CPU/I/O routing and must keep those Tokio runtimes alive +through execution. Resolver errors fail decoding; encoding and the wire format +are unchanged. + +Register the codec with `with_distributed_user_codec` **before** calling +`with_iceberg_integration`, which adds a fixed-runtime codec. Keep codec order +consistent on coordinator and workers. Install runtime extensions in each worker +query's session config; runtime handles are not sent from the coordinator. +This only changes decoded scans, not coordinator-side table planning. + ```bash cargo test -p datafusion-distributed-iceberg ``` diff --git a/iceberg/src/codec.rs b/iceberg/src/codec.rs index 173e261b6..bff2ece7c 100644 --- a/iceberg/src/codec.rs +++ b/iceberg/src/codec.rs @@ -1,3 +1,4 @@ +use std::fmt; use std::sync::Arc; use datafusion::arrow::datatypes::SchemaRef; @@ -19,11 +20,21 @@ use prost::Message; use crate::proto::generated::iceberg as pb; use crate::{IcebergDataSource, IcebergWorkUnitFeed}; +type RuntimeResolver = dyn Fn(&TaskContext) -> Result + Send + Sync; + /// Physical plan codec for [`IcebergDataSource`]. -#[derive(Debug, Clone)] +#[derive(Clone)] pub struct IcebergCodec { storage_factory: Arc, - iceberg_runtime: iceberg::Runtime, + runtime_resolver: Arc, +} + +impl fmt::Debug for IcebergCodec { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_struct("IcebergCodec") + .field("storage_factory", &self.storage_factory) + .finish_non_exhaustive() + } } impl IcebergCodec { @@ -31,10 +42,41 @@ impl IcebergCodec { pub fn new( storage_factory: Arc, iceberg_runtime: iceberg::Runtime, + ) -> Self { + Self::new_with_runtime_resolver(storage_factory, move |_| Ok(iceberg_runtime.clone())) + } + + /// Creates a codec that selects a worker-local runtime for each decoded scan. + /// + /// The resolver receives the decoding task's context, so it can read a runtime + /// from a worker-local session extension. It is not called during encoding; + /// resolver errors fail decoding. CPU/I/O routing is entirely caller-defined, + /// and runtime handles are never serialized. The caller must keep the selected + /// Tokio runtimes alive until scan execution completes. + /// + /// # Example + /// + /// ``` + /// # use std::sync::Arc; + /// # use datafusion::common::exec_datafusion_err; + /// # use datafusion_distributed_iceberg::IcebergCodec; + /// # use iceberg::io::StorageFactory; + /// # fn example(storage_factory: Arc) { + /// let codec = IcebergCodec::new_with_runtime_resolver(storage_factory, |ctx| { + /// ctx.session_config() + /// .get_extension::() + /// .map(|runtime| runtime.as_ref().clone()) + /// .ok_or_else(|| exec_datafusion_err!("missing worker Iceberg runtime")) + /// }); + /// # } + /// ``` + pub fn new_with_runtime_resolver( + storage_factory: Arc, + runtime_resolver: impl Fn(&TaskContext) -> Result + Send + Sync + 'static, ) -> Self { Self { storage_factory, - iceberg_runtime, + runtime_resolver: Arc::new(runtime_resolver), } } } @@ -99,7 +141,7 @@ impl PhysicalExtensionCodec for IcebergCodec { column_stats: None, table_snapshot: None, iceberg_file_io, - iceberg_runtime: self.iceberg_runtime.clone(), + iceberg_runtime: (self.runtime_resolver)(ctx)?, feed, })) } @@ -142,10 +184,16 @@ impl PhysicalExtensionCodec for IcebergCodec { #[cfg(test)] mod tests { - use datafusion::common::Statistics; + use datafusion::common::{Statistics, exec_err}; use datafusion::datasource::source::DataSource; + use datafusion::prelude::SessionConfig; + use datafusion_distributed::{DistributedCodec, DistributedExt}; + use datafusion_proto::physical_plan::AsExecutionPlan; + use datafusion_proto::protobuf::PhysicalPlanNode; + use tokio::runtime::{Builder, Handle, Runtime as TokioRuntime}; use super::*; + use crate::common::df_err; use crate::test_utils::IcebergTestHarness; #[tokio::test] @@ -189,6 +237,125 @@ mod tests { Ok(()) } + #[test] + fn fixed_runtime_is_preserved() -> Result<()> { + let io = runtime()?; + let cpu = runtime()?; + let codec = IcebergCodec::new( + Arc::new(OpenDalResolvingStorageFactory::new()), + iceberg::Runtime::new_with_split(&io, &cpu), + ); + io.block_on(async { + let proto = encoded_scan(&codec).await?; + let ctx = task_context(iceberg::Runtime::new(&io)); + assert_decoded_runtime(&proto, &codec, &ctx, &io, &cpu).await + }) + } + + #[test] + fn default_runtime_is_captured_at_construction() -> Result<()> { + let startup = runtime()?; + let worker = runtime()?; + let codec = { + let _guard = startup.enter(); + IcebergCodec::default() + }; + worker.block_on(async { + let proto = encoded_scan(&codec).await?; + assert_decoded_runtime(&proto, &codec, &TaskContext::default(), &startup, &startup) + .await + }) + } + + #[test] + fn resolves_each_decode_with_query_cpu_and_shared_io() -> Result<()> { + let codec = IcebergCodec::new_with_runtime_resolver( + Arc::new(OpenDalResolvingStorageFactory::new()), + |ctx| { + Ok(ctx + .session_config() + .get_extension::() + .expect("worker runtime extension") + .as_ref() + .clone()) + }, + ); + let config = SessionConfig::new().with_distributed_user_codec(codec); + let codec = DistributedCodec::new_combined_with_user(&config); + let io = runtime()?; + let cpu_a = runtime()?; + let cpu_b = runtime()?; + io.block_on(async { + let proto = encoded_scan(&codec).await?; + let ctx_a = task_context(iceberg::Runtime::new_with_split(&io, &cpu_a)); + let ctx_b = task_context(iceberg::Runtime::new_with_split(&io, &cpu_b)); + assert_decoded_runtime(&proto, &codec, &ctx_a, &io, &cpu_a).await?; + assert_decoded_runtime(&proto, &codec, &ctx_b, &io, &cpu_b).await + }) + } + + #[tokio::test] + async fn propagates_runtime_resolver_errors_only_on_decode() -> Result<()> { + let codec = IcebergCodec::new_with_runtime_resolver( + Arc::new(OpenDalResolvingStorageFactory::new()), + |_| exec_err!("query runtime unavailable"), + ); + let proto = encoded_scan(&codec).await?; + let error = proto + .try_into_physical_plan(&TaskContext::default(), &codec) + .unwrap_err(); + assert!( + error.to_string().contains("query runtime unavailable"), + "{error}" + ); + Ok(()) + } + + fn runtime() -> Result { + Ok(Builder::new_multi_thread() + .worker_threads(1) + .enable_all() + .build()?) + } + + fn task_context(runtime: iceberg::Runtime) -> TaskContext { + TaskContext::default() + .with_session_config(SessionConfig::new().with_extension(Arc::new(runtime))) + } + + async fn encoded_scan(codec: &dyn PhysicalExtensionCodec) -> Result { + let harness = IcebergTestHarness::new().await?; + PhysicalPlanNode::try_from_physical_plan(harness.scan().await?, codec) + } + + async fn assert_decoded_runtime( + proto: &PhysicalPlanNode, + codec: &dyn PhysicalExtensionCodec, + ctx: &TaskContext, + io: &TokioRuntime, + cpu: &TokioRuntime, + ) -> Result<()> { + let decoded = proto.try_into_physical_plan(ctx, codec)?; + let runtime = &iceberg_source(&decoded)?.iceberg_runtime; + assert_eq!( + runtime + .io() + .spawn(async { Handle::current().id() }) + .await + .map_err(df_err)?, + io.handle().id(), + ); + assert_eq!( + runtime + .cpu() + .spawn(async { Handle::current().id() }) + .await + .map_err(df_err)?, + cpu.handle().id(), + ); + Ok(()) + } + fn iceberg_plan(plan: &Arc) -> Result> { if let Some(exec) = plan.downcast_ref::() && exec From 9361198d779db292171a313d7b9345a36b5967df Mon Sep 17 00:00:00 2001 From: Alex Bianchi <75697973+alexanderbianchi@users.noreply.github.com> Date: Wed, 30 Sep 2026 20:23:18 -0400 Subject: [PATCH 2/2] docs: trim Iceberg runtime selection documentation --- iceberg/README.md | 16 ---------------- iceberg/src/codec.rs | 26 +++----------------------- 2 files changed, 3 insertions(+), 39 deletions(-) diff --git a/iceberg/README.md b/iceberg/README.md index 6e927eead..404ec9c3c 100644 --- a/iceberg/README.md +++ b/iceberg/README.md @@ -29,22 +29,6 @@ The default storage factory resolves `file://`, S3 (`s3://`, `s3a://`, `s3n://`), and GCS (`gs://`, `gcs://`) URIs. Use `IcebergIntegrationOptions` to supply custom storage or an Iceberg runtime. -## Decode-time runtime selection - -`IcebergCodec::new` keeps a fixed runtime. For per-query selection, use -`IcebergCodec::new_with_runtime_resolver` (see its Rustdoc example). -Its closure receives the decoding `TaskContext` and can read worker-local session -extensions, for example an `iceberg::Runtime::new_with_split(&io, &query_cpu)`. -The caller controls CPU/I/O routing and must keep those Tokio runtimes alive -through execution. Resolver errors fail decoding; encoding and the wire format -are unchanged. - -Register the codec with `with_distributed_user_codec` **before** calling -`with_iceberg_integration`, which adds a fixed-runtime codec. Keep codec order -consistent on coordinator and workers. Install runtime extensions in each worker -query's session config; runtime handles are not sent from the coordinator. -This only changes decoded scans, not coordinator-side table planning. - ```bash cargo test -p datafusion-distributed-iceberg ``` diff --git a/iceberg/src/codec.rs b/iceberg/src/codec.rs index bff2ece7c..35bdef978 100644 --- a/iceberg/src/codec.rs +++ b/iceberg/src/codec.rs @@ -46,30 +46,10 @@ impl IcebergCodec { Self::new_with_runtime_resolver(storage_factory, move |_| Ok(iceberg_runtime.clone())) } - /// Creates a codec that selects a worker-local runtime for each decoded scan. + /// Selects a runtime from the [`TaskContext`] for each decoded scan. /// - /// The resolver receives the decoding task's context, so it can read a runtime - /// from a worker-local session extension. It is not called during encoding; - /// resolver errors fail decoding. CPU/I/O routing is entirely caller-defined, - /// and runtime handles are never serialized. The caller must keep the selected - /// Tokio runtimes alive until scan execution completes. - /// - /// # Example - /// - /// ``` - /// # use std::sync::Arc; - /// # use datafusion::common::exec_datafusion_err; - /// # use datafusion_distributed_iceberg::IcebergCodec; - /// # use iceberg::io::StorageFactory; - /// # fn example(storage_factory: Arc) { - /// let codec = IcebergCodec::new_with_runtime_resolver(storage_factory, |ctx| { - /// ctx.session_config() - /// .get_extension::() - /// .map(|runtime| runtime.as_ref().clone()) - /// .ok_or_else(|| exec_datafusion_err!("missing worker Iceberg runtime")) - /// }); - /// # } - /// ``` + /// The resolver is not called during encoding; its errors fail decoding. + /// Callers must keep the selected Tokio runtimes alive through execution. pub fn new_with_runtime_resolver( storage_factory: Arc, runtime_resolver: impl Fn(&TaskContext) -> Result + Send + Sync + 'static,