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