Skip to main content

iceberg/writer/base_writer/
equality_delete_writer.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18//! This module provide `EqualityDeleteWriter`.
19
20use 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/// Builder for `EqualityDeleteWriter`.
39#[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    /// Create a new `EqualityDeleteFileWriterBuilder` using a `RollingFileWriterBuilder`.
56    pub fn new(
57        inner: RollingFileWriterBuilder<B, L, F>,
58        config: EqualityDeleteWriterConfig,
59    ) -> Self {
60        Self { inner, config }
61    }
62}
63
64/// Config for `EqualityDeleteWriter`.
65#[derive(Debug)]
66pub struct EqualityDeleteWriterConfig {
67    // Field ids used to determine row equality in equality delete files.
68    equality_ids: Vec<i32>,
69    // Projector used to project the data chunk into specific fields.
70    projector: RecordBatchProjector,
71}
72
73/// Validates equality delete field ids before constructing the writer config.
74///
75/// In addition to rejecting an empty id list, this check also rejects duplicate
76/// ids and ids that do not exist in the table schema. These two checks are
77/// intentionally stricter than the current Java reference implementation, which
78/// only guards against null/empty lists; catching duplicates and missing ids
79/// here produces clearer errors and avoids generating meaningless delete files.
80fn 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    /// Create a new `EqualityDeleteWriterConfig` with equality ids.
105    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            // Equality delete fields follow identifier-field type restrictions,
113            // except optional columns and columns nested under optional structs are allowed.
114            // Project only primitive, non-floating fields; RecordBatchProjector's traversal
115            // keeps fields under maps and lists unreachable.
116            |field| {
117                // Only primitive type is allowed to be used for identifier field ids
118                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    /// Return projected Schema
146    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/// Writer used to write equality delete files.
171#[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        // read the written file
270        let input_file = file_io.new_input(data_file.file_path.clone()).unwrap();
271        // read the written file
272        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        // check data
278        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        // check metadata
284        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        // prepare data
349        // Int, Struct(Int), String, List(Int), Struct(Struct(Int))
350        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        // prepare writer
462        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        // write
476        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        // check
482        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        // Float and Double are not allowed to be used for equality delete
613        assert!(EqualityDeleteWriterConfig::new(vec![0], schema.clone()).is_err());
614        assert!(EqualityDeleteWriterConfig::new(vec![1], schema.clone()).is_err());
615        // Struct is not allowed to be used for equality delete
616        assert!(EqualityDeleteWriterConfig::new(vec![3], schema.clone()).is_err());
617        // Nested field of struct is allowed to be used for equality delete
618        assert!(EqualityDeleteWriterConfig::new(vec![4], schema.clone()).is_ok());
619        // Nested field of map is not allowed to be used for equality delete
620        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        // Nested field of list is not allowed to be used for equality delete
624        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        // prepare data
697        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        // prepare data
791        // Int, Struct(Int), Struct(Struct(Int))
792        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        // null 1            null(struct)
830        // 2    null(struct) null(sub_struct_col)
831        // 3    null(field)  null(sub_sub_col)
832        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        // check
881        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}