diff --git a/Cargo.lock b/Cargo.lock index 3802cc03ff..d3a0ba9d08 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3957,7 +3957,10 @@ dependencies = [ "expect-test", "futures", "iceberg", + "iceberg-catalog-rest", + "mockito", "parquet", + "serde_json", "tempfile", "tokio", "uuid", diff --git a/crates/catalog/rest/public-api.txt b/crates/catalog/rest/public-api.txt index d7e9d95e97..9f30656ffb 100644 --- a/crates/catalog/rest/public-api.txt +++ b/crates/catalog/rest/public-api.txt @@ -289,6 +289,7 @@ pub struct iceberg_catalog_rest::RestCatalogBuilder impl iceberg_catalog_rest::RestCatalogBuilder pub fn iceberg_catalog_rest::RestCatalogBuilder::with_auth_manager(self, auth_manager: alloc::sync::Arc) -> Self pub fn iceberg_catalog_rest::RestCatalogBuilder::with_client(self, client: reqwest::async_impl::client::Client) -> Self +pub fn iceberg_catalog_rest::RestCatalogBuilder::with_page_size(self, page_size: u32) -> Self pub fn iceberg_catalog_rest::RestCatalogBuilder::with_session_context(self, context: iceberg::catalog::session::SessionContext) -> Self impl core::default::Default for iceberg_catalog_rest::RestCatalogBuilder pub fn iceberg_catalog_rest::RestCatalogBuilder::default() -> iceberg_catalog_rest::RestCatalogBuilder @@ -325,6 +326,7 @@ pub fn iceberg_catalog_rest::RestSessionCatalogBuilder::load(self, name: impl co pub fn iceberg_catalog_rest::RestSessionCatalogBuilder::with_auth_manager(self, auth_manager: alloc::sync::Arc) -> Self pub fn iceberg_catalog_rest::RestSessionCatalogBuilder::with_client(self, client: reqwest::async_impl::client::Client) -> Self pub fn iceberg_catalog_rest::RestSessionCatalogBuilder::with_kms_client_factory(self, kms_client_factory: alloc::sync::Arc) -> Self +pub fn iceberg_catalog_rest::RestSessionCatalogBuilder::with_page_size(self, page_size: u32) -> Self pub fn iceberg_catalog_rest::RestSessionCatalogBuilder::with_runtime(self, runtime: iceberg::runtime::Runtime) -> Self pub fn iceberg_catalog_rest::RestSessionCatalogBuilder::with_storage_factory(self, storage_factory: alloc::sync::Arc) -> Self impl core::default::Default for iceberg_catalog_rest::RestSessionCatalogBuilder @@ -381,6 +383,7 @@ pub const iceberg_catalog_rest::AUTH_TYPE_NONE: &str pub const iceberg_catalog_rest::AUTH_TYPE_OAUTH2: &str pub const iceberg_catalog_rest::REST_CATALOG_PROP_AUTH_TYPE: &str pub const iceberg_catalog_rest::REST_CATALOG_PROP_DISABLE_HEADER_REDACTION: &str +pub const iceberg_catalog_rest::REST_CATALOG_PROP_PAGE_SIZE: &str pub const iceberg_catalog_rest::REST_CATALOG_PROP_URI: &str pub const iceberg_catalog_rest::REST_CATALOG_PROP_WAREHOUSE: &str pub trait iceberg_catalog_rest::AuthManager: core::fmt::Debug + core::marker::Send + core::marker::Sync diff --git a/crates/catalog/rest/src/catalog.rs b/crates/catalog/rest/src/catalog.rs index 841d5d33ef..8831cf288e 100644 --- a/crates/catalog/rest/src/catalog.rs +++ b/crates/catalog/rest/src/catalog.rs @@ -20,6 +20,7 @@ use std::collections::{HashMap, HashSet}; use std::fmt::{Debug, Formatter}; use std::future::Future; +use std::num::NonZeroU32; use std::str::FromStr; use std::sync::{Arc, OnceLock}; @@ -56,6 +57,9 @@ use crate::types::{ pub const REST_CATALOG_PROP_URI: &str = "uri"; /// REST catalog warehouse location pub const REST_CATALOG_PROP_WAREHOUSE: &str = "warehouse"; +/// Requested maximum number of results per list response (a positive `u32`). +/// When unset, no `pageSize` query parameter is sent. +pub const REST_CATALOG_PROP_PAGE_SIZE: &str = "rest-page-size"; /// Disable header redaction in error logs and `Debug` output (defaults to /// false for security) pub const REST_CATALOG_PROP_DISABLE_HEADER_REDACTION: &str = "disable-header-redaction"; @@ -113,6 +117,15 @@ impl CatalogBuilder for RestCatalogBuilder { } impl RestCatalogBuilder { + /// Sets the requested page size for namespace and table listings. + /// + /// See [`RestSessionCatalogBuilder::with_page_size`] for configuration + /// precedence and validation. All pages are still returned by list operations. + pub fn with_page_size(mut self, page_size: u32) -> Self { + self.inner = self.inner.with_page_size(page_size); + self + } + /// Configures the catalog with a custom HTTP client. pub fn with_client(mut self, client: Client) -> Self { self.inner = self.inner.with_client(client); @@ -219,6 +232,22 @@ impl RestCatalogConfig { self.props.get("oauth2-server-uri").cloned() } + /// Parse only the effective value, after server defaults and overrides are merged. + fn page_size(&self) -> Result> { + self.props + .get(REST_CATALOG_PROP_PAGE_SIZE) + .map(|value| { + value.parse::().map_err(|err| { + Error::new( + ErrorKind::DataInvalid, + format!("{REST_CATALOG_PROP_PAGE_SIZE} must be a positive u32"), + ) + .with_source(err) + }) + }) + .transpose() + } + fn namespaces_endpoint(&self) -> String { self.url_prefixed(&["namespaces"]) } @@ -423,6 +452,8 @@ struct RestClient { config: RestCatalogConfig, /// Capabilities the server advertises (see [`RestSessionCatalog::supports_endpoint`]). endpoints: HashSet, + /// Validated page size from the merged runtime configuration. + page_size: Option, } impl RestClient { @@ -456,6 +487,7 @@ impl RestClient { _ => crate::endpoint::DEFAULT_ENDPOINTS.clone(), }; let config = user_config.clone().merge_with_config(catalog_config); + let page_size = config.page_size()?; let http_client = http_client.update_with(&config)?; // The manager is handed an unauthenticated client: its own // requests must not be signed by the session it is deriving. @@ -470,6 +502,7 @@ impl RestClient { config, http_client: http_client.with_auth_session(session), endpoints, + page_size, }) } @@ -893,11 +926,16 @@ impl SessionCatalog for RestSessionCatalog { let client = self.client().await?; let endpoint = client.config.namespaces_endpoint(); let mut namespaces = Vec::new(); - let mut next_token = None; + // An empty initial token explicitly opts in to pagination. + let mut next_token = client.page_size.map(|_| String::new()); loop { let mut request = client.http_client.request(Method::GET, endpoint.clone()); + if let Some(page_size) = client.page_size { + request = request.query(&[("pageSize", page_size)]); + } + // Filter on `parent={namespace}` if a parent namespace exists. if let Some(ns) = parent { request = request.query(&[("parent", ns.to_url_string())]); @@ -1076,11 +1114,16 @@ impl SessionCatalog for RestSessionCatalog { let client = self.client().await?; let endpoint = client.config.tables_endpoint(namespace); let mut identifiers = Vec::new(); - let mut next_token = None; + // An empty initial token explicitly opts in to pagination. + let mut next_token = client.page_size.map(|_| String::new()); loop { let mut request = client.http_client.request(Method::GET, endpoint.clone()); + if let Some(page_size) = client.page_size { + request = request.query(&[("pageSize", page_size)]); + } + if let Some(token) = next_token { request = request.query(&[("pageToken", token)]); } @@ -1515,6 +1558,24 @@ impl Default for RestSessionCatalogBuilder { } impl RestSessionCatalogBuilder { + /// Sets the requested page size for namespace and table listings. + /// + /// This overrides server defaults. The `rest-page-size` property passed to + /// [`load`](Self::load) takes precedence over this value, and server overrides + /// take precedence over both. There is no client-side default. + /// + /// The effective value must be a positive `u32`; an invalid value produces + /// [`ErrorKind::DataInvalid`] when the catalog is first used, after the server + /// configuration is fetched. List operations continue to return all results, + /// fetching as many pages as necessary. + pub fn with_page_size(mut self, page_size: u32) -> Self { + self.config.props.insert( + REST_CATALOG_PROP_PAGE_SIZE.to_string(), + page_size.to_string(), + ); + self + } + /// Configures the catalog with a custom HTTP client. pub fn with_client(mut self, client: Client) -> Self { self.config.client = Some(client); @@ -1617,10 +1678,11 @@ impl RestSessionCatalogBuilder { } // Collect other remaining properties - self.config.props = props - .into_iter() - .filter(|(k, _)| k != REST_CATALOG_PROP_URI && k != REST_CATALOG_PROP_WAREHOUSE) - .collect(); + self.config.props.extend( + props + .into_iter() + .filter(|(k, _)| k != REST_CATALOG_PROP_URI && k != REST_CATALOG_PROP_WAREHOUSE), + ); async move { if self.config.name.is_none() { @@ -1706,6 +1768,148 @@ mod tests { ) } + #[tokio::test] + async fn test_page_size_precedence_and_pagination() { + // Server default < builder < load properties < server override. + for (default, builder, property, override_value, expected) in [ + (None, None, None, None, None), + (Some("10"), None, None, None, Some(10)), + (Some("10"), Some(20), None, None, Some(20)), + (None, Some(20), None, None, Some(20)), + (Some("10"), None, Some("30"), None, Some(30)), + (Some("10"), Some(20), Some("30"), None, Some(30)), + (Some("10"), Some(20), Some("30"), Some("40"), Some(40)), + (None, None, None, Some("40"), Some(40)), + (None, Some(1), None, None, Some(1)), + (None, Some(u32::MAX), None, None, Some(u32::MAX)), + // Invalid lower-priority values must not reject a valid effective value. + (Some("invalid"), Some(20), None, None, Some(20)), + (None, Some(0), Some("invalid"), Some("40"), Some(40)), + ] { + let mut server = Server::new_async().await; + let page_size_props = |value: Option<&str>| -> HashMap { + value + .map(|value| (REST_CATALOG_PROP_PAGE_SIZE.to_string(), value.to_string())) + .into_iter() + .collect() + }; + let config_mock = server + .mock("GET", "/v1/config") + .with_status(200) + .with_body( + json!({ + "defaults": page_size_props(default), + "overrides": page_size_props(override_value), + }) + .to_string(), + ) + .create_async() + .await; + + let mut builder_config = RestCatalogBuilder::default(); + if let Some(page_size) = builder { + builder_config = builder_config.with_page_size(page_size); + } + let mut props = page_size_props(property); + props.insert(REST_CATALOG_PROP_URI.to_string(), server.url()); + let catalog = builder_config.load("test", props).await.unwrap(); + + // Check every request, including the empty token that opts in on page 1. + // Also verify that page size composes with a multipart parent namespace. + let mut mocks = Vec::new(); + for (endpoint, parent, first_body, last_body) in [ + ( + "/v1/namespaces", + "parent=parent%1Fchild", + json!({"namespaces": [["ns1"]], "next-page-token": "next/+"}), + json!({"namespaces": [["ns2"]], "next-page-token": null}), + ), + ( + "/v1/namespaces/ns1/tables", + "", + json!({"identifiers": [{"namespace": ["ns1"], "name": "t1"}], "next-page-token": "next/+"}), + json!({"identifiers": [{"namespace": ["ns1"], "name": "t2"}]}), + ), + ] { + for (token, body) in [("", first_body), ("next%2F%2B", last_body)] { + let mut query = Vec::new(); + if let Some(page_size) = expected { + query.push(format!("pageSize={page_size}")); + } + if !parent.is_empty() { + query.push(parent.to_string()); + } + if expected.is_some() || !token.is_empty() { + query.push(format!("pageToken={token}")); + } + let query = if query.is_empty() { + mockito::Matcher::Missing + } else { + mockito::Matcher::Exact(query.join("&")) + }; + mocks.push( + server + .mock("GET", endpoint) + .match_query(query) + .with_status(200) + .with_body(body.to_string()) + .expect(1) + .create_async() + .await, + ); + } + } + + let parent = NamespaceIdent::from_vec(vec!["parent".into(), "child".into()]).unwrap(); + assert_eq!(catalog.list_namespaces(Some(&parent)).await.unwrap(), vec![ + NamespaceIdent::new("ns1".into()), + NamespaceIdent::new("ns2".into()) + ],); + let namespace = NamespaceIdent::new("ns1".into()); + assert_eq!(catalog.list_tables(&namespace).await.unwrap(), vec![ + TableIdent::new(namespace.clone(), "t1".into()), + TableIdent::new(namespace, "t2".into()) + ],); + config_mock.assert_async().await; + for mock in mocks { + mock.assert_async().await; + } + } + } + + #[tokio::test] + async fn test_invalid_page_size() { + for value in ["", "0", "-1", "abc", "1.5", "4294967296"] { + for source in ["defaults", "client", "overrides"] { + let mut server = Server::new_async().await; + let mut config = json!({"defaults": {}, "overrides": {}}); + let mut props = HashMap::from([(REST_CATALOG_PROP_URI.to_string(), server.url())]); + if source == "client" { + props.insert(REST_CATALOG_PROP_PAGE_SIZE.to_string(), value.to_string()); + } else { + config[source][REST_CATALOG_PROP_PAGE_SIZE] = json!(value); + } + let config_mock = server + .mock("GET", "/v1/config") + .with_status(200) + .with_body(config.to_string()) + .create_async() + .await; + let catalog = RestSessionCatalogBuilder::default() + .load("test", props) + .await + .unwrap(); + let err = catalog + .list_namespaces(&SessionContext::empty(), None) + .await + .unwrap_err(); + assert_eq!(err.kind(), ErrorKind::DataInvalid); + assert!(err.to_string().contains(REST_CATALOG_PROP_PAGE_SIZE)); + config_mock.assert_async().await; + } + } + } + #[tokio::test] async fn test_update_config() { let mut server = Server::new_async().await; diff --git a/crates/catalog/rest/src/lib.rs b/crates/catalog/rest/src/lib.rs index 29a171de61..bd1257d895 100644 --- a/crates/catalog/rest/src/lib.rs +++ b/crates/catalog/rest/src/lib.rs @@ -55,6 +55,16 @@ //! } //! ``` //! +//! # Pagination +//! +//! Both builders support `with_page_size(1000)` to request bounded namespace +//! and table list responses. Alternatively, set [`REST_CATALOG_PROP_PAGE_SIZE`] +//! in the properties passed to `load`; this takes precedence over the builder +//! method. Server `/v1/config` defaults have lower priority than client settings, +//! and server overrides have the highest priority. No page size is sent if unset. +//! The effective value must be a positive `u32` and is validated after fetching +//! server configuration. List operations fetch all pages and return all results. +//! //! # Session catalog API //! //! ```rust, no_run diff --git a/crates/integrations/datafusion/Cargo.toml b/crates/integrations/datafusion/Cargo.toml index 11fb7b1d0b..bed4bd0f1d 100644 --- a/crates/integrations/datafusion/Cargo.toml +++ b/crates/integrations/datafusion/Cargo.toml @@ -42,7 +42,10 @@ uuid = { workspace = true } [dev-dependencies] expect-test = { workspace = true } +iceberg-catalog-rest = { workspace = true } +mockito = { workspace = true } parquet = { workspace = true } +serde_json = { workspace = true } tempfile = { workspace = true } [lints] diff --git a/crates/integrations/datafusion/README.md b/crates/integrations/datafusion/README.md index 134a8eff4b..feebd341d9 100644 --- a/crates/integrations/datafusion/README.md +++ b/crates/integrations/datafusion/README.md @@ -20,3 +20,46 @@ # Apache Iceberg DataFusion Integration This crate contains the integration of Apache DataFusion and Apache Iceberg. + +## REST catalog pagination + +Configure pagination on the REST catalog before passing it to +`IcebergCatalogProvider`. The provider loads every namespace and table page; +page size controls REST response sizes, not the number of schemas or tables +visible to DataFusion. + +```rust +use std::collections::HashMap; +use std::sync::Arc; + +use datafusion::prelude::SessionContext; +use iceberg::CatalogBuilder; +use iceberg::io::LocalFsStorageFactory; +use iceberg_catalog_rest::{REST_CATALOG_PROP_URI, RestCatalogBuilder}; +use iceberg_datafusion::IcebergCatalogProvider; + +async fn register_catalog(context: &SessionContext) -> iceberg::Result<()> { + let catalog = RestCatalogBuilder::default() + .with_page_size(1000) + // Use the storage factory appropriate for your table locations. + .with_storage_factory(Arc::new(LocalFsStorageFactory)) + .load( + "rest", + HashMap::from([( + REST_CATALOG_PROP_URI.to_string(), + "http://localhost:8181".to_string(), + )]), + ) + .await?; + let provider = IcebergCatalogProvider::try_new(Arc::new(catalog)).await?; + context.register_catalog("rest", Arc::new(provider)); + Ok(()) +} +``` + +The `rest-page-size` property in the map passed to `load` can also configure the +page size and takes precedence over `with_page_size`. Server `/v1/config` +defaults have lower priority than client configuration, while server overrides +have the highest priority. If none of these sources sets a page size, the client +omits `pageSize` rather than imposing a default. The effective value must be a +positive `u32` and is validated after the lazy server configuration handshake. diff --git a/crates/integrations/datafusion/src/catalog.rs b/crates/integrations/datafusion/src/catalog.rs index 2c6e1ff002..1820c9999f 100644 --- a/crates/integrations/datafusion/src/catalog.rs +++ b/crates/integrations/datafusion/src/catalog.rs @@ -45,6 +45,11 @@ impl IcebergCatalogProvider { /// This method retrieves the list of namespace names /// attempts to create a schema provider for each namespace, and /// collects these providers into a `HashMap`. + /// + /// Listing pagination is handled by the supplied catalog. For a REST catalog, + /// configure `rest-page-size` (or its builder's `with_page_size` method) before + /// passing it here. All namespace and table pages are loaded regardless of + /// the configured page size. pub async fn try_new(client: Arc) -> Result { // TODO: // Schemas and providers should be cached and evicted based on time diff --git a/crates/integrations/datafusion/tests/rest_catalog_pagination_test.rs b/crates/integrations/datafusion/tests/rest_catalog_pagination_test.rs new file mode 100644 index 0000000000..fd772cebbb --- /dev/null +++ b/crates/integrations/datafusion/tests/rest_catalog_pagination_test.rs @@ -0,0 +1,190 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! REST pagination must be transparent to DataFusion catalog and schema discovery. + +use std::collections::HashMap; +use std::sync::Arc; + +use datafusion::arrow::util::pretty::pretty_format_batches; +use datafusion::catalog::CatalogProvider; +use datafusion::prelude::SessionContext; +use expect_test::expect; +use iceberg::CatalogBuilder; +use iceberg::io::LocalFsStorageFactory; +use iceberg::spec::{ + FormatVersion, NestedField, PartitionSpec, PrimitiveType, Schema, SortOrder, + TableMetadataBuilder, Type, +}; +use iceberg_catalog_rest::{REST_CATALOG_PROP_URI, RestCatalogBuilder}; +use iceberg_datafusion::IcebergCatalogProvider; +use mockito::Server; +use serde_json::json; +use tempfile::TempDir; + +#[tokio::test] +async fn test_rest_page_size_through_datafusion() { + let mut server = Server::new_async().await; + let mut mocks = vec![ + server + .mock("GET", "/v1/config") + .with_status(200) + .with_body( + json!({ + "defaults": {"rest-page-size": "3"}, + "overrides": {"rest-page-size": "1"}, + }) + .to_string(), + ) + .create_async() + .await, + ]; + + for (token, namespace, next_token) in [("", "ns1", Some("next")), ("next", "ns2", None)] { + mocks.push( + server + .mock("GET", "/v1/namespaces") + .match_query(format!("pageSize=1&pageToken={token}").as_str()) + .with_status(200) + .with_body( + json!({ + "namespaces": [[namespace]], + "next-page-token": next_token, + }) + .to_string(), + ) + .create_async() + .await, + ); + } + + let temp_dir = TempDir::new().unwrap(); + let schema = Schema::builder() + .with_fields(vec![Arc::new(NestedField::required( + 1, + "id", + Type::Primitive(PrimitiveType::Long), + ))]) + .build() + .unwrap(); + let partition_spec = PartitionSpec::builder(schema.clone()).build().unwrap(); + let sort_order = SortOrder::builder().build(&schema).unwrap(); + let metadata = TableMetadataBuilder::new( + schema, + partition_spec, + sort_order, + temp_dir.path().to_str().unwrap().to_string(), + FormatVersion::V2, + HashMap::new(), + ) + .unwrap() + .build() + .unwrap() + .metadata; + + for namespace in ["ns1", "ns2"] { + let endpoint = format!("/v1/namespaces/{namespace}/tables"); + for (token, table, next_token) in [("", "t1", Some("next")), ("next", "t2", None)] { + mocks.push( + server + .mock("GET", endpoint.as_str()) + .match_query(format!("pageSize=1&pageToken={token}").as_str()) + .with_status(200) + .with_body( + json!({ + "identifiers": [{"namespace": [namespace], "name": table}], + "next-page-token": next_token, + }) + .to_string(), + ) + .create_async() + .await, + ); + mocks.push( + server + .mock("GET", format!("{endpoint}/{table}").as_str()) + .match_query(mockito::Matcher::Missing) + // Information-schema discovery may reload table metadata. + .expect_at_least(1) + .with_status(200) + .with_body( + json!({ + "metadata-location": temp_dir.path().join("metadata.json"), + "metadata": metadata, + }) + .to_string(), + ) + .create_async() + .await, + ); + } + } + + let catalog = RestCatalogBuilder::default() + .with_page_size(2) + .with_storage_factory(Arc::new(LocalFsStorageFactory)) + .load( + "rest", + HashMap::from([(REST_CATALOG_PROP_URI.to_string(), server.url())]), + ) + .await + .unwrap(); + let provider = IcebergCatalogProvider::try_new(Arc::new(catalog)) + .await + .unwrap(); + + let mut schemas = provider.schema_names(); + schemas.sort(); + assert_eq!(schemas, ["ns1", "ns2"]); + for namespace in &schemas { + let schema = provider.schema(namespace).unwrap(); + for table in ["t1", "t2"] { + assert!(schema.table_exist(table)); + assert!(schema.table(table).await.unwrap().is_some()); + } + } + + let context = SessionContext::new_with_config( + datafusion::prelude::SessionConfig::new().with_information_schema(true), + ); + context.register_catalog("rest", Arc::new(provider)); + let batches = context + .sql( + "SELECT table_schema, table_name FROM rest.information_schema.tables \ + WHERE table_catalog = 'rest' AND table_name IN ('t1', 't2') \ + ORDER BY table_schema, table_name", + ) + .await + .unwrap() + .collect() + .await + .unwrap(); + expect![[r#" + +--------------+------------+ + | table_schema | table_name | + +--------------+------------+ + | ns1 | t1 | + | ns1 | t2 | + | ns2 | t1 | + | ns2 | t2 | + +--------------+------------+"#]] + .assert_eq(&pretty_format_batches(&batches).unwrap().to_string()); + + for mock in mocks { + mock.assert_async().await; + } +}