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::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/// Builder for `EqualityDeleteWriter`.
38#[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    /// Create a new `EqualityDeleteFileWriterBuilder` using a `RollingFileWriterBuilder`.
55    pub fn new(
56        inner: RollingFileWriterBuilder<B, L, F>,
57        config: EqualityDeleteWriterConfig,
58    ) -> Self {
59        Self { inner, config }
60    }
61}
62
63/// Config for `EqualityDeleteWriter`.
64#[derive(Debug)]
65pub struct EqualityDeleteWriterConfig {
66    // Field ids used to determine row equality in equality delete files.
67    equality_ids: Vec<i32>,
68    // Projector used to project the data chunk into specific fields.
69    projector: RecordBatchProjector,
70}
71
72/// Validates equality delete field ids before constructing the writer config.
73///
74/// In addition to rejecting an empty id list, this check also rejects duplicate
75/// ids and ids that do not exist in the table schema. These two checks are
76/// intentionally stricter than the current Java reference implementation, which
77/// only guards against null/empty lists; catching duplicates and missing ids
78/// here produces clearer errors and avoids generating meaningless delete files.
79fn 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    /// Create a new `EqualityDeleteWriterConfig` with equality ids.
111    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            // Equality delete fields follow identifier-field type restrictions,
119            // except optional columns and columns nested under optional structs are allowed.
120            // Project only primitive, non-floating fields; RecordBatchProjector's traversal
121            // keeps fields under maps and lists unreachable.
122            |field| {
123                // Only primitive type is allowed to be used for identifier field ids
124                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    /// Return projected Schema
152    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/// Writer used to write equality delete files.
177#[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        // read the written file
280        let input_file = file_io.new_input(data_file.file_path.clone()).unwrap();
281        // read the written file
282        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        // check data
288        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        // check metadata
294        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        // prepare data
359        // Int, Struct(Int), String, List(Int), Struct(Struct(Int))
360        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        // prepare writer
472        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        // write
486        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        // check
492        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        // Float and Double are not allowed to be used for equality delete
623        assert!(EqualityDeleteWriterConfig::new(vec![0], schema.clone()).is_err());
624        assert!(EqualityDeleteWriterConfig::new(vec![1], schema.clone()).is_err());
625        // Struct is not allowed to be used for equality delete
626        assert!(EqualityDeleteWriterConfig::new(vec![3], schema.clone()).is_err());
627        // Nested field of struct is allowed to be used for equality delete
628        assert!(EqualityDeleteWriterConfig::new(vec![4], schema.clone()).is_ok());
629        // Nested field of map is not allowed to be used for equality delete
630        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        // Nested field of list is not allowed to be used for equality delete
634        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        // prepare data
707        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        // prepare data
801        // Int, Struct(Int), Struct(Struct(Int))
802        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        // null 1            null(struct)
840        // 2    null(struct) null(sub_struct_col)
841        // 3    null(field)  null(sub_sub_col)
842        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        // check
891        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}