1use std::collections::HashSet;
21use std::sync::Arc;
22
23use arrow_array::RecordBatch;
24use arrow_schema::{DataType, Field, SchemaRef as ArrowSchemaRef};
25use itertools::Itertools;
26use parquet::arrow::PARQUET_FIELD_ID_META_KEY;
27
28use crate::arrow::record_batch_projector::RecordBatchProjector;
29use crate::arrow::schema_to_arrow_schema;
30use crate::error::invalid_data;
31use crate::spec::{DataFile, PartitionKey, Schema, SchemaRef};
32use crate::writer::file_writer::FileWriterBuilder;
33use crate::writer::file_writer::location_generator::{FileNameGenerator, LocationGenerator};
34use crate::writer::file_writer::rolling_writer::{RollingFileWriter, RollingFileWriterBuilder};
35use crate::writer::{IcebergWriter, IcebergWriterBuilder};
36use crate::{Error, ErrorKind, Result};
37
38#[derive(Debug)]
40pub struct EqualityDeleteFileWriterBuilder<
41 B: FileWriterBuilder,
42 L: LocationGenerator,
43 F: FileNameGenerator,
44> {
45 inner: RollingFileWriterBuilder<B, L, F>,
46 config: EqualityDeleteWriterConfig,
47}
48
49impl<B, L, F> EqualityDeleteFileWriterBuilder<B, L, F>
50where
51 B: FileWriterBuilder,
52 L: LocationGenerator,
53 F: FileNameGenerator,
54{
55 pub fn new(
57 inner: RollingFileWriterBuilder<B, L, F>,
58 config: EqualityDeleteWriterConfig,
59 ) -> Self {
60 Self { inner, config }
61 }
62}
63
64#[derive(Debug)]
66pub struct EqualityDeleteWriterConfig {
67 equality_ids: Vec<i32>,
69 projector: RecordBatchProjector,
71}
72
73fn validate_equality_ids(equality_ids: &[i32], original_schema: &Schema) -> Result<()> {
81 if equality_ids.is_empty() {
82 return Err(invalid_data!(
83 "Equality delete field ids must not be empty."
84 ));
85 }
86
87 let mut seen = HashSet::with_capacity(equality_ids.len());
88 for id in equality_ids {
89 if !seen.insert(*id) {
90 return Err(invalid_data!("Duplicate equality delete field id: {id}")
91 .with_context("field_id", id.to_string()));
92 }
93
94 if original_schema.field_by_id(*id).is_none() {
95 return Err(invalid_data!("Invalid equality delete field id: {id}")
96 .with_context("field_id", id.to_string()));
97 }
98 }
99
100 Ok(())
101}
102
103impl EqualityDeleteWriterConfig {
104 pub fn new(equality_ids: Vec<i32>, original_schema: SchemaRef) -> Result<Self> {
106 validate_equality_ids(&equality_ids, &original_schema)?;
107
108 let original_arrow_schema = Arc::new(schema_to_arrow_schema(&original_schema)?);
109 let projector = RecordBatchProjector::new(
110 original_arrow_schema,
111 &equality_ids,
112 |field| {
117 if field.data_type().is_nested()
119 || matches!(
120 field.data_type(),
121 DataType::Float16 | DataType::Float32 | DataType::Float64
122 )
123 {
124 return Ok(None);
125 }
126 Ok(Some(
127 field
128 .metadata()
129 .get(PARQUET_FIELD_ID_META_KEY)
130 .ok_or_else(|| {
131 Error::new(ErrorKind::Unexpected, "Field metadata is missing.")
132 })?
133 .parse::<i64>()
134 .map_err(|e| Error::new(ErrorKind::Unexpected, e.to_string()))?,
135 ))
136 },
137 |_field: &Field| true,
138 )?;
139 Ok(Self {
140 equality_ids,
141 projector,
142 })
143 }
144
145 pub fn projected_arrow_schema_ref(&self) -> &ArrowSchemaRef {
147 self.projector.projected_schema_ref()
148 }
149}
150
151#[async_trait::async_trait]
152impl<B, L, F> IcebergWriterBuilder for EqualityDeleteFileWriterBuilder<B, L, F>
153where
154 B: FileWriterBuilder,
155 L: LocationGenerator,
156 F: FileNameGenerator,
157{
158 type R = EqualityDeleteFileWriter<B, L, F>;
159
160 async fn build(&self, partition_key: Option<PartitionKey>) -> Result<Self::R> {
161 Ok(EqualityDeleteFileWriter {
162 inner: Some(self.inner.build()),
163 projector: self.config.projector.clone(),
164 equality_ids: self.config.equality_ids.clone(),
165 partition_key,
166 })
167 }
168}
169
170#[derive(Debug)]
172pub struct EqualityDeleteFileWriter<
173 B: FileWriterBuilder,
174 L: LocationGenerator,
175 F: FileNameGenerator,
176> {
177 inner: Option<RollingFileWriter<B, L, F>>,
178 projector: RecordBatchProjector,
179 equality_ids: Vec<i32>,
180 partition_key: Option<PartitionKey>,
181}
182
183#[async_trait::async_trait]
184impl<B, L, F> IcebergWriter for EqualityDeleteFileWriter<B, L, F>
185where
186 B: FileWriterBuilder,
187 L: LocationGenerator,
188 F: FileNameGenerator,
189{
190 async fn write(&mut self, batch: RecordBatch) -> Result<()> {
191 let batch = self.projector.project_batch(batch)?;
192 if let Some(writer) = self.inner.as_mut() {
193 writer.write(&self.partition_key, &batch).await
194 } else {
195 Err(Error::new(
196 ErrorKind::Unexpected,
197 "Equality delete inner writer has been closed.",
198 ))
199 }
200 }
201
202 async fn close(&mut self) -> Result<Vec<DataFile>> {
203 if let Some(writer) = self.inner.take() {
204 writer
205 .close()
206 .await?
207 .into_iter()
208 .map(|mut res| {
209 res.content(crate::spec::DataContentType::EqualityDeletes);
210 res.equality_ids(Some(self.equality_ids.iter().copied().collect_vec()));
211 if let Some(pk) = self.partition_key.as_ref() {
212 res.partition(pk.data().clone());
213 res.partition_spec_id(pk.spec().spec_id());
214 }
215 res.build()
216 .map_err(|e| invalid_data!("Failed to build data file: {e}"))
217 })
218 .collect()
219 } else {
220 Err(Error::new(
221 ErrorKind::Unexpected,
222 "Equality delete inner writer has been closed.",
223 ))
224 }
225 }
226}
227
228#[cfg(test)]
229mod test {
230 use std::collections::HashMap;
231 use std::sync::Arc;
232
233 use arrow_array::types::Int32Type;
234 use arrow_array::{ArrayRef, BooleanArray, Int32Array, Int64Array, RecordBatch, StructArray};
235 use arrow_buffer::NullBuffer;
236 use arrow_schema::{DataType, Field, Fields};
237 use arrow_select::concat::concat_batches;
238 use itertools::Itertools;
239 use parquet::arrow::PARQUET_FIELD_ID_META_KEY;
240 use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
241 use parquet::file::properties::WriterProperties;
242 use tempfile::TempDir;
243 use uuid::Uuid;
244
245 use crate::ErrorKind;
246 use crate::arrow::{arrow_schema_to_schema, schema_to_arrow_schema};
247 use crate::io::FileIO;
248 use crate::spec::{
249 DataFile, DataFileFormat, ListType, MapType, NestedField, PrimitiveType, Schema,
250 StructType, Type,
251 };
252 use crate::writer::base_writer::equality_delete_writer::{
253 EqualityDeleteFileWriterBuilder, EqualityDeleteWriterConfig,
254 };
255 use crate::writer::file_writer::ParquetWriterBuilder;
256 use crate::writer::file_writer::location_generator::{
257 DefaultFileNameGenerator, DefaultLocationGenerator,
258 };
259 use crate::writer::file_writer::rolling_writer::RollingFileWriterBuilder;
260 use crate::writer::{IcebergWriter, IcebergWriterBuilder};
261
262 async fn check_parquet_data_file_with_equality_delete_write(
263 file_io: &FileIO,
264 data_file: &DataFile,
265 batch: &RecordBatch,
266 ) {
267 assert_eq!(data_file.file_format, DataFileFormat::Parquet);
268
269 let input_file = file_io.new_input(data_file.file_path.clone()).unwrap();
271 let input_content = input_file.read().await.unwrap();
273 let reader_builder =
274 ParquetRecordBatchReaderBuilder::try_new(input_content.clone()).unwrap();
275 let metadata = reader_builder.metadata().clone();
276
277 let reader = reader_builder.build().unwrap();
279 let batches = reader.map(|batch| batch.unwrap()).collect::<Vec<_>>();
280 let res = concat_batches(&batch.schema(), &batches).unwrap();
281 assert_eq!(*batch, res);
282
283 let expect_column_num = batch.num_columns();
285
286 assert_eq!(
287 data_file.record_count,
288 metadata
289 .row_groups()
290 .iter()
291 .map(|group| group.num_rows())
292 .sum::<i64>() as u64
293 );
294
295 assert_eq!(data_file.file_size_in_bytes, input_content.len() as u64);
296
297 assert_eq!(data_file.column_sizes.len(), expect_column_num);
298
299 for (index, id) in data_file.column_sizes().keys().sorted().enumerate() {
300 metadata
301 .row_groups()
302 .iter()
303 .map(|group| group.columns())
304 .for_each(|column| {
305 assert_eq!(
306 *data_file.column_sizes.get(id).unwrap() as i64,
307 column.get(index).unwrap().compressed_size()
308 );
309 });
310 }
311
312 assert_eq!(data_file.value_counts.len(), expect_column_num);
313 data_file.value_counts.iter().for_each(|(_, &v)| {
314 let expect = metadata
315 .row_groups()
316 .iter()
317 .map(|group| group.num_rows())
318 .sum::<i64>() as u64;
319 assert_eq!(v, expect);
320 });
321
322 for (index, id) in data_file.null_value_counts().keys().enumerate() {
323 let expect = batch.column(index).null_count() as u64;
324 assert_eq!(*data_file.null_value_counts.get(id).unwrap(), expect);
325 }
326
327 let split_offsets = data_file
328 .split_offsets
329 .as_ref()
330 .expect("split_offsets should be set");
331 assert_eq!(split_offsets.len(), metadata.num_row_groups());
332 split_offsets.iter().enumerate().for_each(|(i, &v)| {
333 let expect = metadata.row_groups()[i].file_offset().unwrap();
334 assert_eq!(v, expect);
335 });
336 }
337
338 #[tokio::test]
339 async fn test_equality_delete_writer() -> Result<(), anyhow::Error> {
340 let temp_dir = TempDir::new().unwrap();
341 let file_io = FileIO::new_with_fs();
342 let location_gen = DefaultLocationGenerator::with_data_location(
343 temp_dir.path().to_str().unwrap().to_string(),
344 );
345 let file_name_gen =
346 DefaultFileNameGenerator::new("test".to_string(), None, DataFileFormat::Parquet);
347
348 let schema = Schema::builder()
351 .with_schema_id(1)
352 .with_fields(vec![
353 NestedField::required(0, "col0", Type::Primitive(PrimitiveType::Int)).into(),
354 NestedField::required(
355 1,
356 "col1",
357 Type::Struct(StructType::new(vec![
358 NestedField::required(5, "sub_col", Type::Primitive(PrimitiveType::Int))
359 .into(),
360 ])),
361 )
362 .into(),
363 NestedField::required(2, "col2", Type::Primitive(PrimitiveType::String)).into(),
364 NestedField::required(
365 3,
366 "col3",
367 Type::List(ListType::new(
368 NestedField::required(6, "element", Type::Primitive(PrimitiveType::Int))
369 .into(),
370 )),
371 )
372 .into(),
373 NestedField::required(
374 4,
375 "col4",
376 Type::Struct(StructType::new(vec![
377 NestedField::required(
378 7,
379 "sub_col",
380 Type::Struct(StructType::new(vec![
381 NestedField::required(
382 8,
383 "sub_sub_col",
384 Type::Primitive(PrimitiveType::Int),
385 )
386 .into(),
387 ])),
388 )
389 .into(),
390 ])),
391 )
392 .into(),
393 ])
394 .build()
395 .unwrap();
396 let arrow_schema = Arc::new(schema_to_arrow_schema(&schema).unwrap());
397 let col0 = Arc::new(Int32Array::from_iter_values(vec![1; 1024])) as ArrayRef;
398 let col1 = Arc::new(StructArray::new(
399 if let DataType::Struct(fields) = arrow_schema.fields.get(1).unwrap().data_type() {
400 fields.clone()
401 } else {
402 unreachable!()
403 },
404 vec![Arc::new(Int32Array::from_iter_values(vec![1; 1024]))],
405 None,
406 ));
407 let col2 = Arc::new(arrow_array::StringArray::from_iter_values(vec![
408 "test";
409 1024
410 ])) as ArrayRef;
411 let col3 = Arc::new({
412 let list_parts = arrow_array::ListArray::from_iter_primitive::<Int32Type, _, _>(vec![
413 Some(
414 vec![Some(1),]
415 );
416 1024
417 ])
418 .into_parts();
419 arrow_array::ListArray::new(
420 if let DataType::List(field) = arrow_schema.fields.get(3).unwrap().data_type() {
421 field.clone()
422 } else {
423 unreachable!()
424 },
425 list_parts.1,
426 list_parts.2,
427 list_parts.3,
428 )
429 }) as ArrayRef;
430 let col4 = Arc::new(StructArray::new(
431 if let DataType::Struct(fields) = arrow_schema.fields.get(4).unwrap().data_type() {
432 fields.clone()
433 } else {
434 unreachable!()
435 },
436 vec![Arc::new(StructArray::new(
437 if let DataType::Struct(fields) = arrow_schema.fields.get(4).unwrap().data_type() {
438 if let DataType::Struct(fields) = fields.first().unwrap().data_type() {
439 fields.clone()
440 } else {
441 unreachable!()
442 }
443 } else {
444 unreachable!()
445 },
446 vec![Arc::new(Int32Array::from_iter_values(vec![1; 1024]))],
447 None,
448 ))],
449 None,
450 ));
451 let columns = vec![col0, col1, col2, col3, col4];
452 let to_write = RecordBatch::try_new(arrow_schema.clone(), columns).unwrap();
453
454 let equality_ids = vec![0_i32, 8];
455 let equality_config =
456 EqualityDeleteWriterConfig::new(equality_ids, Arc::new(schema)).unwrap();
457 let delete_schema =
458 arrow_schema_to_schema(equality_config.projected_arrow_schema_ref()).unwrap();
459 let projector = equality_config.projector.clone();
460
461 let pb =
463 ParquetWriterBuilder::new(WriterProperties::builder().build(), Arc::new(delete_schema));
464 let rolling_writer_builder = RollingFileWriterBuilder::new_with_default_file_size(
465 pb,
466 file_io.clone(),
467 location_gen,
468 file_name_gen,
469 );
470 let mut equality_delete_writer =
471 EqualityDeleteFileWriterBuilder::new(rolling_writer_builder, equality_config)
472 .build(None)
473 .await?;
474
475 equality_delete_writer.write(to_write.clone()).await?;
477 let res = equality_delete_writer.close().await?;
478 assert_eq!(res.len(), 1);
479 let data_file = res.into_iter().next().unwrap();
480
481 let to_write_projected = projector.project_batch(to_write)?;
483 check_parquet_data_file_with_equality_delete_write(
484 &file_io,
485 &data_file,
486 &to_write_projected,
487 )
488 .await;
489 Ok(())
490 }
491
492 fn equality_id_validation_schema() -> Arc<Schema> {
493 Arc::new(
494 Schema::builder()
495 .with_schema_id(1)
496 .with_fields(vec![
497 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
498 NestedField::optional(2, "name", Type::Primitive(PrimitiveType::String)).into(),
499 ])
500 .build()
501 .unwrap(),
502 )
503 }
504
505 #[test]
506 fn test_equality_delete_rejects_empty_equality_ids() {
507 let err =
508 EqualityDeleteWriterConfig::new(vec![], equality_id_validation_schema()).unwrap_err();
509
510 assert_eq!(err.kind(), ErrorKind::DataInvalid);
511 assert!(
512 err.to_string()
513 .contains("Equality delete field ids must not be empty."),
514 "{err}"
515 );
516 }
517
518 #[test]
519 fn test_equality_delete_rejects_duplicate_equality_ids() {
520 let err = EqualityDeleteWriterConfig::new(vec![1, 1], equality_id_validation_schema())
521 .unwrap_err();
522
523 assert_eq!(err.kind(), ErrorKind::DataInvalid);
524 assert!(
525 err.to_string()
526 .contains("Duplicate equality delete field id: 1"),
527 "{err}"
528 );
529 assert!(err.to_string().contains("field_id: 1"), "{err}");
530 }
531
532 #[test]
533 fn test_equality_delete_rejects_missing_equality_id() {
534 let err =
535 EqualityDeleteWriterConfig::new(vec![99], equality_id_validation_schema()).unwrap_err();
536
537 assert_eq!(err.kind(), ErrorKind::DataInvalid);
538 assert!(
539 err.to_string()
540 .contains("Invalid equality delete field id: 99"),
541 "{err}"
542 );
543 assert!(err.to_string().contains("field_id: 99"), "{err}");
544 }
545
546 #[tokio::test]
547 async fn test_equality_delete_unreachable_column() -> Result<(), anyhow::Error> {
548 let schema = Arc::new(
549 Schema::builder()
550 .with_schema_id(1)
551 .with_fields(vec![
552 NestedField::required(0, "col0", Type::Primitive(PrimitiveType::Float)).into(),
553 NestedField::required(1, "col1", Type::Primitive(PrimitiveType::Double)).into(),
554 NestedField::optional(2, "col2", Type::Primitive(PrimitiveType::Int)).into(),
555 NestedField::required(
556 3,
557 "col3",
558 Type::Struct(StructType::new(vec![
559 NestedField::required(
560 4,
561 "sub_col",
562 Type::Primitive(PrimitiveType::Int),
563 )
564 .into(),
565 ])),
566 )
567 .into(),
568 NestedField::optional(
569 5,
570 "col4",
571 Type::Struct(StructType::new(vec![
572 NestedField::required(
573 6,
574 "sub_col2",
575 Type::Primitive(PrimitiveType::Int),
576 )
577 .into(),
578 ])),
579 )
580 .into(),
581 NestedField::required(
582 7,
583 "col5",
584 Type::Map(MapType::new(
585 Arc::new(NestedField::required(
586 8,
587 "key",
588 Type::Primitive(PrimitiveType::String),
589 )),
590 Arc::new(NestedField::required(
591 9,
592 "value",
593 Type::Primitive(PrimitiveType::Int),
594 )),
595 )),
596 )
597 .into(),
598 NestedField::required(
599 10,
600 "col6",
601 Type::List(ListType::new(Arc::new(NestedField::required(
602 11,
603 "element",
604 Type::Primitive(PrimitiveType::Int),
605 )))),
606 )
607 .into(),
608 ])
609 .build()
610 .unwrap(),
611 );
612 assert!(EqualityDeleteWriterConfig::new(vec![0], schema.clone()).is_err());
614 assert!(EqualityDeleteWriterConfig::new(vec![1], schema.clone()).is_err());
615 assert!(EqualityDeleteWriterConfig::new(vec![3], schema.clone()).is_err());
617 assert!(EqualityDeleteWriterConfig::new(vec![4], schema.clone()).is_ok());
619 assert!(EqualityDeleteWriterConfig::new(vec![7], schema.clone()).is_err());
621 assert!(EqualityDeleteWriterConfig::new(vec![8], schema.clone()).is_err());
622 assert!(EqualityDeleteWriterConfig::new(vec![9], schema.clone()).is_err());
623 assert!(EqualityDeleteWriterConfig::new(vec![10], schema.clone()).is_err());
625 assert!(EqualityDeleteWriterConfig::new(vec![11], schema.clone()).is_err());
626
627 Ok(())
628 }
629
630 #[tokio::test]
631 async fn test_equality_delete_with_primitive_type() -> Result<(), anyhow::Error> {
632 let temp_dir = TempDir::new().unwrap();
633 let file_io = FileIO::new_with_fs();
634 let location_gen = DefaultLocationGenerator::with_data_location(
635 temp_dir.path().to_str().unwrap().to_string(),
636 );
637 let file_name_gen =
638 DefaultFileNameGenerator::new("test".to_string(), None, DataFileFormat::Parquet);
639
640 let schema = Arc::new(
641 Schema::builder()
642 .with_schema_id(1)
643 .with_fields(vec![
644 NestedField::required(0, "col0", Type::Primitive(PrimitiveType::Boolean))
645 .into(),
646 NestedField::required(1, "col1", Type::Primitive(PrimitiveType::Int)).into(),
647 NestedField::required(2, "col2", Type::Primitive(PrimitiveType::Long)).into(),
648 NestedField::required(
649 3,
650 "col3",
651 Type::Primitive(PrimitiveType::Decimal {
652 precision: 38,
653 scale: 5,
654 }),
655 )
656 .into(),
657 NestedField::required(4, "col4", Type::Primitive(PrimitiveType::Date)).into(),
658 NestedField::required(5, "col5", Type::Primitive(PrimitiveType::Time)).into(),
659 NestedField::required(6, "col6", Type::Primitive(PrimitiveType::Timestamp))
660 .into(),
661 NestedField::required(7, "col7", Type::Primitive(PrimitiveType::Timestamptz))
662 .into(),
663 NestedField::required(8, "col8", Type::Primitive(PrimitiveType::TimestampNs))
664 .into(),
665 NestedField::required(9, "col9", Type::Primitive(PrimitiveType::TimestamptzNs))
666 .into(),
667 NestedField::required(10, "col10", Type::Primitive(PrimitiveType::String))
668 .into(),
669 NestedField::required(11, "col11", Type::Primitive(PrimitiveType::Uuid)).into(),
670 NestedField::required(12, "col12", Type::Primitive(PrimitiveType::Fixed(10)))
671 .into(),
672 NestedField::required(13, "col13", Type::Primitive(PrimitiveType::Binary))
673 .into(),
674 ])
675 .build()
676 .unwrap(),
677 );
678 let equality_ids = vec![0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13];
679 let config = EqualityDeleteWriterConfig::new(equality_ids, schema.clone()).unwrap();
680 let delete_arrow_schema = config.projected_arrow_schema_ref().clone();
681 let delete_schema = arrow_schema_to_schema(&delete_arrow_schema).unwrap();
682
683 let pb =
684 ParquetWriterBuilder::new(WriterProperties::builder().build(), Arc::new(delete_schema));
685 let rolling_writer_builder = RollingFileWriterBuilder::new_with_default_file_size(
686 pb,
687 file_io.clone(),
688 location_gen,
689 file_name_gen,
690 );
691 let mut equality_delete_writer =
692 EqualityDeleteFileWriterBuilder::new(rolling_writer_builder, config)
693 .build(None)
694 .await?;
695
696 let col0 = Arc::new(BooleanArray::from(vec![
698 Some(true),
699 Some(false),
700 Some(true),
701 ])) as ArrayRef;
702 let col1 = Arc::new(Int32Array::from(vec![Some(1), Some(2), Some(4)])) as ArrayRef;
703 let col2 = Arc::new(Int64Array::from(vec![Some(1), Some(2), Some(4)])) as ArrayRef;
704 let col3 = Arc::new(
705 arrow_array::Decimal128Array::from(vec![Some(1), Some(2), Some(4)])
706 .with_precision_and_scale(38, 5)
707 .unwrap(),
708 ) as ArrayRef;
709 let col4 = Arc::new(arrow_array::Date32Array::from(vec![
710 Some(0),
711 Some(1),
712 Some(3),
713 ])) as ArrayRef;
714 let col5 = Arc::new(arrow_array::Time64MicrosecondArray::from(vec![
715 Some(0),
716 Some(1),
717 Some(3),
718 ])) as ArrayRef;
719 let col6 = Arc::new(arrow_array::TimestampMicrosecondArray::from(vec![
720 Some(0),
721 Some(1),
722 Some(3),
723 ])) as ArrayRef;
724 let col7 = Arc::new(
725 arrow_array::TimestampMicrosecondArray::from(vec![Some(0), Some(1), Some(3)])
726 .with_timezone_utc(),
727 ) as ArrayRef;
728 let col8 = Arc::new(arrow_array::TimestampNanosecondArray::from(vec![
729 Some(0),
730 Some(1),
731 Some(3),
732 ])) as ArrayRef;
733 let col9 = Arc::new(
734 arrow_array::TimestampNanosecondArray::from(vec![Some(0), Some(1), Some(3)])
735 .with_timezone_utc(),
736 ) as ArrayRef;
737 let col10 = Arc::new(arrow_array::StringArray::from(vec![
738 Some("a"),
739 Some("b"),
740 Some("d"),
741 ])) as ArrayRef;
742 let col11 = Arc::new(
743 arrow_array::FixedSizeBinaryArray::try_from_sparse_iter_with_size(
744 vec![
745 Some(Uuid::from_u128(0).as_bytes().to_vec()),
746 Some(Uuid::from_u128(1).as_bytes().to_vec()),
747 Some(Uuid::from_u128(3).as_bytes().to_vec()),
748 ]
749 .into_iter(),
750 16,
751 )
752 .unwrap(),
753 ) as ArrayRef;
754 let col12 = Arc::new(
755 arrow_array::FixedSizeBinaryArray::try_from_sparse_iter_with_size(
756 vec![
757 Some(vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10]),
758 Some(vec![11, 12, 13, 14, 15, 16, 17, 18, 19, 20]),
759 Some(vec![21, 22, 23, 24, 25, 26, 27, 28, 29, 30]),
760 ]
761 .into_iter(),
762 10,
763 )
764 .unwrap(),
765 ) as ArrayRef;
766 let col13 = Arc::new(arrow_array::LargeBinaryArray::from_opt_vec(vec![
767 Some(b"one"),
768 Some(b""),
769 Some(b"zzzz"),
770 ])) as ArrayRef;
771 let to_write = RecordBatch::try_new(delete_arrow_schema.clone(), vec![
772 col0, col1, col2, col3, col4, col5, col6, col7, col8, col9, col10, col11, col12, col13,
773 ])
774 .unwrap();
775 equality_delete_writer.write(to_write.clone()).await?;
776 let res = equality_delete_writer.close().await?;
777 assert_eq!(res.len(), 1);
778 check_parquet_data_file_with_equality_delete_write(
779 &file_io,
780 &res.into_iter().next().unwrap(),
781 &to_write,
782 )
783 .await;
784
785 Ok(())
786 }
787
788 #[tokio::test]
789 async fn test_equality_delete_with_nullable_field() -> Result<(), anyhow::Error> {
790 let schema = Schema::builder()
793 .with_schema_id(1)
794 .with_fields(vec![
795 NestedField::optional(0, "col0", Type::Primitive(PrimitiveType::Int)).into(),
796 NestedField::optional(
797 1,
798 "col1",
799 Type::Struct(StructType::new(vec![
800 NestedField::optional(2, "sub_col", Type::Primitive(PrimitiveType::Int))
801 .into(),
802 ])),
803 )
804 .into(),
805 NestedField::optional(
806 3,
807 "col2",
808 Type::Struct(StructType::new(vec![
809 NestedField::optional(
810 4,
811 "sub_struct_col",
812 Type::Struct(StructType::new(vec![
813 NestedField::optional(
814 5,
815 "sub_sub_col",
816 Type::Primitive(PrimitiveType::Int),
817 )
818 .into(),
819 ])),
820 )
821 .into(),
822 ])),
823 )
824 .into(),
825 ])
826 .build()
827 .unwrap();
828 let arrow_schema = Arc::new(schema_to_arrow_schema(&schema).unwrap());
829 let col0 = Arc::new(Int32Array::from(vec![None, Some(2), Some(3)])) as ArrayRef;
833 let col1 = {
834 let nulls = NullBuffer::from(vec![true, false, true]);
835 Arc::new(StructArray::new(
836 if let DataType::Struct(fields) = arrow_schema.fields.get(1).unwrap().data_type() {
837 fields.clone()
838 } else {
839 unreachable!()
840 },
841 vec![Arc::new(Int32Array::from(vec![Some(1), Some(2), None]))],
842 Some(nulls),
843 ))
844 };
845 let col2 = {
846 let inner_col = {
847 let nulls = NullBuffer::from(vec![true, false, true]);
848 Arc::new(StructArray::new(
849 Fields::from(vec![
850 Field::new("sub_sub_col", DataType::Int32, true).with_metadata(
851 HashMap::from([(
852 PARQUET_FIELD_ID_META_KEY.to_string(),
853 "5".to_string(),
854 )]),
855 ),
856 ]),
857 vec![Arc::new(Int32Array::from(vec![Some(1), Some(2), None]))],
858 Some(nulls),
859 ))
860 };
861 let nulls = NullBuffer::from(vec![false, true, true]);
862 Arc::new(StructArray::new(
863 if let DataType::Struct(fields) = arrow_schema.fields.get(2).unwrap().data_type() {
864 fields.clone()
865 } else {
866 unreachable!()
867 },
868 vec![inner_col],
869 Some(nulls),
870 ))
871 };
872 let columns = vec![col0, col1, col2];
873
874 let to_write = RecordBatch::try_new(arrow_schema.clone(), columns).unwrap();
875 let equality_ids = vec![0_i32, 2, 5];
876 let equality_config =
877 EqualityDeleteWriterConfig::new(equality_ids, Arc::new(schema)).unwrap();
878 let projector = equality_config.projector.clone();
879
880 let to_write_projected = projector.project_batch(to_write)?;
882 let expect_batch =
883 RecordBatch::try_new(equality_config.projected_arrow_schema_ref().clone(), vec![
884 Arc::new(Int32Array::from(vec![None, Some(2), Some(3)])) as ArrayRef,
885 Arc::new(Int32Array::from(vec![Some(1), None, None])) as ArrayRef,
886 Arc::new(Int32Array::from(vec![None, None, None])) as ArrayRef,
887 ])
888 .unwrap();
889 assert_eq!(to_write_projected, expect_batch);
890 Ok(())
891 }
892}