Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 23 additions & 1 deletion Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions iceberg/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -27,8 +27,10 @@ delegate = "0.13"
prost = "0.14.1"
tokio = { version = "1.48", features = ["full"] }
tokio-stream = "0.1"
dashmap = "6.0.1"

[dev-dependencies]
iceberg-catalog-rest = { git = "https://github.com/apache/iceberg-rust.git", rev = "4d83bc77dc10cff851a3ef427c451ff977bc3643" }
insta = "1"
parquet = "59"
tokio = { version = "1", features = ["macros", "rt-multi-thread"] }
49 changes: 49 additions & 0 deletions iceberg/src/catalog.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
// Iceberg's Catalog backed into DataFusion catalog

use std::{collections::HashMap, sync::Arc};

use datafusion::{
catalog::{CatalogProvider, SchemaProvider},
error::Result,
};
use futures::future::try_join_all;
use iceberg::{Catalog, Runtime};

use crate::{common::df_err, schema_provider::IcebergSchemaProvider};

#[derive(Debug)]
pub struct IcebergCatalog {
schemas: HashMap<String, Arc<dyn SchemaProvider>>,
}

impl IcebergCatalog {
pub async fn try_new(catalog: Arc<dyn Catalog>, iceberg_runtime: Runtime) -> Result<Self> {
let namespaces = catalog.list_namespaces(None).await.map_err(df_err)?;

let schema_providers = try_join_all(namespaces.iter().map(|ns| {
IcebergSchemaProvider::try_new(catalog.clone(), ns.clone(), iceberg_runtime.clone())
}))
.await?;

let schemas: HashMap<String, Arc<dyn SchemaProvider>> = namespaces
.into_iter()
.zip(schema_providers)
.map(|(name, provider)| {
let key = name.as_ref().join(".");
(key, Arc::new(provider) as Arc<dyn SchemaProvider>)
})
.collect();

Ok(Self { schemas })
}
}

impl CatalogProvider for IcebergCatalog {
fn schema(&self, name: &str) -> Option<Arc<dyn SchemaProvider>> {
self.schemas.get(name).cloned()
}

fn schema_names(&self) -> Vec<String> {
self.schemas.keys().cloned().collect()
}
}
4 changes: 4 additions & 0 deletions iceberg/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,12 +6,14 @@
//! integration. It deliberately contains no distributed execution adaptation
//! and no Iceberg write or commit implementation.

mod catalog;
mod common;
mod config;
mod data_source;
mod distributed_desired_task_count_handler;
mod iceberg_ext;
mod proto;
mod schema_provider;
mod table_provider;
mod work_unit_feed;
mod work_unit_wire;
Expand All @@ -20,12 +22,14 @@ mod codec;
#[doc(hidden)]
pub mod test_utils;

pub use catalog::IcebergCatalog;
pub use codec::IcebergCodec;
pub use config::IcebergConfig;
pub use data_source::IcebergDataSource;
pub use distributed_desired_task_count_handler::iceberg_desired_task_count;
pub use iceberg_ext::IcebergExt;
pub use iceberg_ext::IcebergIntegrationOptions;
pub use schema_provider::IcebergSchemaProvider;
pub use table_provider::IcebergCatalogTableProvider;
pub use table_provider::IcebergStaticTableProvider;
pub use table_provider::IcebergTableProviderFactory;
Expand Down
75 changes: 75 additions & 0 deletions iceberg/src/schema_provider.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,75 @@
use std::sync::Arc;

use dashmap::DashMap;
use datafusion::{
catalog::{SchemaProvider, TableProvider},
error::Result,
};
use futures::future::try_join_all;
use iceberg::{Catalog, NamespaceIdent, Runtime};

use crate::{IcebergCatalogTableProvider, common::df_err};

#[derive(Debug)]
pub struct IcebergSchemaProvider {
// Using Arc + DashMap for cheap clones of the tables in this schema
tables: Arc<DashMap<String, Arc<IcebergCatalogTableProvider>>>,
}

impl IcebergSchemaProvider {
pub async fn try_new(
catalog: Arc<dyn Catalog>,
namespace: NamespaceIdent,
iceberg_runtime: Runtime,
) -> Result<Self> {
let table_names: Vec<_> = catalog
.list_tables(&namespace)
.await
.map_err(df_err)?
.into_iter()
.map(|ident| ident.name().to_string())
.collect();

let table_providers = try_join_all(
table_names
.iter()
.map(|name| {
IcebergCatalogTableProvider::try_new(
catalog.clone(),
namespace.clone(),
name,
iceberg_runtime.clone(),
)
})
.collect::<Vec<_>>(),
)
.await?;

let tables = Arc::new(DashMap::new());

// Getting a map of: <table_name, table_provider>
for (name, provider) in table_names.into_iter().zip(table_providers) {
tables.insert(name, Arc::new(provider));
}

Ok(IcebergSchemaProvider { tables })
}
}

#[async_trait::async_trait]
impl SchemaProvider for IcebergSchemaProvider {
async fn table(&self, name: &str) -> Result<Option<Arc<dyn TableProvider>>> {
Ok(self
.tables
.get(name)
.map(|provider| provider.clone() as Arc<dyn TableProvider>))
}

fn table_names(&self) -> Vec<String> {
self.tables.iter().map(|k| k.key().clone()).collect()
}

fn table_exist(&self, name: &str) -> bool {
self.tables.contains_key(name)
}
}
85 changes: 79 additions & 6 deletions iceberg/src/test_utils/harness.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,19 +25,26 @@ use iceberg::io::{
FileMetadata, FileRead, FileWrite, InputFile, LocalFsStorage, OutputFile, Storage,
StorageConfig, StorageFactory,
};
use iceberg::memory::{MEMORY_CATALOG_WAREHOUSE, MemoryCatalogBuilder};
use iceberg::spec::TableMetadata;
use iceberg::{Error, ErrorKind, Result as IcebergResult};
use iceberg::{
Catalog, CatalogBuilder, Error, ErrorKind, NamespaceIdent, Result as IcebergResult, TableIdent,
};
use serde::{Deserialize, Serialize};

use super::taxi_metadata;
use crate::{IcebergExt, IcebergIntegrationOptions, iceberg_desired_task_count};
use crate::common::df_err;
use crate::{IcebergCatalog, IcebergExt, IcebergIntegrationOptions, iceberg_desired_task_count};

pub const FIXTURE_CATALOG: &str = "iceberg";
pub const FIXTURE_NAMESPACE: &str = "nyc";
pub const FIXTURE_URI: &str = "s3://iceberg-test/warehouse/taxi";
const WAREHOUSE_URI: &str = "s3://iceberg-test/warehouse/";
const FIXTURE_METADATA_URI: &str = "s3://iceberg-test/warehouse/taxi/metadata/v1.metadata.json";

pub struct IcebergTestHarness {
ctx: SessionContext,
iceberg_catalog: Option<Arc<dyn iceberg::Catalog>>,
}

impl IcebergTestHarness {
Expand All @@ -61,6 +68,14 @@ impl IcebergTestHarness {
self.ctx.sql(sql).await?.create_physical_plan().await
}

/// Returns the in-memory Iceberg catalog backing a fixture built with
/// [`IcebergTestHarnessBuilder::with_catalog`].
pub fn iceberg_catalog(&self) -> Result<Arc<dyn Catalog>> {
self.iceberg_catalog.clone().ok_or_else(|| {
DataFusionError::Plan("the fixture was not built with a catalog".to_string())
})
}

/// Returns the fixture scan without SQL optimization, including for an empty table.
pub async fn scan(&self) -> Result<Arc<dyn ExecutionPlan>> {
self.ctx
Expand Down Expand Up @@ -94,6 +109,7 @@ pub struct IcebergTestHarnessBuilder {
metadata: TableMetadata,
table_options: BTreeMap<String, String>,
files: HashMap<String, Vec<u8>>,
catalog: bool,
#[cfg(feature = "integration")]
workers: Option<usize>,
}
Expand All @@ -108,6 +124,7 @@ impl Default for IcebergTestHarnessBuilder {
metadata: taxi_metadata(),
table_options: BTreeMap::new(),
files: HashMap::new(),
catalog: false,
#[cfg(feature = "integration")]
workers: None,
}
Expand Down Expand Up @@ -144,6 +161,14 @@ impl IcebergTestHarnessBuilder {
self
}

/// Registers the fixture through an in-memory Iceberg catalog exposed as
/// [`FIXTURE_CATALOG`].[`FIXTURE_NAMESPACE`] instead of `CREATE EXTERNAL TABLE`.
/// Table options are not applied in this mode.
pub fn with_catalog(mut self) -> Self {
self.catalog = true;
self
}

/// Enables distributed planning with logical workers backed by in-memory gRPC.
#[cfg(feature = "integration")]
pub fn with_workers(mut self, workers: usize) -> Self {
Expand All @@ -162,12 +187,19 @@ impl IcebergTestHarnessBuilder {
..FixtureStorageFactory::default()
};
let options = IcebergIntegrationOptions {
storage_factory: Arc::new(storage_factory),
storage_factory: Arc::new(storage_factory.clone()),
iceberg_runtime: iceberg::Runtime::current(),
};
let state = self
let iceberg_runtime = options.iceberg_runtime.clone();
let mut state = self
.session_builder
.with_iceberg_integration(options.clone());
if self.catalog {
// Makes `taxi` resolve through the registered Iceberg catalog.
let config = state.config().get_or_insert_default();
*config = std::mem::take(config)
.with_default_catalog_and_schema(FIXTURE_CATALOG, FIXTURE_NAMESPACE);
}
#[cfg(feature = "integration")]
let state = if let Some(workers) = self.workers {
let resolver =
Expand All @@ -186,6 +218,11 @@ impl IcebergTestHarnessBuilder {
state
};
let ctx = SessionContext::new_with_state(state.build());
if self.catalog && !self.table_options.is_empty() {
return Err(DataFusionError::Plan(
"table options are not supported with a catalog-backed fixture".to_string(),
));
}
let mut statement = format!(
"CREATE EXTERNAL TABLE taxi STORED AS ICEBERG \
LOCATION '{FIXTURE_METADATA_URI}'"
Expand All @@ -205,8 +242,44 @@ impl IcebergTestHarnessBuilder {
.join(", ");
statement.push_str(&format!(" OPTIONS ({options})"));
}
ctx.sql(&statement).await?.collect().await?;
Ok(IcebergTestHarness { ctx })

let mut iceberg_catalog = None;
if self.catalog {
let catalog = MemoryCatalogBuilder::default()
.with_storage_factory(Arc::new(storage_factory))
.with_runtime(iceberg_runtime.clone())
.load(
"memory",
HashMap::from([(
MEMORY_CATALOG_WAREHOUSE.to_string(),
WAREHOUSE_URI.to_string(),
)]),
)
.await
.map_err(df_err)?;
let namespace = NamespaceIdent::new(FIXTURE_NAMESPACE.to_string());
catalog
.create_namespace(&namespace, HashMap::new())
.await
.map_err(df_err)?;
catalog
.register_table(
&TableIdent::new(namespace, "taxi".to_string()),
FIXTURE_METADATA_URI.to_string(),
)
.await
.map_err(df_err)?;
let catalog: Arc<dyn Catalog> = Arc::new(catalog);
let provider = IcebergCatalog::try_new(catalog.clone(), iceberg_runtime).await?;
ctx.register_catalog(FIXTURE_CATALOG, Arc::new(provider));
iceberg_catalog = Some(catalog);
} else {
ctx.sql(&statement).await?.collect().await?;
}
Ok(IcebergTestHarness {
ctx,
iceberg_catalog,
})
}
}

Expand Down
Loading
Loading