mbutrovich commented on code in PR #45:
URL: https://github.com/apache/datafusion-iceberg/pull/45#discussion_r4189157875
##########
crates/datafusion/src/schema.rs:
##########
@@ -218,32 +209,56 @@ impl SchemaProvider for IcebergSchemaProvider {
let tables = self.tables.clone();
let table_name = name.to_string();
- // Use tokio's spawn_blocking to handle the async work on a blocking
thread pool
- let result = tokio::task::spawn_blocking(move || {
- let rt = tokio::runtime::Handle::current();
- rt.block_on(async move {
- let table_ident = TableIdent::new(namespace,
table_name.clone());
+ run_on_catalog_runtime("drop-table", async move {
+ let table_ident = TableIdent::new(namespace, table_name.clone());
- // Drop the table from the Iceberg catalog
- catalog
- .drop_table(&table_ident)
- .await
- .map_err(to_datafusion_error)?;
+ // Drop the table from the Iceberg catalog
+ catalog
+ .drop_table(&table_ident)
+ .await
+ .map_err(to_datafusion_error)?;
- // Remove from local cache and return the removed provider
- let removed = tables
- .remove(&table_name)
- .map(|(_, table)| table as Arc<dyn TableProvider>);
+ // Remove from local cache and return the removed provider
+ let removed = tables
+ .remove(&table_name)
+ .map(|(_, table)| table as Arc<dyn TableProvider>);
- Ok(removed)
- })
- });
-
- futures::executor::block_on(result)
- .map_err(|e| exec_datafusion_err!("Failed to drop Iceberg table:
{e}"))?
+ Ok(removed)
+ })
}
}
+/// Runs an async catalog operation on its own Tokio runtime thread.
+///
+/// `SchemaProvider` exposes synchronous registration methods. Blocking the
caller's
+/// runtime while polling an async catalog operation can starve its I/O and
timer
+/// drivers, and relying on `Handle::current()` also fails for synchronous
callers
+/// outside a Tokio runtime.
+fn run_on_catalog_runtime<T, F>(operation: &'static str, future: F) ->
Result<T>
+where
+ T: Send + 'static,
+ F: Future<Output = Result<T>> + Send + 'static,
+{
+ std::thread::Builder::new()
+ .name(format!("iceberg-{operation}"))
+ .spawn(move || {
+ let runtime = tokio::runtime::Builder::new_current_thread()
+ .enable_all()
+ .build()
+ .map_err(|error| {
+ exec_datafusion_err!(
+ "Failed to create Tokio runtime for {operation}:
{error}"
+ )
+ })?;
+ runtime.block_on(future)
+ })
+ .map_err(|error| {
+ exec_datafusion_err!("Failed to spawn {operation} thread: {error}")
+ })?
+ .join()
+ .map_err(|_| exec_datafusion_err!("{operation} thread panicked"))?
+}
Review Comment:
Does this fix the hang for a catalog whose I/O was set up on the caller's
runtime? If I'm reading it right, the caller still blocks in `join()`, and on a
[current-thread
runtime](https://docs.rs/tokio/1.53.1/tokio/runtime/struct.Builder.html#method.new_current_thread)
(one where only the thread inside `Runtime::block_on` runs tasks and drives
the I/O and timer drivers) nothing else can run that runtime's tasks while it
waits. The `Handle::block_on` docs describe the same limit:
> When this is used on a `current_thread` runtime, only the
`Runtime::block_on` method can drive the IO and timer drivers, but the
`Handle::block_on` method cannot drive them.
The new runtime only helps for resources it creates itself, like the
`tokio::time::sleep` in `DelayedCatalog`. Real catalogs have already done I/O
on the caller's runtime by the time a DDL statement runs.
[`IcebergCatalogProvider::try_new`](https://github.com/apache/datafusion-iceberg/blob/9a959d635f21d3158cc65e9f64027893de9945b2/crates/datafusion/src/catalog.rs#L50-L85)
lists namespaces and tables and loads every table. The REST catalog's
[`HttpClient`](https://github.com/apache/iceberg-rust/blob/8cb2adeddb4f8da1ca7bd86ca303337f011c12e2/crates/catalog/rest/src/client.rs#L37-L38)
wraps a `reqwest::Client`, which keeps those connections alive in a pool, and
hyper drives each pooled connection with a task spawned on the runtime that
opened it. `CREATE TABLE` then reuses a pooled connection from the new runtime
and waits on a task that only the blocked caller thread can poll.
I tried two things to check this:
1. A standalone program with the versions in this repo's `Cargo.lock`
(reqwest 0.12.28, tokio 1.53.1). It sends one request from a current-thread
runtime, then sends a second request with the same client through a copy of
`run_on_catalog_runtime`. The second request never completes (no result after 5
seconds). The same program with a multi-thread caller completes.
2. In this PR's tests, I gave `DelayedCatalog` a `Handle::current()`
captured at construction and changed `create_table` to
`self.1.spawn(tokio::time::sleep(..)).await` instead of sleeping inline. That
stands in for a task on the caller's runtime.
`test_register_table_with_async_catalog_on_current_thread_runtime` then fails
at the head commit with "registration deadlocked the current-thread runtime:
Timeout".
What do you think about choosing the strategy from `Handle::try_current()`
instead of always spawning a thread?
- On a multi-thread runtime,
[`block_in_place`](https://docs.rs/tokio/1.53.1/tokio/task/fn.block_in_place.html)
plus `Handle::block_on` runs the catalog call on the caller's runtime, which
the other workers keep driving. This is the case `datafusion-cli`, the
playground, and the sqllogictest harness hit. It doesn't spawn a thread or
build a runtime per statement, and the catalog's connections stay on the
runtime that owns them. The reqwest program above completes with this approach,
with a pooled connection and a single worker thread.
- With no runtime, the calling thread can build a current-thread runtime and
call `block_on` itself, so it doesn't need a second thread either.
- On a current-thread runtime there's no way to block the caller and still
drive its I/O. One option is to return an error there, which #37's expected
behavior allows ("complete or return an error on either runtime flavor"). The
cost is that the existing `#[tokio::test]` tests in this file, which run on a
current-thread runtime with `MemoryCatalog` and pass today, would need `flavor
= "multi_thread"`. The other option is to keep this PR's separate thread for
that case only, document that a catalog with I/O on the caller's runtime still
hangs, and keep #37 open, or file a tracking issue for the async boundary
discussed in section 3 of #24. Which would you prefer?
With the first two branches the future no longer needs to be `Send +
'static`, so `register_table` and `deregister_table` could drop their clones of
`catalog`, `namespace`, and `tables`. Here's a sketch that I compiled against
this branch. `test_table_registration_without_caller_runtime` and a
`multi_thread` version of the registration test pass with it:
```rust
fn block_on_catalog<T>(operation: &str, future: impl Future<Output =
Result<T>>) -> Result<T> {
match Handle::try_current() {
Ok(handle) if handle.runtime_flavor() == RuntimeFlavor::MultiThread
=> {
block_in_place(|| handle.block_on(future))
}
Ok(_) => exec_err!("{operation} requires a multi-thread Tokio
runtime"),
Err(_) => tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.map_err(|e| exec_datafusion_err!("Failed to create Tokio
runtime: {e}"))?
.block_on(future),
}
}
```
##########
crates/datafusion/src/schema.rs:
##########
@@ -306,6 +325,193 @@ mod tests {
(provider, temp_dir)
}
+ #[derive(Debug)]
+ struct DelayedCatalog(Arc<dyn Catalog>);
+
+ #[async_trait::async_trait]
+ impl Catalog for DelayedCatalog {
+ async fn list_namespaces(
+ &self,
+ parent: Option<&NamespaceIdent>,
+ ) -> iceberg::Result<Vec<NamespaceIdent>> {
+ self.0.list_namespaces(parent).await
+ }
+
+ async fn create_namespace(
+ &self,
+ namespace: &NamespaceIdent,
+ properties: HashMap<String, String>,
+ ) -> iceberg::Result<Namespace> {
+ self.0.create_namespace(namespace, properties).await
+ }
+
+ async fn get_namespace(
+ &self,
+ namespace: &NamespaceIdent,
+ ) -> iceberg::Result<Namespace> {
+ self.0.get_namespace(namespace).await
+ }
+
+ async fn namespace_exists(
+ &self,
+ namespace: &NamespaceIdent,
+ ) -> iceberg::Result<bool> {
+ self.0.namespace_exists(namespace).await
+ }
+
+ async fn update_namespace(
+ &self,
+ namespace: &NamespaceIdent,
+ properties: HashMap<String, String>,
+ ) -> iceberg::Result<()> {
+ self.0.update_namespace(namespace, properties).await
+ }
+
+ async fn drop_namespace(
+ &self,
+ namespace: &NamespaceIdent,
+ ) -> iceberg::Result<()> {
+ self.0.drop_namespace(namespace).await
+ }
+
+ async fn list_tables(
+ &self,
+ namespace: &NamespaceIdent,
+ ) -> iceberg::Result<Vec<TableIdent>> {
+ self.0.list_tables(namespace).await
+ }
+
+ async fn create_table(
+ &self,
+ namespace: &NamespaceIdent,
+ creation: TableCreation,
+ ) -> iceberg::Result<Table> {
+ tokio::time::sleep(std::time::Duration::from_millis(5)).await;
+ self.0.create_table(namespace, creation).await
+ }
+
+ async fn load_table(&self, table: &TableIdent) ->
iceberg::Result<Table> {
+ self.0.load_table(table).await
+ }
+
+ async fn drop_table(&self, table: &TableIdent) -> iceberg::Result<()> {
+ tokio::time::sleep(std::time::Duration::from_millis(5)).await;
+ self.0.drop_table(table).await
+ }
+
+ async fn purge_table(&self, table: &TableIdent) -> iceberg::Result<()>
{
+ self.0.purge_table(table).await
+ }
+
+ async fn table_exists(&self, table: &TableIdent) ->
iceberg::Result<bool> {
+ self.0.table_exists(table).await
+ }
+
+ async fn rename_table(
+ &self,
+ src: &TableIdent,
+ dest: &TableIdent,
+ ) -> iceberg::Result<()> {
+ self.0.rename_table(src, dest).await
+ }
+
+ async fn register_table(
+ &self,
+ table: &TableIdent,
+ metadata_location: String,
+ ) -> iceberg::Result<Table> {
+ self.0.register_table(table, metadata_location).await
+ }
+
+ async fn update_table(&self, commit: TableCommit) ->
iceberg::Result<Table> {
+ self.0.update_table(commit).await
+ }
+ }
+
+ async fn create_delayed_test_schema_provider() -> (IcebergSchemaProvider,
TempDir) {
+ let temp_dir = TempDir::new().unwrap();
+ let warehouse_path = temp_dir.path().to_str().unwrap().to_string();
+ let catalog = MemoryCatalogBuilder::default()
+ .load(
+ "memory",
+ HashMap::from([(MEMORY_CATALOG_WAREHOUSE.to_string(),
warehouse_path)]),
+ )
+ .await
+ .unwrap();
+ let namespace = NamespaceIdent::new("test_ns".to_string());
+ catalog
+ .create_namespace(&namespace, HashMap::new())
+ .await
+ .unwrap();
+ let catalog: Arc<dyn Catalog> =
Arc::new(DelayedCatalog(Arc::new(catalog)));
+ let provider = IcebergSchemaProvider::try_new(catalog, namespace)
+ .await
+ .unwrap();
+ (provider, temp_dir)
+ }
+
+ #[test]
+ fn test_register_table_with_async_catalog_on_current_thread_runtime() {
+ let (sender, receiver) = std::sync::mpsc::channel();
+ std::thread::spawn(move || {
+ let runtime = tokio::runtime::Builder::new_current_thread()
+ .enable_all()
+ .build()
+ .unwrap();
+ runtime.block_on(async move {
+ let (schema_provider, _temp_dir) =
+ create_delayed_test_schema_provider().await;
+ let arrow_schema = Arc::new(ArrowSchema::new(vec![Field::new(
+ "id",
+ DataType::Int32,
+ false,
+ )]));
+ let empty_batch = RecordBatch::new_empty(arrow_schema.clone());
+ let mem_table =
+ MemTable::try_new(arrow_schema,
vec![vec![empty_batch]]).unwrap();
+
+ let result = schema_provider
+ .register_table("async_table".to_string(),
Arc::new(mem_table))
+ .and_then(|_| {
+
schema_provider.deregister_table("async_table").map(|_| ())
+ })
+ .map_err(|error| error.to_string());
+ let _ = sender.send(result);
+ });
+ });
+
+ let result = receiver
+ .recv_timeout(std::time::Duration::from_secs(5))
+ .expect("registration deadlocked the current-thread runtime");
+ assert!(
+ result.is_ok(),
+ "expected table registration and deregistration to complete:
{result:?}"
+ );
+ }
Review Comment:
Could the tests cover a catalog whose work depends on the caller's runtime,
like the `Handle`-capturing variant of `DelayedCatalog` described in the other
comment? This test passes because the sleep is created on whichever runtime
polls it, so it doesn't exercise the REST-style case from #37. Depending on how
the current-thread case is resolved, that test would assert either the error or
the documented limitation. Could you also add a `#[tokio::test(flavor =
"multi_thread")]` case for register and deregister? That's the flavor most
callers use, and with a per-flavor strategy it takes a different code path from
the two tests here.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]