| // Licensed to the Apache Software Foundation (ASF) under one |
| // or more contributor license agreements. See the NOTICE file |
| // distributed with this work for additional information |
| // regarding copyright ownership. The ASF licenses this file |
| // to you under the Apache License, Version 2.0 (the |
| // "License"); you may not use this file except in compliance |
| // with the License. You may obtain a copy of the License at |
| // |
| // http://www.apache.org/licenses/LICENSE-2.0 |
| // |
| // Unless required by applicable law or agreed to in writing, |
| // software distributed under the License is distributed on an |
| // "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY |
| // KIND, either express or implied. See the License for the |
| // specific language governing permissions and limitations |
| // under the License. |
| |
| //! Integration tests for rest catalog. |
| |
| mod common; |
| |
| use std::sync::Arc; |
| |
| use arrow_array::{ArrayRef, BooleanArray, Int32Array, RecordBatch, StringArray}; |
| use common::{random_ns, test_schema}; |
| use futures::TryStreamExt; |
| use iceberg::transaction::{ApplyTransactionAction, Transaction}; |
| use iceberg::writer::base_writer::data_file_writer::DataFileWriterBuilder; |
| use iceberg::writer::file_writer::ParquetWriterBuilder; |
| use iceberg::writer::file_writer::location_generator::{ |
| DefaultFileNameGenerator, DefaultLocationGenerator, |
| }; |
| use iceberg::writer::file_writer::rolling_writer::RollingFileWriterBuilder; |
| use iceberg::writer::{IcebergWriter, IcebergWriterBuilder}; |
| use iceberg::{Catalog, CatalogBuilder, TableCreation}; |
| use iceberg_catalog_rest::RestCatalogBuilder; |
| use iceberg_integration_tests::get_test_fixture; |
| use iceberg_storage_opendal::OpenDalStorageFactory; |
| use parquet::file::properties::WriterProperties; |
| |
| #[tokio::test] |
| async fn test_append_data_file_conflict() { |
| let fixture = get_test_fixture(); |
| let rest_catalog = RestCatalogBuilder::default() |
| .with_storage_factory(Arc::new(OpenDalStorageFactory::S3 { |
| configured_scheme: "s3".to_string(), |
| customized_credential_load: None, |
| })) |
| .load("rest", fixture.catalog_config.clone()) |
| .await |
| .unwrap(); |
| let ns = random_ns().await; |
| let schema = test_schema(); |
| |
| let table_creation = TableCreation::builder() |
| .name("t1".to_string()) |
| .schema(schema.clone()) |
| .build(); |
| |
| let table = rest_catalog |
| .create_table(ns.name(), table_creation) |
| .await |
| .unwrap(); |
| |
| // Create the writer and write the data |
| let schema: Arc<arrow_schema::Schema> = Arc::new( |
| table |
| .metadata() |
| .current_schema() |
| .as_ref() |
| .try_into() |
| .unwrap(), |
| ); |
| let location_generator = DefaultLocationGenerator::new(table.metadata().clone()).unwrap(); |
| let file_name_generator = DefaultFileNameGenerator::new( |
| "test".to_string(), |
| None, |
| iceberg::spec::DataFileFormat::Parquet, |
| ); |
| let parquet_writer_builder = ParquetWriterBuilder::new( |
| WriterProperties::default(), |
| table.metadata().current_schema().clone(), |
| ); |
| let rolling_file_writer_builder = RollingFileWriterBuilder::new_with_default_file_size( |
| parquet_writer_builder, |
| table.file_io().clone(), |
| location_generator.clone(), |
| file_name_generator.clone(), |
| ); |
| let data_file_writer_builder = DataFileWriterBuilder::new(rolling_file_writer_builder); |
| let mut data_file_writer = data_file_writer_builder.build(None).await.unwrap(); |
| let col1 = StringArray::from(vec![Some("foo"), Some("bar"), None, Some("baz")]); |
| let col2 = Int32Array::from(vec![Some(1), Some(2), Some(3), Some(4)]); |
| let col3 = BooleanArray::from(vec![Some(true), Some(false), None, Some(false)]); |
| let batch = RecordBatch::try_new(schema.clone(), vec![ |
| Arc::new(col1) as ArrayRef, |
| Arc::new(col2) as ArrayRef, |
| Arc::new(col3) as ArrayRef, |
| ]) |
| .unwrap(); |
| data_file_writer.write(batch.clone()).await.unwrap(); |
| let data_file = data_file_writer.close().await.unwrap(); |
| |
| // start two transaction and commit one of them |
| let tx1 = Transaction::new(&table); |
| let append_action = tx1.fast_append().add_data_files(data_file.clone()); |
| let tx1 = append_action.apply(tx1).unwrap(); |
| |
| let tx2 = Transaction::new(&table); |
| let append_action = tx2.fast_append().add_data_files(data_file.clone()); |
| let tx2 = append_action.apply(tx2).unwrap(); |
| let table = tx2 |
| .commit(&rest_catalog) |
| .await |
| .expect("The first commit should not fail."); |
| |
| // check result |
| let batch_stream = table |
| .scan() |
| .select_all() |
| .build() |
| .unwrap() |
| .to_arrow() |
| .await |
| .unwrap(); |
| let batches: Vec<_> = batch_stream.try_collect().await.unwrap(); |
| assert_eq!(batches.len(), 1); |
| assert_eq!(batches[0], batch); |
| |
| // another commit should fail |
| assert!(tx1.commit(&rest_catalog).await.is_err()); |
| } |