use std::collections::HashMap;
use std::sync::Arc;
use iceberg::spec::{NestedFieldRef, Transform, Type};
use iceberg::{Catalog as _, CatalogBuilder as _, TableIdent};
use iceberg_catalog_rest::{RestCatalog, RestCatalogBuilder};
use iceberg_storage_opendal::OpenDalStorageFactory;
use tracing::{debug, info};
use super::partition::{PartitionFieldInfo, PartitionSpec};
use crate::types::arrow::FieldIdMap;
use crate::types::config::IcebergConfig;
use crate::types::constants::{DATA_ROW_EXPECTED_COLUMNS, METADATA_ROW_EXPECTED_COLUMNS};
use crate::types::error::{IcebergError, Result};
fn expected_columns_for_table(table_name: &str) -> Option<&'static [&'static str]> {
if table_name.ends_with("metadata_row") {
Some(METADATA_ROW_EXPECTED_COLUMNS)
} else if table_name.ends_with("data_row") {
Some(DATA_ROW_EXPECTED_COLUMNS)
} else {
None
}
}
pub async fn discover_field_ids(config: &IcebergConfig, table_name: &str) -> Result<FieldIdMap> {
info!(
catalog = config.catalog_uri.as_str(),
table = table_name,
"discovering field IDs"
);
let catalog = build_catalog(config).await?;
let table_ident = parse_table_ident(table_name)?;
let table =
catalog
.load_table(&table_ident)
.await
.map_err(|source| IcebergError::TableLoad {
table_name: table_name.to_owned(),
source: Box::new(source),
})?;
let schema = table.metadata().current_schema();
let mut mapping = HashMap::new();
for field in schema.as_struct().fields() {
collect_field_ids(field, &mut mapping, "");
}
debug!(
table = table_name,
fields = mapping.len(),
"discovered field IDs"
);
let field_id_map = FieldIdMap::new(table_name.to_owned(), mapping);
if let Some(expected) = expected_columns_for_table(table_name) {
field_id_map.validate_columns(expected)?;
}
Ok(field_id_map)
}
pub async fn discover_partition_spec(
config: &IcebergConfig,
table_name: &str,
) -> Result<PartitionSpec> {
info!(
catalog = config.catalog_uri.as_str(),
table = table_name,
"discovering partition spec"
);
let catalog = build_catalog(config).await?;
let table_ident = parse_table_ident(table_name)?;
let table =
catalog
.load_table(&table_ident)
.await
.map_err(|source| IcebergError::TableLoad {
table_name: table_name.to_owned(),
source: Box::new(source),
})?;
let schema = table.metadata().current_schema();
let spec = table.metadata().default_partition_spec();
let mut fields = Vec::new();
for partition_field in spec.fields() {
let source_field = schema
.field_by_id(partition_field.source_id)
.ok_or_else(|| IcebergError::PartitionFieldNotFound {
field_id: partition_field.source_id,
table_name: table_name.to_owned(),
})?;
let transform_name = transform_to_string(partition_field.transform);
fields.push(PartitionFieldInfo::new(
&source_field.name,
&partition_field.name,
&transform_name,
));
}
debug!(
table = table_name,
partition_fields = fields.len(),
"discovered partition spec"
);
Ok(PartitionSpec::new(table_name, fields))
}
pub async fn discover_table_metadata(
config: &IcebergConfig,
table_name: &str,
) -> Result<(FieldIdMap, PartitionSpec)> {
info!(
catalog = config.catalog_uri.as_str(),
table = table_name,
"discovering table metadata"
);
let catalog = build_catalog(config).await?;
let table_ident = parse_table_ident(table_name)?;
let table =
catalog
.load_table(&table_ident)
.await
.map_err(|source| IcebergError::TableLoad {
table_name: table_name.to_owned(),
source: Box::new(source),
})?;
let schema = table.metadata().current_schema();
let mut mapping = HashMap::new();
for field in schema.as_struct().fields() {
collect_field_ids(field, &mut mapping, "");
}
let field_id_map = FieldIdMap::new(table_name.to_owned(), mapping);
if let Some(expected) = expected_columns_for_table(table_name) {
field_id_map.validate_columns(expected)?;
}
let spec = table.metadata().default_partition_spec();
let mut fields = Vec::new();
for partition_field in spec.fields() {
let source_field = schema
.field_by_id(partition_field.source_id)
.ok_or_else(|| IcebergError::PartitionFieldNotFound {
field_id: partition_field.source_id,
table_name: table_name.to_owned(),
})?;
fields.push(PartitionFieldInfo::new(
&source_field.name,
&partition_field.name,
&transform_to_string(partition_field.transform),
));
}
let partition_spec = PartitionSpec::new(table_name, fields);
debug!(
table = table_name,
field_ids = field_id_map.len(),
partition_fields = partition_spec.fields().len(),
"discovered table metadata"
);
Ok((field_id_map, partition_spec))
}
async fn build_catalog(config: &IcebergConfig) -> Result<RestCatalog> {
let mut props: HashMap<String, String> = HashMap::new();
props.insert("uri".to_owned(), config.catalog_uri.clone());
if let Some(warehouse) = &config.catalog_warehouse {
props.insert("warehouse".to_owned(), warehouse.clone());
}
if let Some(credential) = &config.credential {
props.insert("credential".to_owned(), credential.clone());
}
if let Some(oauth2_server_uri) = &config.oauth2_server_uri {
props.insert("oauth2-server-uri".to_owned(), oauth2_server_uri.clone());
}
if let Some(scope) = &config.scope {
props.insert("scope".to_owned(), scope.clone());
}
let storage_factory = Arc::new(OpenDalStorageFactory::S3 {
configured_scheme: "s3".to_owned(),
customized_credential_load: None,
});
RestCatalogBuilder::default()
.with_storage_factory(storage_factory)
.load("rest", props)
.await
.map_err(|source| {
IcebergError::CatalogConnection {
catalog_uri: config.catalog_uri.clone(),
source: Box::new(source),
}
.into()
})
}
fn collect_field_ids(field: &NestedFieldRef, mapping: &mut HashMap<String, i32>, path: &str) {
let full_path = if path.is_empty() {
field.name.clone()
} else {
format!("{path}.{}", field.name)
};
mapping.insert(full_path.clone(), field.id);
collect_nested_type_ids(&field.field_type, mapping, &full_path);
}
fn collect_nested_type_ids(field_type: &Type, mapping: &mut HashMap<String, i32>, path: &str) {
match field_type {
Type::List(list_type) => {
let element_path = format!("{path}.element");
mapping.insert(element_path.clone(), list_type.element_field.id);
collect_nested_type_ids(&list_type.element_field.field_type, mapping, &element_path);
}
Type::Struct(struct_type) => {
for nested_field in struct_type.fields() {
collect_field_ids(nested_field, mapping, path);
}
}
Type::Map(map_type) => {
let key_path = format!("{path}.key");
let value_path = format!("{path}.value");
mapping.insert(key_path.clone(), map_type.key_field.id);
mapping.insert(value_path.clone(), map_type.value_field.id);
collect_nested_type_ids(&map_type.key_field.field_type, mapping, &key_path);
collect_nested_type_ids(&map_type.value_field.field_type, mapping, &value_path);
}
Type::Primitive(_) => {}
}
}
fn parse_table_ident(table_name: &str) -> Result<TableIdent> {
TableIdent::from_strs(table_name.split('.')).map_err(|source| {
IcebergError::InvalidTableName {
table_name: table_name.to_owned(),
source: Box::new(source),
}
.into()
})
}
fn transform_to_string(transform: Transform) -> String {
match transform {
Transform::Identity => "identity".to_owned(),
Transform::Day => "day".to_owned(),
Transform::Hour => "hour".to_owned(),
Transform::Month => "month".to_owned(),
Transform::Year => "year".to_owned(),
Transform::Bucket(n) => format!("bucket[{n}]"),
Transform::Truncate(n) => format!("truncate[{n}]"),
Transform::Void => "void".to_owned(),
Transform::Unknown => "unknown".to_owned(),
}
}