This is an automated email from the ASF dual-hosted git repository.
mbutrovich pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/datafusion-iceberg.git
The following commit(s) were added to refs/heads/main by this push:
new 697c1da refactor: use DataFusion errors in integration (#7)
697c1da is described below
commit 697c1da74adc3645419267880512ac91ad601ffa
Author: Gabriel <[email protected]>
AuthorDate: Wed Sep 23 22:58:54 2026 +0200
refactor: use DataFusion errors in integration (#7)
* refactor: use DataFusion errors in integration
* test: update DataFusion error expectations
---------
Co-authored-by: Matt Butrovich <[email protected]>
---
Cargo.lock | 1 -
crates/datafusion/Cargo.toml | 1 -
crates/datafusion/src/catalog.rs | 7 +-
crates/datafusion/src/error.rs | 105 +++++++++++++++++++--
crates/datafusion/src/physical_plan/commit.rs | 44 ++++-----
crates/datafusion/src/physical_plan/project.rs | 17 ++--
crates/datafusion/src/physical_plan/repartition.rs | 11 ++-
crates/datafusion/src/physical_plan/scan.rs | 14 +--
crates/datafusion/src/physical_plan/sort.rs | 18 ++--
crates/datafusion/src/physical_plan/write.rs | 60 ++++++------
crates/datafusion/src/schema.rs | 72 +++++---------
crates/datafusion/src/table/metadata_table.rs | 6 +-
crates/datafusion/src/table/mod.rs | 89 ++++++++---------
.../datafusion/src/table/table_provider_factory.rs | 38 ++++----
crates/datafusion/src/task_writer.rs | 40 +++++---
.../tests/integration_datafusion_test.rs | 27 +++---
.../slts/df_test/timestamp_predicate_pushdown.slt | 4 +-
17 files changed, 315 insertions(+), 239 deletions(-)
diff --git a/Cargo.lock b/Cargo.lock
index 0191132..60a2093 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -2111,7 +2111,6 @@ dependencies = [
name = "datafusion-iceberg"
version = "0.10.1"
dependencies = [
- "anyhow",
"async-trait",
"dashmap",
"datafusion",
diff --git a/crates/datafusion/Cargo.toml b/crates/datafusion/Cargo.toml
index 9647aed..b4c3b92 100644
--- a/crates/datafusion/Cargo.toml
+++ b/crates/datafusion/Cargo.toml
@@ -29,7 +29,6 @@ license = { workspace = true }
repository = { workspace = true }
[dependencies]
-anyhow = { workspace = true }
async-trait = { workspace = true }
dashmap = "6"
datafusion = { workspace = true }
diff --git a/crates/datafusion/src/catalog.rs b/crates/datafusion/src/catalog.rs
index 2c6e1ff..a96f731 100644
--- a/crates/datafusion/src/catalog.rs
+++ b/crates/datafusion/src/catalog.rs
@@ -19,10 +19,12 @@ use std::collections::HashMap;
use std::sync::Arc;
use datafusion::catalog::{CatalogProvider, SchemaProvider};
+use datafusion::error::Result;
use futures::future::try_join_all;
-use iceberg::{Catalog, NamespaceIdent, Result};
+use iceberg::{Catalog, NamespaceIdent};
use crate::schema::IcebergSchemaProvider;
+use crate::to_datafusion_error;
/// Provides an interface to manage and access multiple schemas
/// within an Iceberg [`Catalog`].
@@ -51,7 +53,8 @@ impl IcebergCatalogProvider {
// As of right now; schemas might become stale.
let schema_names: Vec<_> = client
.list_namespaces(None)
- .await?
+ .await
+ .map_err(to_datafusion_error)?
.iter()
.flat_map(|ns| ns.as_ref().clone())
.collect();
diff --git a/crates/datafusion/src/error.rs b/crates/datafusion/src/error.rs
index 5bc4701..9733521 100644
--- a/crates/datafusion/src/error.rs
+++ b/crates/datafusion/src/error.rs
@@ -15,18 +15,109 @@
// specific language governing permissions and limitations
// under the License.
-use anyhow::anyhow;
+use datafusion::error::DataFusionError;
use iceberg::{Error, ErrorKind};
/// Converts a datafusion error into an iceberg error.
-pub fn from_datafusion_error(error: datafusion::error::DataFusionError) ->
Error {
+pub fn from_datafusion_error(error: DataFusionError) -> Error {
+ let fallback_message = error.to_string();
+ let DataFusionError::Context(context, error) = error else {
+ return unexpected_datafusion_error(fallback_message);
+ };
+
+ let Some(kind) = parse_iceberg_error_kind(&context) else {
+ return unexpected_datafusion_error(fallback_message);
+ };
+ let DataFusionError::Execution(message) = error.as_ref() else {
+ return unexpected_datafusion_error(fallback_message);
+ };
+
+ Error::new(kind, strip_error_kind(kind, message))
+}
+
+/// Converts an iceberg error into a datafusion error.
+pub fn to_datafusion_error(error: Error) -> DataFusionError {
+ DataFusionError::Context(
+ format!("IcebergError({})", error.kind()),
+ Box::new(DataFusionError::Execution(error.to_string())),
+ )
+}
+
+fn unexpected_datafusion_error(message: String) -> Error {
Error::new(
ErrorKind::Unexpected,
- "Operation failed for hitting datafusion error".to_string(),
+ format!("DataFusion execution failed: {message}"),
)
- .with_source(anyhow!("datafusion error: {error:?}"))
}
-/// Converts an iceberg error into a datafusion error.
-pub fn to_datafusion_error(error: Error) -> datafusion::error::DataFusionError
{
- datafusion::error::DataFusionError::External(error.into())
+
+fn parse_iceberg_error_kind(context: &str) -> Option<ErrorKind> {
+ let kind = context.strip_prefix("IcebergError(")?.strip_suffix(')')?;
+
+ match kind {
+ "PreconditionFailed" => Some(ErrorKind::PreconditionFailed),
+ "Unexpected" => Some(ErrorKind::Unexpected),
+ "DataInvalid" => Some(ErrorKind::DataInvalid),
+ "NamespaceAlreadyExists" => Some(ErrorKind::NamespaceAlreadyExists),
+ "TableAlreadyExists" => Some(ErrorKind::TableAlreadyExists),
+ "NamespaceNotFound" => Some(ErrorKind::NamespaceNotFound),
+ "TableNotFound" => Some(ErrorKind::TableNotFound),
+ "FeatureUnsupported" => Some(ErrorKind::FeatureUnsupported),
+ "CatalogCommitConflicts" => Some(ErrorKind::CatalogCommitConflicts),
+ _ => None,
+ }
+}
+
+fn strip_error_kind(kind: ErrorKind, message: &str) -> String {
+ let kind = kind.into_static();
+ if message == kind {
+ String::new()
+ } else {
+ message
+ .strip_prefix(&format!("{kind} => "))
+ .unwrap_or(message)
+ .to_string()
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use datafusion::error::DataFusionError;
+ use iceberg::{Error, ErrorKind};
+
+ use super::{from_datafusion_error, to_datafusion_error};
+
+ #[test]
+ fn roundtrips_iceberg_error_kind_and_message() {
+ let error = Error::new(ErrorKind::DataInvalid, "invalid manifest");
+ let roundtripped = from_datafusion_error(to_datafusion_error(error));
+
+ assert_eq!(roundtripped.kind(), ErrorKind::DataInvalid);
+ assert_eq!(roundtripped.to_string(), "DataInvalid => invalid
manifest");
+ }
+
+ #[test]
+ fn encodes_iceberg_errors_with_native_datafusion_variants() {
+ let error =
+ to_datafusion_error(Error::new(ErrorKind::DataInvalid, "invalid
manifest"));
+
+ assert!(matches!(
+ error,
+ DataFusionError::Context(context, inner)
+ if context == "IcebergError(DataInvalid)"
+ && matches!(inner.as_ref(),
DataFusionError::Execution(message) if message == "DataInvalid => invalid
manifest")
+ ));
+ }
+
+ #[test]
+ fn maps_non_iceberg_datafusion_errors_to_unexpected() {
+ let error = from_datafusion_error(DataFusionError::Execution(
+ "worker failed".to_string(),
+ ));
+
+ assert_eq!(error.kind(), ErrorKind::Unexpected);
+ assert_eq!(
+ error.to_string(),
+ "Unexpected => DataFusion execution failed: Execution error:
worker failed"
+ );
+ }
}
diff --git a/crates/datafusion/src/physical_plan/commit.rs
b/crates/datafusion/src/physical_plan/commit.rs
index d080f61..dcbee3a 100644
--- a/crates/datafusion/src/physical_plan/commit.rs
+++ b/crates/datafusion/src/physical_plan/commit.rs
@@ -22,8 +22,10 @@ use datafusion::arrow::array::{Array, ArrayRef, RecordBatch,
StringArray, UInt64
use datafusion::arrow::datatypes::{
DataType, Field, Schema as ArrowSchema, SchemaRef as ArrowSchemaRef,
};
-use datafusion::common::tree_node::TreeNodeRecursion;
-use datafusion::common::{DataFusionError, Result as DFResult};
+use datafusion::common::{
+ internal_datafusion_err, internal_err, tree_node::TreeNodeRecursion,
+};
+use datafusion::error::{DataFusionError, Result};
use datafusion::execution::{SendableRecordBatchStream, TaskContext};
use datafusion::physical_expr::{EquivalenceProperties, Partitioning,
PhysicalExpr};
use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType};
@@ -84,7 +86,7 @@ impl IcebergCommitExec {
}
// Create a record batch with just the count of rows written
- fn make_count_batch(count: u64) -> DFResult<RecordBatch> {
+ fn make_count_batch(count: u64) -> Result<RecordBatch> {
let count_array = Arc::new(UInt64Array::from(vec![count])) as ArrayRef;
RecordBatch::try_from_iter_with_nullable(vec![("count", count_array,
false)])
@@ -142,8 +144,8 @@ impl ExecutionPlan for IcebergCommitExec {
fn apply_expressions(
&self,
- _f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) ->
DFResult<TreeNodeRecursion>,
- ) -> DFResult<TreeNodeRecursion> {
+ _f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) ->
Result<TreeNodeRecursion>,
+ ) -> Result<TreeNodeRecursion> {
Ok(TreeNodeRecursion::Continue)
}
@@ -163,12 +165,12 @@ impl ExecutionPlan for IcebergCommitExec {
fn with_new_children(
self: Arc<Self>,
children: Vec<Arc<dyn ExecutionPlan>>,
- ) -> DFResult<Arc<dyn ExecutionPlan>> {
+ ) -> Result<Arc<dyn ExecutionPlan>> {
if children.len() != 1 {
- return Err(DataFusionError::Internal(format!(
+ return internal_err!(
"IcebergCommitExec expects exactly one child, but provided {}",
children.len()
- )));
+ );
}
Ok(Arc::new(IcebergCommitExec::new(
@@ -183,12 +185,12 @@ impl ExecutionPlan for IcebergCommitExec {
&self,
partition: usize,
context: Arc<TaskContext>,
- ) -> DFResult<SendableRecordBatchStream> {
+ ) -> Result<SendableRecordBatchStream> {
// IcebergCommitExec only has one partition (partition 0)
if partition != 0 {
- return Err(DataFusionError::Internal(format!(
+ return internal_err!(
"IcebergCommitExec only has one partition, but got partition
{partition}"
- )));
+ );
}
let table = self.table.clone();
@@ -215,15 +217,15 @@ impl ExecutionPlan for IcebergCommitExec {
let files_array = batch
.column_by_name(DATA_FILES_COL_NAME)
.ok_or_else(|| {
- DataFusionError::Internal(
- "Expected 'data_files' column in input
batch".to_string(),
+ internal_datafusion_err!(
+ "Expected 'data_files' column in input batch"
)
})?
.as_any()
.downcast_ref::<StringArray>()
.ok_or_else(|| {
- DataFusionError::Internal(
- "Expected 'data_files' column to be
StringArray".to_string(),
+ internal_datafusion_err!(
+ "Expected 'data_files' column to be StringArray"
)
})?;
@@ -231,7 +233,7 @@ impl ExecutionPlan for IcebergCommitExec {
let batch_files: Vec<DataFile> = files_array
.into_iter()
.flatten()
- .map(|f| -> DFResult<DataFile> {
+ .map(|f| -> Result<DataFile> {
// Parse JSON to DataFileSerde and convert to DataFile
deserialize_data_file_from_json(
f,
@@ -241,7 +243,7 @@ impl ExecutionPlan for IcebergCommitExec {
)
.map_err(to_datafusion_error)
})
- .collect::<datafusion::common::Result<_>>()?;
+ .collect::<Result<_>>()?;
// add record_counts from the current batch to total record
count
total_record_count +=
@@ -361,15 +363,15 @@ mod tests {
fn apply_expressions(
&self,
- _f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) ->
DFResult<TreeNodeRecursion>,
- ) -> DFResult<TreeNodeRecursion> {
+ _f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) ->
Result<TreeNodeRecursion>,
+ ) -> Result<TreeNodeRecursion> {
Ok(TreeNodeRecursion::Continue)
}
fn with_new_children(
self: Arc<Self>,
_children: Vec<Arc<dyn ExecutionPlan>>,
- ) -> datafusion::common::Result<Arc<dyn ExecutionPlan>> {
+ ) -> Result<Arc<dyn ExecutionPlan>> {
Ok(self)
}
@@ -377,7 +379,7 @@ mod tests {
&self,
_partition: usize,
_context: Arc<TaskContext>,
- ) -> datafusion::common::Result<SendableRecordBatchStream> {
+ ) -> Result<SendableRecordBatchStream> {
// Create a record batch with the serialized data files
let array =
Arc::new(StringArray::from(self.data_files_json.clone())) as
ArrayRef;
diff --git a/crates/datafusion/src/physical_plan/project.rs
b/crates/datafusion/src/physical_plan/project.rs
index d8ab9b3..f127d8b 100644
--- a/crates/datafusion/src/physical_plan/project.rs
+++ b/crates/datafusion/src/physical_plan/project.rs
@@ -21,7 +21,8 @@ use std::sync::Arc;
use datafusion::arrow::array::RecordBatch;
use datafusion::arrow::datatypes::{DataType, Schema as ArrowSchema};
-use datafusion::common::{DataFusionError, Result as DFResult};
+use datafusion::common::plan_err;
+use datafusion::error::Result;
use datafusion::physical_expr::PhysicalExpr;
use datafusion::physical_expr::expressions::Column;
use datafusion::physical_plan::projection::ProjectionExec;
@@ -51,7 +52,7 @@ use crate::to_datafusion_error;
pub fn project_with_partition(
input: Arc<dyn ExecutionPlan>,
table: &Table,
-) -> DFResult<Arc<dyn ExecutionPlan>> {
+) -> Result<Arc<dyn ExecutionPlan>> {
let metadata = table.metadata();
let partition_spec = metadata.default_partition_spec();
let table_schema = metadata.current_schema();
@@ -72,11 +73,11 @@ pub fn project_with_partition(
.map_err(to_datafusion_error)?;
if input_schema_cleaned != expected_schema_cleaned {
- return Err(DataFusionError::Plan(format!(
+ return plan_err!(
"Input schema does not match Iceberg table schema.\n\
Expected schema: {expected_schema_cleaned}\n\
Input schema: {input_schema_cleaned}"
- )));
+ );
}
let calculator =
@@ -129,15 +130,15 @@ impl PartialEq for PartitionExpr {
impl Eq for PartitionExpr {}
impl PhysicalExpr for PartitionExpr {
- fn data_type(&self, _input_schema: &ArrowSchema) -> DFResult<DataType> {
+ fn data_type(&self, _input_schema: &ArrowSchema) -> Result<DataType> {
Ok(self.calculator.partition_arrow_type().clone())
}
- fn nullable(&self, _input_schema: &ArrowSchema) -> DFResult<bool> {
+ fn nullable(&self, _input_schema: &ArrowSchema) -> Result<bool> {
Ok(false)
}
- fn evaluate(&self, batch: &RecordBatch) -> DFResult<ColumnarValue> {
+ fn evaluate(&self, batch: &RecordBatch) -> Result<ColumnarValue> {
let array = self
.calculator
.calculate(batch)
@@ -152,7 +153,7 @@ impl PhysicalExpr for PartitionExpr {
fn with_new_children(
self: Arc<Self>,
_children: Vec<Arc<dyn PhysicalExpr>>,
- ) -> DFResult<Arc<dyn PhysicalExpr>> {
+ ) -> Result<Arc<dyn PhysicalExpr>> {
Ok(self)
}
diff --git a/crates/datafusion/src/physical_plan/repartition.rs
b/crates/datafusion/src/physical_plan/repartition.rs
index e6cbc13..0dbda0c 100644
--- a/crates/datafusion/src/physical_plan/repartition.rs
+++ b/crates/datafusion/src/physical_plan/repartition.rs
@@ -18,7 +18,8 @@
use std::num::NonZeroUsize;
use std::sync::Arc;
-use datafusion::error::{DataFusionError, Result as DFResult};
+use datafusion::common::plan_err;
+use datafusion::error::Result;
use datafusion::physical_expr::PhysicalExpr;
use datafusion::physical_plan::expressions::Column;
use datafusion::physical_plan::repartition::RepartitionExec;
@@ -90,7 +91,7 @@ pub(crate) fn repartition(
input: Arc<dyn ExecutionPlan>,
table_metadata: TableMetadataRef,
target_partitions: NonZeroUsize,
-) -> DFResult<Arc<dyn ExecutionPlan>> {
+) -> Result<Arc<dyn ExecutionPlan>> {
let partitioning_strategy =
determine_partitioning_strategy(&input, &table_metadata,
target_partitions)?;
@@ -125,7 +126,7 @@ fn determine_partitioning_strategy(
input: &Arc<dyn ExecutionPlan>,
table_metadata: &TableMetadata,
target_partitions: NonZeroUsize,
-) -> DFResult<Partitioning> {
+) -> Result<Partitioning> {
let partition_spec = table_metadata.default_partition_spec();
let input_schema = input.schema();
let target_partition_count = target_partitions.get();
@@ -158,10 +159,10 @@ fn determine_partitioning_strategy(
}
// Case 2: Partitioned table missing _partition column (normally this
should not happen)
- (true, Err(_)) => Err(DataFusionError::Plan(format!(
+ (true, Err(_)) => plan_err!(
"Partitioned table input missing
{PROJECTED_PARTITION_VALUE_COLUMN} column. \
Ensure projection happens before repartitioning."
- ))),
+ ),
// Case 3: Unpartitioned table, always use RoundRobinBatch
(false, _) =>
Ok(Partitioning::RoundRobinBatch(target_partition_count)),
diff --git a/crates/datafusion/src/physical_plan/scan.rs
b/crates/datafusion/src/physical_plan/scan.rs
index 1796304..b9b15c5 100644
--- a/crates/datafusion/src/physical_plan/scan.rs
+++ b/crates/datafusion/src/physical_plan/scan.rs
@@ -22,7 +22,7 @@ use std::vec;
use datafusion::arrow::array::RecordBatch;
use datafusion::arrow::datatypes::SchemaRef as ArrowSchemaRef;
use datafusion::common::tree_node::TreeNodeRecursion;
-use datafusion::error::Result as DFResult;
+use datafusion::error::Result;
use datafusion::execution::{SendableRecordBatchStream, TaskContext};
use datafusion::physical_expr::{EquivalenceProperties, PhysicalExpr};
use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType};
@@ -128,15 +128,15 @@ impl ExecutionPlan for IcebergTableScan {
fn apply_expressions(
&self,
- _f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) ->
DFResult<TreeNodeRecursion>,
- ) -> DFResult<TreeNodeRecursion> {
+ _f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) ->
Result<TreeNodeRecursion>,
+ ) -> Result<TreeNodeRecursion> {
Ok(TreeNodeRecursion::Continue)
}
fn with_new_children(
self: Arc<Self>,
_children: Vec<Arc<dyn ExecutionPlan>>,
- ) -> DFResult<Arc<dyn ExecutionPlan>> {
+ ) -> Result<Arc<dyn ExecutionPlan>> {
Ok(self)
}
@@ -148,7 +148,7 @@ impl ExecutionPlan for IcebergTableScan {
&self,
_partition: usize,
_context: Arc<TaskContext>,
- ) -> DFResult<SendableRecordBatchStream> {
+ ) -> Result<SendableRecordBatchStream> {
let fut = get_batch_stream(
self.table.clone(),
self.snapshot_id,
@@ -158,7 +158,7 @@ impl ExecutionPlan for IcebergTableScan {
let stream = futures::stream::once(fut).try_flatten();
// Apply limit if specified
- let limited_stream: Pin<Box<dyn Stream<Item = DFResult<RecordBatch>> +
Send>> =
+ let limited_stream: Pin<Box<dyn Stream<Item = Result<RecordBatch>> +
Send>> =
if let Some(limit) = self.limit {
let mut remaining = limit;
Box::pin(stream.try_filter_map(move |batch| {
@@ -217,7 +217,7 @@ async fn get_batch_stream(
snapshot_id: Option<i64>,
column_names: Option<Vec<String>>,
predicates: Option<Predicate>,
-) -> DFResult<Pin<Box<dyn Stream<Item = DFResult<RecordBatch>> + Send>>> {
+) -> Result<Pin<Box<dyn Stream<Item = Result<RecordBatch>> + Send>>> {
let scan_builder = match snapshot_id {
Some(snapshot_id) => table.scan().snapshot_id(snapshot_id),
None => table.scan(),
diff --git a/crates/datafusion/src/physical_plan/sort.rs
b/crates/datafusion/src/physical_plan/sort.rs
index fce9210..1993e08 100644
--- a/crates/datafusion/src/physical_plan/sort.rs
+++ b/crates/datafusion/src/physical_plan/sort.rs
@@ -20,8 +20,8 @@
use std::sync::Arc;
use datafusion::arrow::compute::SortOptions;
-use datafusion::common::Result as DFResult;
-use datafusion::error::DataFusionError;
+use datafusion::common::plan_datafusion_err;
+use datafusion::error::Result;
use datafusion::physical_expr::{LexOrdering, PhysicalSortExpr};
use datafusion::physical_plan::ExecutionPlan;
use datafusion::physical_plan::expressions::Column;
@@ -44,17 +44,15 @@ use iceberg::arrow::PROJECTED_PARTITION_VALUE_COLUMN;
/// * `Err` - If the partition column is not found
pub(crate) fn sort_by_partition(
input: Arc<dyn ExecutionPlan>,
-) -> DFResult<Arc<dyn ExecutionPlan>> {
+) -> Result<Arc<dyn ExecutionPlan>> {
let schema = input.schema();
// Find the partition column in the schema
let (partition_column_index, _partition_field) = schema
.column_with_name(PROJECTED_PARTITION_VALUE_COLUMN)
- .ok_or_else(|| {
- DataFusionError::Plan(format!(
- "Partition column '{PROJECTED_PARTITION_VALUE_COLUMN}' not
found in schema. Ensure the plan has been extended with partition values using
project_with_partition."
- ))
- })?;
+ .ok_or_else(|| plan_datafusion_err!(
+ "Partition column '{PROJECTED_PARTITION_VALUE_COLUMN}' not found
in schema. Ensure the plan has been extended with partition values using
project_with_partition."
+ ))?;
// Create a single sort expression for the partition column
let column_expr = Arc::new(Column::new(
@@ -70,9 +68,7 @@ pub(crate) fn sort_by_partition(
// Create a SortExec with preserve_partitioning=true to ensure the output
partitioning
// is the same as the input partitioning, and the data is sorted within
each partition
let lex_ordering = LexOrdering::new(vec![sort_expr]).ok_or_else(|| {
- DataFusionError::Plan(
- "Failed to create LexOrdering from sort expression".to_string(),
- )
+ plan_datafusion_err!("Failed to create LexOrdering from sort
expression")
})?;
let sort_exec = SortExec::new(lex_ordering,
input).with_preserve_partitioning(true);
diff --git a/crates/datafusion/src/physical_plan/write.rs
b/crates/datafusion/src/physical_plan/write.rs
index c5a4206..96c0ac0 100644
--- a/crates/datafusion/src/physical_plan/write.rs
+++ b/crates/datafusion/src/physical_plan/write.rs
@@ -23,9 +23,9 @@ use datafusion::arrow::array::{ArrayRef, RecordBatch,
StringArray};
use datafusion::arrow::datatypes::{
DataType, Field, Schema as ArrowSchema, SchemaRef as ArrowSchemaRef,
};
-use datafusion::common::Result as DFResult;
-use datafusion::common::tree_node::TreeNodeRecursion;
+use datafusion::common::{internal_err, not_impl_err,
tree_node::TreeNodeRecursion};
use datafusion::error::DataFusionError;
+use datafusion::error::Result;
use datafusion::execution::{SendableRecordBatchStream, TaskContext};
use datafusion::physical_expr::{EquivalenceProperties, Partitioning,
PhysicalExpr};
use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType};
@@ -44,7 +44,6 @@ use iceberg::writer::file_writer::location_generator::{
DefaultFileNameGenerator, DefaultLocationGenerator,
};
use iceberg::writer::file_writer::rolling_writer::RollingFileWriterBuilder;
-use iceberg::{Error, ErrorKind};
use uuid::Uuid;
use crate::physical_plan::DATA_FILES_COL_NAME;
@@ -94,7 +93,7 @@ impl IcebergWriteExec {
}
// Create a record batch with serialized data files
- fn make_result_batch(data_files: Vec<String>) -> DFResult<RecordBatch> {
+ fn make_result_batch(data_files: Vec<String>) -> Result<RecordBatch> {
let files_array = Arc::new(StringArray::from(data_files)) as ArrayRef;
RecordBatch::try_new(Self::make_result_schema(),
vec![files_array]).map_err(|e| {
@@ -161,20 +160,20 @@ impl ExecutionPlan for IcebergWriteExec {
fn apply_expressions(
&self,
- _f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) ->
DFResult<TreeNodeRecursion>,
- ) -> DFResult<TreeNodeRecursion> {
+ _f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) ->
Result<TreeNodeRecursion>,
+ ) -> Result<TreeNodeRecursion> {
Ok(TreeNodeRecursion::Continue)
}
fn with_new_children(
self: Arc<Self>,
children: Vec<Arc<dyn ExecutionPlan>>,
- ) -> DFResult<Arc<dyn ExecutionPlan>> {
+ ) -> Result<Arc<dyn ExecutionPlan>> {
if children.len() != 1 {
- return Err(DataFusionError::Internal(format!(
+ return internal_err!(
"IcebergWriteExec expects exactly one child, but provided {}",
children.len()
- )));
+ );
}
Ok(Arc::new(Self::new(
@@ -208,7 +207,7 @@ impl ExecutionPlan for IcebergWriteExec {
&self,
partition: usize,
context: Arc<TaskContext>,
- ) -> DFResult<SendableRecordBatchStream> {
+ ) -> Result<SendableRecordBatchStream> {
let partition_type =
self.table.metadata().default_partition_type().clone();
let format_version = self.table.metadata().format_version();
@@ -222,12 +221,9 @@ impl ExecutionPlan for IcebergWriteExec {
let file_format = DataFileFormat::from_str(&write_format_default)
.map_err(to_datafusion_error)?;
if file_format != DataFileFormat::Parquet {
- return Err(to_datafusion_error(Error::new(
- ErrorKind::FeatureUnsupported,
- format!(
- "File format {file_format} is not supported for
insert_into yet!"
- ),
- )));
+ return not_impl_err!(
+ "File format {file_format} is not supported for insert_into
yet!"
+ );
}
// Build the writer from the already-parsed table properties so it
honors
@@ -275,8 +271,7 @@ impl ExecutionPlan for IcebergWriteExec {
fanout_enabled,
schema.clone(),
partition_spec,
- )
- .map_err(to_datafusion_error)?;
+ )?;
// Get input data
let data = execute_input_stream(
@@ -293,13 +288,10 @@ impl ExecutionPlan for IcebergWriteExec {
while let Some(batch) = input_stream.next().await {
let batch = batch?;
- task_writer
- .write(batch)
- .await
- .map_err(to_datafusion_error)?;
+ task_writer.write(batch).await?;
}
- let data_files =
task_writer.close().await.map_err(to_datafusion_error)?;
+ let data_files = task_writer.close().await?;
// Convert builders to data files and then to JSON strings
let data_files_strs: Vec<String> = data_files
@@ -312,7 +304,7 @@ impl ExecutionPlan for IcebergWriteExec {
)
.map_err(to_datafusion_error)
})
- .collect::<DFResult<Vec<String>>>()?;
+ .collect::<Result<Vec<String>>>()?;
Self::make_result_batch(data_files_strs)
})
@@ -335,7 +327,7 @@ mod tests {
use datafusion::arrow::datatypes::{
DataType, Field, Schema as ArrowSchema, SchemaRef as ArrowSchemaRef,
};
- use datafusion::common::Result as DFResult;
+ use datafusion::error::Result;
use datafusion::execution::{SendableRecordBatchStream, TaskContext};
use datafusion::physical_expr::{EquivalenceProperties, Partitioning};
use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType};
@@ -350,7 +342,8 @@ mod tests {
deserialize_data_file_from_json,
};
use iceberg::{
- Catalog, CatalogBuilder, MemoryCatalog, NamespaceIdent, Result,
TableCreation,
+ Catalog, CatalogBuilder, Error, ErrorKind, MemoryCatalog,
NamespaceIdent,
+ Result as IcebergResult, TableCreation,
};
use parquet::arrow::PARQUET_FIELD_ID_META_KEY;
use tempfile::TempDir;
@@ -414,15 +407,15 @@ mod tests {
fn apply_expressions(
&self,
- _f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) ->
DFResult<TreeNodeRecursion>,
- ) -> DFResult<TreeNodeRecursion> {
+ _f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) ->
Result<TreeNodeRecursion>,
+ ) -> Result<TreeNodeRecursion> {
Ok(TreeNodeRecursion::Continue)
}
fn with_new_children(
self: Arc<Self>,
_children: Vec<Arc<dyn ExecutionPlan>>,
- ) -> DFResult<Arc<dyn ExecutionPlan>> {
+ ) -> Result<Arc<dyn ExecutionPlan>> {
Ok(self)
}
@@ -430,7 +423,7 @@ mod tests {
&self,
_partition: usize,
_context: Arc<TaskContext>,
- ) -> DFResult<SendableRecordBatchStream> {
+ ) -> Result<SendableRecordBatchStream> {
let batches = self.batches.clone();
let stream = stream::iter(batches.into_iter().map(Ok));
Ok(Box::pin(RecordBatchStreamAdapter::new(
@@ -458,7 +451,7 @@ mod tests {
}
/// Helper function to create a test table schema
- fn get_test_schema() -> Result<Schema> {
+ fn get_test_schema() -> IcebergResult<Schema> {
Schema::builder()
.with_schema_id(0)
.with_fields(vec![
@@ -485,7 +478,7 @@ mod tests {
}
#[tokio::test]
- async fn test_iceberg_write_exec() -> Result<()> {
+ async fn test_iceberg_write_exec() -> Result<(), Box<dyn
std::error::Error>> {
// 1. Set up test environment
let iceberg_catalog = get_iceberg_catalog().await;
let namespace = NamespaceIdent::new("test_namespace".to_string());
@@ -654,7 +647,8 @@ mod tests {
}
#[tokio::test]
- async fn test_iceberg_write_exec_advertises_result_schema() -> Result<()> {
+ async fn test_iceberg_write_exec_advertises_result_schema()
+ -> Result<(), Box<dyn std::error::Error>> {
let iceberg_catalog = get_iceberg_catalog().await;
let namespace = NamespaceIdent::new("test_namespace".to_string());
iceberg_catalog
diff --git a/crates/datafusion/src/schema.rs b/crates/datafusion/src/schema.rs
index 04bea0b..6bc4b4d 100644
--- a/crates/datafusion/src/schema.rs
+++ b/crates/datafusion/src/schema.rs
@@ -20,18 +20,17 @@ use std::sync::Arc;
use async_trait::async_trait;
use dashmap::DashMap;
use datafusion::catalog::SchemaProvider;
+use datafusion::common::{exec_datafusion_err, exec_err, plan_datafusion_err};
use datafusion::datasource::TableProvider;
-use datafusion::error::{DataFusionError, Result as DFResult};
+use datafusion::error::Result;
use datafusion::execution::TaskContext;
use datafusion::prelude::SessionContext;
-use futures::StreamExt;
+use futures::TryStreamExt;
use futures::future::try_join_all;
use iceberg::arrow::arrow_schema_to_schema_auto_assign_ids;
use iceberg::inspect::MetadataTableType;
use iceberg::spec::FormatVersion;
-use iceberg::{
- Catalog, Error, ErrorKind, NamespaceIdent, Result, TableCreation,
TableIdent,
-};
+use iceberg::{Catalog, NamespaceIdent, TableCreation, TableIdent};
use crate::table::IcebergTableProvider;
use crate::to_datafusion_error;
@@ -69,7 +68,8 @@ impl IcebergSchemaProvider {
// As of right now; tables might become stale.
let table_names: Vec<_> = client
.list_tables(&namespace)
- .await?
+ .await
+ .map_err(to_datafusion_error)?
.iter()
.map(|tbl| tbl.name().to_string())
.collect();
@@ -122,15 +122,12 @@ impl SchemaProvider for IcebergSchemaProvider {
}
}
- async fn table(&self, name: &str) -> DFResult<Option<Arc<dyn
TableProvider>>> {
+ async fn table(&self, name: &str) -> Result<Option<Arc<dyn
TableProvider>>> {
if let Some((table_name, metadata_table_name)) = name.split_once('$') {
let metadata_table_type =
MetadataTableType::try_from(metadata_table_name)
- .map_err(DataFusionError::Plan)?;
+ .map_err(|e| plan_datafusion_err!("{e}"))?;
if let Some(table) = self.tables.get(table_name) {
- let metadata_table = table
- .metadata_table(metadata_table_type)
- .await
- .map_err(to_datafusion_error)?;
+ let metadata_table =
table.metadata_table(metadata_table_type).await?;
return Ok(Some(Arc::new(metadata_table)));
} else {
return Ok(None);
@@ -147,12 +144,10 @@ impl SchemaProvider for IcebergSchemaProvider {
&self,
name: String,
table: Arc<dyn TableProvider>,
- ) -> DFResult<Option<Arc<dyn TableProvider>>> {
+ ) -> Result<Option<Arc<dyn TableProvider>>> {
// Check if table already exists
if self.table_exist(name.as_str()) {
- return Err(DataFusionError::Execution(format!(
- "Table {name} already exists"
- )));
+ return exec_err!("Table {name} already exists");
}
// Convert DataFusion schema to Iceberg schema
@@ -184,9 +179,7 @@ impl SchemaProvider for IcebergSchemaProvider {
let rt = tokio::runtime::Handle::current();
rt.block_on(async move {
// Verify the input table is empty - CREATE TABLE only accepts
schema definition
- ensure_table_is_empty(&table)
- .await
- .map_err(to_datafusion_error)?;
+ ensure_table_is_empty(&table).await?;
catalog
.create_table(&namespace, table_creation)
@@ -199,8 +192,7 @@ impl SchemaProvider for IcebergSchemaProvider {
namespace.clone(),
name_clone.clone(),
)
- .await
- .map_err(to_datafusion_error)?;
+ .await?;
// Store the new table provider
tables.insert(name_clone, Arc::new(table_provider));
@@ -211,12 +203,11 @@ impl SchemaProvider for IcebergSchemaProvider {
// Block on the spawned task to get the result
// This is safe because spawn_blocking moves the blocking to a
dedicated thread pool
- futures::executor::block_on(result).map_err(|e| {
- DataFusionError::Execution(format!("Failed to create Iceberg
table: {e}"))
- })?
+ futures::executor::block_on(result)
+ .map_err(|e| exec_datafusion_err!("Failed to create Iceberg table:
{e}"))?
}
- fn deregister_table(&self, name: &str) -> DFResult<Option<Arc<dyn
TableProvider>>> {
+ fn deregister_table(&self, name: &str) -> Result<Option<Arc<dyn
TableProvider>>> {
// Check if table exists
if !self.table_exist(name) {
return Ok(None);
@@ -248,9 +239,8 @@ impl SchemaProvider for IcebergSchemaProvider {
})
});
- futures::executor::block_on(result).map_err(|e| {
- DataFusionError::Execution(format!("Failed to drop Iceberg table:
{e}"))
- })?
+ futures::executor::block_on(result)
+ .map_err(|e| exec_datafusion_err!("Failed to drop Iceberg table:
{e}"))?
}
}
@@ -258,32 +248,16 @@ impl SchemaProvider for IcebergSchemaProvider {
/// Returns an error if the table has any rows.
async fn ensure_table_is_empty(table: &Arc<dyn TableProvider>) -> Result<()> {
let session_ctx = SessionContext::new();
- let exec_plan = table
- .scan(&session_ctx.state(), None, &[], Some(1))
- .await
- .map_err(|e| {
- Error::new(ErrorKind::Unexpected, format!("Failed to scan table:
{e}"))
- })?;
+ let exec_plan = table.scan(&session_ctx.state(), None, &[],
Some(1)).await?;
let task_ctx = Arc::new(TaskContext::default());
- let stream = exec_plan.execute(0, task_ctx).map_err(|e| {
- Error::new(
- ErrorKind::Unexpected,
- format!("Failed to execute scan: {e}"),
- )
- })?;
+ let stream = exec_plan.execute(0, task_ctx)?;
- let batches: Vec<_> = stream.collect().await;
- let has_data = batches
- .into_iter()
- .filter_map(|r| r.ok())
- .any(|batch| batch.num_rows() > 0);
+ let batches: Vec<_> = stream.try_collect().await?;
+ let has_data = batches.iter().any(|batch| batch.num_rows() > 0);
if has_data {
- return Err(Error::new(
- ErrorKind::Unexpected,
- "register_table does not support tables with data.",
- ));
+ return exec_err!("register_table does not support tables with data.");
}
Ok(())
diff --git a/crates/datafusion/src/table/metadata_table.rs
b/crates/datafusion/src/table/metadata_table.rs
index 56d0c10..e9ef9fa 100644
--- a/crates/datafusion/src/table/metadata_table.rs
+++ b/crates/datafusion/src/table/metadata_table.rs
@@ -22,7 +22,7 @@ use datafusion::arrow::array::RecordBatch;
use datafusion::arrow::datatypes::SchemaRef as ArrowSchemaRef;
use datafusion::catalog::Session;
use datafusion::datasource::{TableProvider, TableType};
-use datafusion::error::Result as DFResult;
+use datafusion::error::Result;
use datafusion::logical_expr::Expr;
use datafusion::physical_plan::ExecutionPlan;
use futures::TryStreamExt;
@@ -64,13 +64,13 @@ impl TableProvider for IcebergMetadataTableProvider {
_projection: Option<&Vec<usize>>,
_filters: &[Expr],
_limit: Option<usize>,
- ) -> DFResult<Arc<dyn ExecutionPlan>> {
+ ) -> Result<Arc<dyn ExecutionPlan>> {
Ok(Arc::new(IcebergMetadataScan::new(self.clone())))
}
}
impl IcebergMetadataTableProvider {
- pub async fn scan(self) -> DFResult<BoxStream<'static,
DFResult<RecordBatch>>> {
+ pub async fn scan(self) -> Result<BoxStream<'static, Result<RecordBatch>>>
{
let metadata_table = self.table.inspect();
let stream = match self.r#type {
MetadataTableType::Snapshots =>
metadata_table.snapshots().scan().await,
diff --git a/crates/datafusion/src/table/mod.rs
b/crates/datafusion/src/table/mod.rs
index a1db3c8..69f4b3a 100644
--- a/crates/datafusion/src/table/mod.rs
+++ b/crates/datafusion/src/table/mod.rs
@@ -34,9 +34,9 @@ use std::sync::Arc;
use async_trait::async_trait;
use datafusion::arrow::datatypes::SchemaRef as ArrowSchemaRef;
use datafusion::catalog::Session;
-use datafusion::common::DataFusionError;
+use datafusion::common::{config_datafusion_err, exec_datafusion_err,
not_impl_err};
use datafusion::datasource::{TableProvider, TableType};
-use datafusion::error::Result as DFResult;
+use datafusion::error::Result;
use datafusion::logical_expr::dml::InsertOp;
use datafusion::logical_expr::{Expr, TableProviderFilterPushDown};
use datafusion::physical_plan::ExecutionPlan;
@@ -45,7 +45,7 @@ use iceberg::arrow::schema_to_arrow_schema;
use iceberg::inspect::MetadataTableType;
use iceberg::spec::TableProperties;
use iceberg::table::Table;
-use iceberg::{Catalog, Error, ErrorKind, NamespaceIdent, Result, TableIdent};
+use iceberg::{Catalog, NamespaceIdent, TableIdent};
use metadata_table::IcebergMetadataTableProvider;
use crate::error::to_datafusion_error;
@@ -87,8 +87,14 @@ impl IcebergTableProvider {
let table_ident = TableIdent::new(namespace, name.into());
// Load table once to get initial schema
- let table = catalog.load_table(&table_ident).await?;
- let schema =
Arc::new(schema_to_arrow_schema(table.metadata().current_schema())?);
+ let table = catalog
+ .load_table(&table_ident)
+ .await
+ .map_err(to_datafusion_error)?;
+ let schema = Arc::new(
+ schema_to_arrow_schema(table.metadata().current_schema())
+ .map_err(to_datafusion_error)?,
+ );
Ok(IcebergTableProvider {
catalog,
@@ -102,7 +108,11 @@ impl IcebergTableProvider {
r#type: MetadataTableType,
) -> Result<IcebergMetadataTableProvider> {
// Load fresh table metadata for metadata table access
- let table = self.catalog.load_table(&self.table_ident).await?;
+ let table = self
+ .catalog
+ .load_table(&self.table_ident)
+ .await
+ .map_err(to_datafusion_error)?;
Ok(IcebergMetadataTableProvider { table, r#type })
}
}
@@ -123,7 +133,7 @@ impl TableProvider for IcebergTableProvider {
projection: Option<&Vec<usize>>,
filters: &[Expr],
limit: Option<usize>,
- ) -> DFResult<Arc<dyn ExecutionPlan>> {
+ ) -> Result<Arc<dyn ExecutionPlan>> {
// Load fresh table metadata from catalog
let table = self
.catalog
@@ -145,7 +155,7 @@ impl TableProvider for IcebergTableProvider {
fn supports_filters_pushdown(
&self,
filters: &[&Expr],
- ) -> DFResult<Vec<TableProviderFilterPushDown>> {
+ ) -> Result<Vec<TableProviderFilterPushDown>> {
// Push down all filters, as a single source of truth, the scanner
will drop the filters which couldn't be push down
Ok(vec![TableProviderFilterPushDown::Inexact; filters.len()])
}
@@ -155,11 +165,11 @@ impl TableProvider for IcebergTableProvider {
state: &dyn Session,
input: Arc<dyn ExecutionPlan>,
_insert_op: InsertOp,
- ) -> DFResult<Arc<dyn ExecutionPlan>> {
+ ) -> Result<Arc<dyn ExecutionPlan>> {
if _insert_op != InsertOp::Append {
- return Err(DataFusionError::NotImplemented(format!(
+ return not_impl_err!(
"IcebergTableProvider supports only append inserts, got
{_insert_op}"
- )));
+ );
}
// Load fresh table metadata from catalog
@@ -181,9 +191,7 @@ impl TableProvider for IcebergTableProvider {
// Step 2: Repartition for parallel processing
let target_partitions =
NonZeroUsize::new(state.config().target_partitions())
.ok_or_else(|| {
- DataFusionError::Configuration(
- "target_partitions must be greater than 0".to_string(),
- )
+ config_datafusion_err!("target_partitions must be greater than
0")
})?;
let repartitioned_plan =
@@ -195,19 +203,12 @@ impl TableProvider for IcebergTableProvider {
.properties()
.get(TableProperties::PROPERTY_DATAFUSION_WRITE_FANOUT_ENABLED)
.map(|value| {
- value
- .parse::<bool>()
- .map_err(|e| {
- Error::new(
- ErrorKind::DataInvalid,
- format!(
- "Invalid value for {}, expected 'true' or
'false'",
-
TableProperties::PROPERTY_DATAFUSION_WRITE_FANOUT_ENABLED
- ),
- )
- .with_source(e)
- })
- .map_err(to_datafusion_error)
+ value.parse::<bool>().map_err(|e| {
+ config_datafusion_err!(
+ "Invalid value for {}, expected 'true' or 'false':
{e}",
+
TableProperties::PROPERTY_DATAFUSION_WRITE_FANOUT_ENABLED
+ )
+ })
})
.transpose()?
.unwrap_or(TableProperties::PROPERTY_DATAFUSION_WRITE_FANOUT_ENABLED_DEFAULT);
@@ -255,7 +256,10 @@ impl IcebergStaticTableProvider {
///
/// Uses the table's current snapshot for all queries. Does not support
write operations.
pub async fn try_new_from_table(table: Table) -> Result<Self> {
- let schema =
Arc::new(schema_to_arrow_schema(table.metadata().current_schema())?);
+ let schema = Arc::new(
+ schema_to_arrow_schema(table.metadata().current_schema())
+ .map_err(to_datafusion_error)?,
+ );
Ok(IcebergStaticTableProvider {
table,
snapshot_id: None,
@@ -275,16 +279,16 @@ impl IcebergStaticTableProvider {
.metadata()
.snapshot_by_id(snapshot_id)
.ok_or_else(|| {
- Error::new(
- ErrorKind::Unexpected,
- format!(
- "snapshot id {snapshot_id} not found in table {}",
- table.identifier().name()
- ),
+ exec_datafusion_err!(
+ "snapshot id {snapshot_id} not found in table {}",
+ table.identifier().name()
)
})?;
- let table_schema = snapshot.schema(table.metadata())?;
- let schema = Arc::new(schema_to_arrow_schema(&table_schema)?);
+ let table_schema = snapshot
+ .schema(table.metadata())
+ .map_err(to_datafusion_error)?;
+ let schema =
+
Arc::new(schema_to_arrow_schema(&table_schema).map_err(to_datafusion_error)?);
Ok(IcebergStaticTableProvider {
table,
snapshot_id: Some(snapshot_id),
@@ -309,7 +313,7 @@ impl TableProvider for IcebergStaticTableProvider {
projection: Option<&Vec<usize>>,
filters: &[Expr],
limit: Option<usize>,
- ) -> DFResult<Arc<dyn ExecutionPlan>> {
+ ) -> Result<Arc<dyn ExecutionPlan>> {
// Use cached table (no refresh)
Ok(Arc::new(IcebergTableScan::new(
self.table.clone(),
@@ -324,7 +328,7 @@ impl TableProvider for IcebergStaticTableProvider {
fn supports_filters_pushdown(
&self,
filters: &[&Expr],
- ) -> DFResult<Vec<TableProviderFilterPushDown>> {
+ ) -> Result<Vec<TableProviderFilterPushDown>> {
// Push down all filters, as a single source of truth, the scanner
will drop the filters which couldn't be push down
Ok(vec![TableProviderFilterPushDown::Inexact; filters.len()])
}
@@ -334,13 +338,11 @@ impl TableProvider for IcebergStaticTableProvider {
_state: &dyn Session,
_input: Arc<dyn ExecutionPlan>,
_insert_op: InsertOp,
- ) -> DFResult<Arc<dyn ExecutionPlan>> {
- Err(to_datafusion_error(Error::new(
- ErrorKind::FeatureUnsupported,
+ ) -> Result<Arc<dyn ExecutionPlan>> {
+ not_impl_err!(
"Write operations are not supported on IcebergStaticTableProvider.
\
Use IcebergTableProvider with a catalog for write support."
- .to_string(),
- )))
+ )
}
}
@@ -350,6 +352,7 @@ mod tests {
use std::sync::Arc;
use datafusion::common::Column;
+ use datafusion::error::DataFusionError;
use datafusion::physical_plan::ExecutionPlan;
use datafusion::prelude::SessionContext;
use iceberg::io::FileIO;
diff --git a/crates/datafusion/src/table/table_provider_factory.rs
b/crates/datafusion/src/table/table_provider_factory.rs
index 656854d..1e7704b 100644
--- a/crates/datafusion/src/table/table_provider_factory.rs
+++ b/crates/datafusion/src/table/table_provider_factory.rs
@@ -21,12 +21,12 @@ use std::sync::Arc;
use async_trait::async_trait;
use datafusion::catalog::{Session, TableProvider, TableProviderFactory};
-use datafusion::common::TableReference;
-use datafusion::error::Result as DFResult;
+use datafusion::common::{TableReference, not_impl_err};
+use datafusion::error::Result;
use datafusion::logical_expr::CreateExternalTable;
+use iceberg::TableIdent;
use iceberg::io::{FileIOBuilder, LocalFsStorageFactory, StorageFactory};
use iceberg::table::StaticTable;
-use iceberg::{Error, ErrorKind, Result, TableIdent};
use super::IcebergStaticTableProvider;
use crate::to_datafusion_error;
@@ -122,8 +122,8 @@ impl TableProviderFactory for IcebergTableProviderFactory {
&self,
_state: &dyn Session,
cmd: &CreateExternalTable,
- ) -> DFResult<Arc<dyn TableProvider>> {
- let metadata_file_path = check_cmd(cmd).map_err(to_datafusion_error)?;
+ ) -> Result<Arc<dyn TableProvider>> {
+ let metadata_file_path = check_cmd(cmd)?;
let table_name = &cmd.name;
let options = &cmd.options;
@@ -141,13 +141,10 @@ impl TableProviderFactory for IcebergTableProviderFactory
{
options,
storage_factory,
)
- .await
- .map_err(to_datafusion_error)?
+ .await?
.into_table();
- let provider = IcebergStaticTableProvider::try_new_from_table(table)
- .await
- .map_err(to_datafusion_error)?;
+ let provider =
IcebergStaticTableProvider::try_new_from_table(table).await?;
Ok(Arc::new(provider))
}
@@ -171,18 +168,16 @@ fn check_cmd(cmd: &CreateExternalTable) -> Result<&str> {
|| !column_defaults.is_empty();
if is_invalid {
- return Err(Error::new(
- ErrorKind::FeatureUnsupported,
- "Currently we only support reading existing icebergs tables in
external table command. To create new table, please use catalog provider.",
- ));
+ return not_impl_err!(
+ "Currently we only support reading existing icebergs tables in
external table command. To create new table, please use catalog provider."
+ );
}
match cmd.locations.as_slice() {
[location] => Ok(location),
- _ => Err(Error::new(
- ErrorKind::FeatureUnsupported,
- "Iceberg external tables require exactly one metadata location.",
- )),
+ _ => not_impl_err!(
+ "Iceberg external tables require exactly one metadata location."
+ ),
}
}
@@ -213,11 +208,14 @@ async fn create_static_table(
props: &HashMap<String, String>,
storage_factory: Arc<dyn StorageFactory>,
) -> Result<StaticTable> {
- let table_ident = TableIdent::from_strs(table_name.to_vec())?;
+ let table_ident =
+
TableIdent::from_strs(table_name.to_vec()).map_err(to_datafusion_error)?;
let file_io = FileIOBuilder::new(storage_factory)
.with_props(props)
.build();
- StaticTable::from_metadata_file(metadata_file_path, table_ident,
file_io).await
+ StaticTable::from_metadata_file(metadata_file_path, table_ident, file_io)
+ .await
+ .map_err(to_datafusion_error)
}
#[cfg(test)]
diff --git a/crates/datafusion/src/task_writer.rs
b/crates/datafusion/src/task_writer.rs
index 3e7c713..45df287 100644
--- a/crates/datafusion/src/task_writer.rs
+++ b/crates/datafusion/src/task_writer.rs
@@ -21,7 +21,7 @@
//! of RecordBatch data to Iceberg tables.
use datafusion::arrow::array::RecordBatch;
-use iceberg::Result;
+use datafusion::error::Result;
use iceberg::arrow::RecordBatchPartitionSplitter;
use iceberg::spec::{DataFile, PartitionSpecRef, SchemaRef};
use iceberg::writer::IcebergWriterBuilder;
@@ -30,6 +30,8 @@ use
iceberg::writer::partitioning::clustered_writer::ClusteredWriter;
use iceberg::writer::partitioning::fanout_writer::FanoutWriter;
use iceberg::writer::partitioning::unpartitioned_writer::UnpartitionedWriter;
+use crate::to_datafusion_error;
+
/// High-level writer for DataFusion that handles partitioning and routing of
RecordBatch data.
///
/// TaskWriter coordinates writing data to Iceberg tables by:
@@ -138,7 +140,8 @@ impl<B: IcebergWriterBuilder> TaskWriter<B> {
RecordBatchPartitionSplitter::try_new_with_precomputed_values(
schema.clone(),
partition_spec.clone(),
- )?,
+ )
+ .map_err(to_datafusion_error)?,
)
} else {
None
@@ -183,7 +186,7 @@ impl<B: IcebergWriterBuilder> TaskWriter<B> {
match &mut self.writer {
SupportedWriter::Unpartitioned(writer) => {
// Unpartitioned: write directly without splitting
- writer.write(batch).await
+ writer.write(batch).await.map_err(to_datafusion_error)
}
SupportedWriter::Fanout(writer) => {
Self::write_partitioned_batches(writer,
&self.partition_splitter, &batch)
@@ -220,11 +223,14 @@ impl<B: IcebergWriterBuilder> TaskWriter<B> {
let splitter = partition_splitter
.as_ref()
.expect("Partition splitter should be initialized");
- let partitioned_batches = splitter.split(batch)?;
+ let partitioned_batches =
splitter.split(batch).map_err(to_datafusion_error)?;
// Write each partition
for (partition_key, partition_batch) in partitioned_batches {
- writer.write(partition_key, partition_batch).await?;
+ writer
+ .write(partition_key, partition_batch)
+ .await
+ .map_err(to_datafusion_error)?;
}
Ok(())
@@ -254,9 +260,15 @@ impl<B: IcebergWriterBuilder> TaskWriter<B> {
/// ```
pub async fn close(self) -> Result<Vec<DataFile>> {
match self.writer {
- SupportedWriter::Unpartitioned(writer) => writer.close().await,
- SupportedWriter::Fanout(writer) => writer.close().await,
- SupportedWriter::Clustered(writer) => writer.close().await,
+ SupportedWriter::Unpartitioned(writer) => {
+ writer.close().await.map_err(to_datafusion_error)
+ }
+ SupportedWriter::Fanout(writer) => {
+ writer.close().await.map_err(to_datafusion_error)
+ }
+ SupportedWriter::Clustered(writer) => {
+ writer.close().await.map_err(to_datafusion_error)
+ }
}
}
}
@@ -287,7 +299,7 @@ mod tests {
use super::*;
- fn create_test_schema() -> Result<Arc<iceberg::spec::Schema>> {
+ fn create_test_schema() -> iceberg::Result<Arc<iceberg::spec::Schema>> {
Ok(Arc::new(
iceberg::spec::Schema::builder()
.with_schema_id(1)
@@ -358,7 +370,7 @@ mod tests {
fn create_writer_builder(
temp_dir: &TempDir,
schema: Arc<iceberg::spec::Schema>,
- ) -> Result<
+ ) -> iceberg::Result<
DataFileWriterBuilder<
ParquetWriterBuilder,
DefaultLocationGenerator,
@@ -386,7 +398,7 @@ mod tests {
}
#[tokio::test]
- async fn test_task_writer_unpartitioned() -> Result<()> {
+ async fn test_task_writer_unpartitioned() -> Result<(), Box<dyn
std::error::Error>> {
let temp_dir = TempDir::new()?;
let schema = create_test_schema()?;
let arrow_schema = create_arrow_schema();
@@ -453,7 +465,8 @@ mod tests {
}
#[tokio::test]
- async fn test_task_writer_partitioned_fanout() -> Result<()> {
+ async fn test_task_writer_partitioned_fanout()
+ -> Result<(), Box<dyn std::error::Error>> {
let temp_dir = TempDir::new()?;
let schema = create_test_schema()?;
let arrow_schema = create_arrow_schema_with_partition();
@@ -504,7 +517,8 @@ mod tests {
}
#[tokio::test]
- async fn test_task_writer_partitioned_clustered() -> Result<()> {
+ async fn test_task_writer_partitioned_clustered()
+ -> Result<(), Box<dyn std::error::Error>> {
let temp_dir = TempDir::new()?;
let schema = create_test_schema()?;
let arrow_schema = create_arrow_schema_with_partition();
diff --git a/crates/datafusion/tests/integration_datafusion_test.rs
b/crates/datafusion/tests/integration_datafusion_test.rs
index 41753ab..dcea6a5 100644
--- a/crates/datafusion/tests/integration_datafusion_test.rs
+++ b/crates/datafusion/tests/integration_datafusion_test.rs
@@ -18,6 +18,7 @@
//! Integration tests for Iceberg Datafusion with Hive Metastore.
use std::collections::HashMap;
+use std::error::Error;
use std::sync::Arc;
use std::vec;
@@ -34,8 +35,8 @@ use iceberg::spec::{
};
use iceberg::test_utils::check_record_batches;
use iceberg::{
- Catalog, CatalogBuilder, MemoryCatalog, NamespaceIdent, Result,
TableCreation,
- TableIdent,
+ Catalog, CatalogBuilder, MemoryCatalog, NamespaceIdent, Result as
IcebergResult,
+ TableCreation, TableIdent,
};
use tempfile::TempDir;
@@ -65,7 +66,7 @@ fn get_struct_type() -> StructType {
async fn set_test_namespace(
catalog: &MemoryCatalog,
namespace: &NamespaceIdent,
-) -> Result<()> {
+) -> IcebergResult<()> {
let properties = HashMap::new();
catalog.create_namespace(namespace, properties).await?;
@@ -77,7 +78,7 @@ fn get_table_creation(
location: impl ToString,
name: impl ToString,
schema: Option<Schema>,
-) -> Result<TableCreation> {
+) -> IcebergResult<TableCreation> {
let schema = match schema {
None => Schema::builder()
.with_schema_id(0)
@@ -102,7 +103,7 @@ fn get_table_creation(
}
#[tokio::test]
-async fn test_provider_plan_stream_schema() -> Result<()> {
+async fn test_provider_plan_stream_schema() -> Result<(), Box<dyn Error>> {
let iceberg_catalog = get_iceberg_catalog().await;
let namespace =
NamespaceIdent::new("test_provider_get_table_schema".to_string());
set_test_namespace(&iceberg_catalog, &namespace).await?;
@@ -155,7 +156,7 @@ async fn test_provider_plan_stream_schema() -> Result<()> {
}
#[tokio::test]
-async fn test_provider_list_table_names() -> Result<()> {
+async fn test_provider_list_table_names() -> Result<(), Box<dyn Error>> {
let iceberg_catalog = get_iceberg_catalog().await;
let namespace =
NamespaceIdent::new("test_provider_list_table_names".to_string());
set_test_namespace(&iceberg_catalog, &namespace).await?;
@@ -188,7 +189,7 @@ async fn test_provider_list_table_names() -> Result<()> {
}
#[tokio::test]
-async fn test_provider_list_schema_names() -> Result<()> {
+async fn test_provider_list_schema_names() -> Result<(), Box<dyn Error>> {
let iceberg_catalog = get_iceberg_catalog().await;
let namespace =
NamespaceIdent::new("test_provider_list_schema_names".to_string());
set_test_namespace(&iceberg_catalog, &namespace).await?;
@@ -213,7 +214,7 @@ async fn test_provider_list_schema_names() -> Result<()> {
}
#[tokio::test]
-async fn test_table_projection() -> Result<()> {
+async fn test_table_projection() -> Result<(), Box<dyn Error>> {
let iceberg_catalog = get_iceberg_catalog().await;
let namespace = NamespaceIdent::new("ns".to_string());
set_test_namespace(&iceberg_catalog, &namespace).await?;
@@ -282,7 +283,7 @@ async fn test_table_projection() -> Result<()> {
}
#[tokio::test]
-async fn test_table_predict_pushdown() -> Result<()> {
+async fn test_table_predict_pushdown() -> Result<(), Box<dyn Error>> {
let iceberg_catalog = get_iceberg_catalog().await;
let namespace = NamespaceIdent::new("ns".to_string());
set_test_namespace(&iceberg_catalog, &namespace).await?;
@@ -328,7 +329,7 @@ async fn test_table_predict_pushdown() -> Result<()> {
}
#[tokio::test]
-async fn test_metadata_table() -> Result<()> {
+async fn test_metadata_table() -> Result<(), Box<dyn Error>> {
let iceberg_catalog = get_iceberg_catalog().await;
let namespace = NamespaceIdent::new("ns".to_string());
set_test_namespace(&iceberg_catalog, &namespace).await?;
@@ -455,7 +456,7 @@ async fn test_metadata_table() -> Result<()> {
}
#[tokio::test]
-async fn test_insert_into() -> Result<()> {
+async fn test_insert_into() -> Result<(), Box<dyn Error>> {
let iceberg_catalog = get_iceberg_catalog().await;
let namespace = NamespaceIdent::new("test_insert_into".to_string());
set_test_namespace(&iceberg_catalog, &namespace).await?;
@@ -581,7 +582,7 @@ fn get_nested_struct_type() -> StructType {
}
#[tokio::test]
-async fn test_insert_into_nested() -> Result<()> {
+async fn test_insert_into_nested() -> Result<(), Box<dyn Error>> {
let iceberg_catalog = get_iceberg_catalog().await;
let namespace = NamespaceIdent::new("test_insert_nested".to_string());
set_test_namespace(&iceberg_catalog, &namespace).await?;
@@ -839,7 +840,7 @@ async fn test_insert_into_nested() -> Result<()> {
}
#[tokio::test]
-async fn test_insert_into_partitioned() -> Result<()> {
+async fn test_insert_into_partitioned() -> Result<(), Box<dyn Error>> {
let iceberg_catalog = get_iceberg_catalog().await;
let namespace = NamespaceIdent::new("test_partitioned_write".to_string());
set_test_namespace(&iceberg_catalog, &namespace).await?;
diff --git
a/crates/sqllogictest/testdata/slts/df_test/timestamp_predicate_pushdown.slt
b/crates/sqllogictest/testdata/slts/df_test/timestamp_predicate_pushdown.slt
index ffa7417..8b1f833 100644
--- a/crates/sqllogictest/testdata/slts/df_test/timestamp_predicate_pushdown.slt
+++ b/crates/sqllogictest/testdata/slts/df_test/timestamp_predicate_pushdown.slt
@@ -174,10 +174,10 @@ DROP TABLE default.default.test_timestamp_micros
# Test with TIMESTAMP(3) - millisecond precision
# This should fail because Iceberg doesn't support millisecond precision
-statement error DataFusion error: External error: DataInvalid => Unsupported
Arrow data type: Timestamp\(ms\)
+statement error (?s)DataFusion error: IcebergError\(DataInvalid\).*Unsupported
Arrow data type: Timestamp\(ms\)
CREATE TABLE default.default.test_timestamp_millis (id INT NOT NULL, ts
TIMESTAMP(3))
# Test with TIMESTAMP(0) - second precision
# This should fail because Iceberg doesn't support second precision
-statement error DataFusion error: External error: DataInvalid => Unsupported
Arrow data type: Timestamp\(s\)
+statement error (?s)DataFusion error: IcebergError\(DataInvalid\).*Unsupported
Arrow data type: Timestamp\(s\)
CREATE TABLE default.default.test_timestamp_seconds (id INT NOT NULL, ts
TIMESTAMP(0))
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]