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]

Reply via email to