apache / apache/iceberg-rust

Add Tokio Runtime Handle Configuration for OpenDAL Executor

Open
#1,945 4 comments 0 reactions 0 assignees View on GitHub
enhancement
Dominant language
Rust
Stars
1.4k
Forks
567
Avg merge
2d 2h
Merged PRs (30d)
93

Description

## Summary

Add support for configuring a custom Tokio runtime handle for OpenDAL operations, enabling proper runtime segregation in DataFusion applications that separate CPU-bound and I/O-bound workloads across different Tokio runtimes.

## Motivation

Applications currently must wrap Iceberg table providers with manual runtime switching:

```rust
// Current workaround: manually spawn all operations on IO runtime
let wrapped_table = IOTableProviderWrapper::new(
iceberg_table_provider,
io_runtime_handle
);
```

This approach has a couple downsides:
- Poor seperation: Parquet fetching, decoding, and filter pushdown all happen on the I/O pool
- Complex code: Requires wrapper abstractions around every table provider

### Existing Solutions for gRPC

Similar runtime segregation requirements have been solved elegantly in other contexts. For example, Tonic's gRPC client allows configuring a custom executor:

```rust
endpoint.executor(MaybeHandleExecutor(io_handle))
```

This approach allows HTTP/network I/O to run on the I/O runtime while application logic runs on the CPU runtime.

## Proposed Solution

### API Design

Leverage iceberg-rust's existing `Extensions` mechanism to allow users to configure a Tokio runtime handle:

```rust
// In iceberg-rust: crates/iceberg/src/io/file_io.rs

/// Runtime handle for executing async I/O operations.
/// When provided, OpenDAL operations will use this runtime for spawning tasks.
#[derive(Clone, Debug)]
pub struct RuntimeHandle(pub tokio::runtime::Handle);

impl RuntimeHandle {
/// Create a new RuntimeHandle from a Tokio runtime handle
pub fn new(handle: tokio::runtime::Handle) -> Self {
Self(handle)
}

/// Get the current runtime handle
pub fn current() -> Self {
Self(tokio::runtime::Handle::current())
}
}
```

### Usage Example

```rust
// Create dedicated I/O runtime
let io_runtime = tokio::runtime::Builder::new_multi_thread()
.worker_threads(8)
.thread_name("io-pool")
.enable_io()
.enable_time()
.build()?;

// Configure FileIO with runtime handle
let file_io = FileIOBuilder::new("s3")
.with_extension(RuntimeHandle::new(io_runtime.handle().clone()))
.with_props(s3_config)
.build()?;

// Or configure via catalog
let catalog = RestCatalogBuilder::new()
.with_file_io_extension(RuntimeHandle::new(io_runtime.handle().clone()))
.with_props(catalog_config)
.build()?;
```

### Implementation Approach

1. **Create Custom OpenDAL Executor**

Implement OpenDAL's `Execute` trait with a custom Tokio executor:

```rust
// In crates/iceberg/src/io/storage.rs

#[derive(Clone)]
struct CustomTokioExecutor {
handle: tokio::runtime::Handle,
}

impl opendal::Execute for CustomTokioExecutor {
fn execute(&self, f: Pin + Send>>) {
self.handle.spawn(f);
}
}
```

2. **Extract RuntimeHandle from Extensions**

Modify storage backend builders to check for `RuntimeHandle` in extensions:

```rust
// In crates/iceberg/src/io/storage.rs

pub(crate) fn build(file_io_builder: FileIOBuilder) -> Result {
let (scheme_str, props, extensions) = file_io_builder.into_parts();

// Extract runtime handle if provided
let executor = if let Some(runtime_handle) = extensions.get::() {
let exec = CustomTokioExecutor {
handle: Arc::unwrap_or_clone(runtime_handle).0
};
Some(opendal::Executor::with(exec))
} else {
None // Use OpenDAL default
};

// ... storage initialization
}
```

3. **Apply Executor to Operators**

Configure OpenDAL operators with the custom executor:

```rust
let mut operator = Operator::new(builder)?.finish();

if let Some(executor) = executor {
operator = operator.with_executor(executor);
}

operator = operator.layer(RetryLayer::new());
```

Contributor guide

Open the contributing guide

Research direction

Start in crates/iceberg/src/io/file_io.rs to inspect Extensions and FileIOBuilder, then read crates/iceberg/src/io/storage.rs and verify the OpenDAL Execute, Executor, and operator APIs used there. Done means a configured RuntimeHandle can be passed through FileIO or catalog setup and causes OpenDAL operations to use that handle while retaining the default when none is provided.

Written by the indexing model from the issue text.

Assessment

Tech stack
rust
Domain
backend
Issue type
Feature
Difficulty
4/5
Estimated time
3-5 days
Activity status
Quiet
Clarity
Mostly clear
Newbie friendliness
52/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.