Skip to content

Commit 93879b2

Browse files
committed
fix(datafusion): reject unsupported insert operations
1 parent ac92ec9 commit 93879b2

1 file changed

Lines changed: 61 additions & 1 deletion

File tree

  • crates/integrations/datafusion/src/table

crates/integrations/datafusion/src/table/mod.rs

Lines changed: 61 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -159,8 +159,14 @@ impl TableProvider for IcebergTableProvider {
159159
&self,
160160
state: &dyn Session,
161161
input: Arc<dyn ExecutionPlan>,
162-
_insert_op: InsertOp,
162+
insert_op: InsertOp,
163163
) -> DFResult<Arc<dyn ExecutionPlan>> {
164+
if insert_op != InsertOp::Append {
165+
return Err(DataFusionError::NotImplemented(format!(
166+
"IcebergTableProvider supports only append inserts, got {insert_op}"
167+
)));
168+
}
169+
164170
// Load fresh table metadata from catalog
165171
let table = self
166172
.catalog
@@ -710,6 +716,60 @@ mod tests {
710716
false
711717
}
712718

719+
#[tokio::test]
720+
async fn test_catalog_backed_provider_rejects_overwrite_insert() {
721+
use datafusion::physical_plan::empty::EmptyExec;
722+
723+
let (catalog, namespace, table_name, _temp_dir) = get_test_catalog_and_table().await;
724+
let provider = IcebergTableProvider::try_new(catalog, namespace, table_name)
725+
.await
726+
.unwrap();
727+
let ctx = SessionContext::new();
728+
let input = Arc::new(EmptyExec::new(provider.schema())) as Arc<dyn ExecutionPlan>;
729+
730+
let error = provider
731+
.insert_into(&ctx.state(), input, InsertOp::Overwrite)
732+
.await
733+
.expect_err("overwrite inserts should be rejected");
734+
735+
assert!(
736+
matches!(
737+
error,
738+
DataFusionError::NotImplemented(ref message)
739+
if message
740+
== "IcebergTableProvider supports only append inserts, got Insert Overwrite"
741+
),
742+
"unexpected error: {error}"
743+
);
744+
}
745+
746+
#[tokio::test]
747+
async fn test_catalog_backed_provider_rejects_replace_insert() {
748+
use datafusion::physical_plan::empty::EmptyExec;
749+
750+
let (catalog, namespace, table_name, _temp_dir) = get_test_catalog_and_table().await;
751+
let provider = IcebergTableProvider::try_new(catalog, namespace, table_name)
752+
.await
753+
.unwrap();
754+
let ctx = SessionContext::new();
755+
let input = Arc::new(EmptyExec::new(provider.schema())) as Arc<dyn ExecutionPlan>;
756+
757+
let error = provider
758+
.insert_into(&ctx.state(), input, InsertOp::Replace)
759+
.await
760+
.expect_err("replace inserts should be rejected");
761+
762+
assert!(
763+
matches!(
764+
error,
765+
DataFusionError::NotImplemented(ref message)
766+
if message
767+
== "IcebergTableProvider supports only append inserts, got Replace Into"
768+
),
769+
"unexpected error: {error}"
770+
);
771+
}
772+
713773
#[tokio::test]
714774
async fn test_insert_plan_fanout_enabled_no_sort() {
715775
use datafusion::datasource::TableProvider;

0 commit comments

Comments
 (0)