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]

Reply via email to