Skip to main content

iceberg/writer/base_writer/
position_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 provides `PositionDeleteFileWriter`.
19//!
20//! A position delete file has two required columns: `file_path` (`string`, field id
21//! [`RESERVED_FIELD_ID_DELETE_FILE_PATH`]) and `pos` (`long`, field id
22//! [`RESERVED_FIELD_ID_DELETE_FILE_POS`]). The spec also allows an optional `row` column
23//! that inlines the deleted row's values; this writer does not support it yet, so batches
24//! must contain only the two required columns. The writer takes batches already shaped as
25//! those two columns (see [`position_delete_schema`]) and sets
26//! [`DataContentType::PositionDeletes`] on the output. It does not sort its input; see
27//! [`PositionDeleteFileWriter::write`].
28//!
29//! Position delete files are a v2 construct. v3 replaces them with deletion vectors and
30//! forbids adding new position delete files, so callers must not route v3 writes here.
31//! This base writer has no format-version gate by design; that gating belongs at the
32//! transaction/commit layer.
33
34use std::sync::Arc;
35
36use arrow_array::RecordBatch;
37use arrow_schema::{DataType, Field};
38use once_cell::sync::Lazy;
39use parquet::arrow::PARQUET_FIELD_ID_META_KEY;
40
41use crate::error::invalid_data;
42use crate::metadata_columns::{
43    RESERVED_FIELD_ID_DELETE_FILE_PATH, RESERVED_FIELD_ID_DELETE_FILE_POS, delete_file_path_field,
44    delete_file_pos_field,
45};
46use crate::spec::{DataContentType, DataFile, PartitionKey, Schema, SchemaRef};
47use crate::writer::file_writer::FileWriterBuilder;
48use crate::writer::file_writer::location_generator::{FileNameGenerator, LocationGenerator};
49use crate::writer::file_writer::rolling_writer::{RollingFileWriter, RollingFileWriterBuilder};
50use crate::writer::{IcebergWriter, IcebergWriterBuilder};
51use crate::{Error, ErrorKind, Result};
52
53/// The canonical Iceberg schema of a position delete file: the required `file_path`
54/// (`string`) and `pos` (`long`) columns with their reserved field ids.
55static POSITION_DELETE_SCHEMA: Lazy<SchemaRef> = Lazy::new(|| {
56    Arc::new(
57        Schema::builder()
58            .with_fields(vec![
59                delete_file_path_field().clone(),
60                delete_file_pos_field().clone(),
61            ])
62            .build()
63            .expect("position delete schema is statically valid"),
64    )
65});
66
67/// [`POSITION_DELETE_SCHEMA`] converted to Arrow, keeping the reserved field ids in each
68/// field's Parquet field-id metadata. This is the crate-internal Arrow projection of a
69/// position delete file, shared with
70/// [`PositionDeletes`](super::position_delete_input::PositionDeletes).
71static POSITION_DELETE_ARROW_SCHEMA: Lazy<arrow_schema::SchemaRef> = Lazy::new(|| {
72    Arc::new(
73        crate::arrow::schema_to_arrow_schema(&POSITION_DELETE_SCHEMA)
74            .expect("position delete arrow schema is statically valid"),
75    )
76});
77
78/// Returns the canonical Iceberg schema of a position delete file.
79///
80/// Use this to build the [`ParquetWriterBuilder`](crate::writer::file_writer::ParquetWriterBuilder)
81/// that backs a [`PositionDeleteFileWriter`], so the written file matches the
82/// spec exactly.
83pub fn position_delete_schema() -> SchemaRef {
84    POSITION_DELETE_SCHEMA.clone()
85}
86
87/// Returns the canonical Arrow schema of a position delete file.
88///
89/// A cheap [`Arc`] clone of the shared static, used by
90/// [`PositionDeletes`](super::position_delete_input::PositionDeletes).
91pub(crate) fn position_delete_arrow_schema() -> arrow_schema::SchemaRef {
92    POSITION_DELETE_ARROW_SCHEMA.clone()
93}
94
95/// Reads a field's Iceberg field id from its Parquet field-id metadata.
96fn field_id(field: &Field) -> Result<i32> {
97    field
98        .metadata()
99        .get(PARQUET_FIELD_ID_META_KEY)
100        .ok_or_else(|| {
101            invalid_data!(
102                "Position delete column `{}` is missing its Iceberg field id metadata.",
103                field.name()
104            )
105        })?
106        .parse::<i32>()
107        .map_err(|e| {
108            invalid_data!(
109                "Position delete column `{}` has an invalid field id: {e}",
110                field.name()
111            )
112        })
113}
114
115/// Validates that a batch is a position delete file: the `file_path` (`Utf8`) and
116/// `pos` (`Int64`) columns, in order, with the two reserved field ids. Checking it
117/// here gives a clear error before the batch reaches the Parquet writer.
118fn validate_position_delete_batch(batch: &RecordBatch) -> Result<()> {
119    let fields = batch.schema_ref().fields();
120    if fields.len() != 2 {
121        return Err(invalid_data!(
122            "This writer supports only the two required position delete columns (`file_path`, `pos`); \
123             batches with a different column count (e.g. including the optional `row` column) are not supported. Got {} columns.",
124            fields.len()
125        ));
126    }
127
128    let path = &fields[0];
129    let path_id = field_id(path)?;
130    if path_id != RESERVED_FIELD_ID_DELETE_FILE_PATH {
131        return Err(invalid_data!(
132            "The first position delete column must be `file_path` (field id {RESERVED_FIELD_ID_DELETE_FILE_PATH}), but got field id {path_id}."
133        ));
134    }
135    // The canonical schema maps Iceberg `string` to `Utf8` and the file writer is
136    // configured with it, so a `LargeUtf8` column has to be cast to `Utf8` first.
137    if path.data_type() != &DataType::Utf8 {
138        return Err(invalid_data!(
139            "The position delete `file_path` column must be Utf8 (cast it first); got {:?}.",
140            path.data_type()
141        ));
142    }
143    // Required column: a nullable field could write nulls under a required schema.
144    if path.is_nullable() {
145        return Err(invalid_data!(
146            "The position delete `file_path` column must be required (non-nullable)."
147        ));
148    }
149
150    let pos = &fields[1];
151    let pos_id = field_id(pos)?;
152    if pos_id != RESERVED_FIELD_ID_DELETE_FILE_POS {
153        return Err(invalid_data!(
154            "The second position delete column must be `pos` (field id {RESERVED_FIELD_ID_DELETE_FILE_POS}), but got field id {pos_id}."
155        ));
156    }
157    if pos.data_type() != &DataType::Int64 {
158        return Err(invalid_data!(
159            "The position delete `pos` column must be Int64, but got {:?}.",
160            pos.data_type()
161        ));
162    }
163    if pos.is_nullable() {
164        return Err(invalid_data!(
165            "The position delete `pos` column must be required (non-nullable)."
166        ));
167    }
168
169    Ok(())
170}
171
172/// Builder for [`PositionDeleteFileWriter`].
173#[derive(Debug)]
174pub struct PositionDeleteFileWriterBuilder<
175    B: FileWriterBuilder,
176    L: LocationGenerator,
177    F: FileNameGenerator,
178> {
179    inner: RollingFileWriterBuilder<B, L, F>,
180}
181
182impl<B, L, F> PositionDeleteFileWriterBuilder<B, L, F>
183where
184    B: FileWriterBuilder,
185    L: LocationGenerator,
186    F: FileNameGenerator,
187{
188    /// Create a new `PositionDeleteFileWriterBuilder` using a `RollingFileWriterBuilder`.
189    ///
190    /// The `RollingFileWriterBuilder` must be backed by a file writer configured
191    /// with the [`position_delete_schema`]; the per-batch validation in
192    /// [`PositionDeleteFileWriter::write`] guards against a mismatched batch, but
193    /// the caller is responsible for wiring the same schema into the file writer.
194    pub fn new(inner: RollingFileWriterBuilder<B, L, F>) -> Self {
195        Self { inner }
196    }
197}
198
199#[async_trait::async_trait]
200impl<B, L, F> IcebergWriterBuilder for PositionDeleteFileWriterBuilder<B, L, F>
201where
202    B: FileWriterBuilder,
203    L: LocationGenerator,
204    F: FileNameGenerator,
205{
206    type R = PositionDeleteFileWriter<B, L, F>;
207
208    async fn build(&self, partition_key: Option<PartitionKey>) -> Result<Self::R> {
209        Ok(PositionDeleteFileWriter {
210            inner: Some(self.inner.build()),
211            partition_key,
212        })
213    }
214}
215
216/// Writer used to write position delete files within one spec/partition.
217#[derive(Debug)]
218pub struct PositionDeleteFileWriter<
219    B: FileWriterBuilder,
220    L: LocationGenerator,
221    F: FileNameGenerator,
222> {
223    inner: Option<RollingFileWriter<B, L, F>>,
224    partition_key: Option<PartitionKey>,
225}
226
227#[async_trait::async_trait]
228impl<B, L, F> IcebergWriter for PositionDeleteFileWriter<B, L, F>
229where
230    B: FileWriterBuilder,
231    L: LocationGenerator,
232    F: FileNameGenerator,
233{
234    /// Writes a batch of `(file_path, pos)` records; the shape is validated on every
235    /// call.
236    ///
237    /// The writer does not sort its input. Position delete files must be sorted by
238    /// `file_path` then `pos`, so the caller must supply rows in that order across all
239    /// `write` calls; a sorting writer will remove this requirement.
240    async fn write(&mut self, batch: RecordBatch) -> Result<()> {
241        // Reject a closed writer before validating the batch.
242        let Some(writer) = self.inner.as_mut() else {
243            return Err(Error::new(
244                ErrorKind::Unexpected,
245                "Position delete writer is already closed; cannot write.",
246            ));
247        };
248        validate_position_delete_batch(&batch)?;
249        writer.write(&self.partition_key, &batch).await
250    }
251
252    async fn close(&mut self) -> Result<Vec<DataFile>> {
253        if let Some(writer) = self.inner.take() {
254            writer
255                .close()
256                .await?
257                .into_iter()
258                .map(|mut res| {
259                    res.content(DataContentType::PositionDeletes);
260                    // sort_order_id stays null, as the spec requires for position deletes.
261                    if let Some(pk) = self.partition_key.as_ref() {
262                        res.partition(pk.data().clone());
263                        res.partition_spec_id(pk.spec().spec_id());
264                    }
265                    res.build()
266                        .map_err(|e| invalid_data!("Failed to build position delete file: {e}"))
267                })
268                .collect()
269        } else {
270            Err(Error::new(
271                ErrorKind::Unexpected,
272                "Position delete writer is already closed.",
273            ))
274        }
275    }
276}
277
278#[cfg(test)]
279mod test {
280    use std::collections::HashMap;
281    use std::sync::Arc;
282
283    use arrow_array::{Int32Array, Int64Array, LargeStringArray, RecordBatch, StringArray};
284    use arrow_schema::{DataType, Field};
285    use arrow_select::concat::concat_batches;
286    use parquet::arrow::PARQUET_FIELD_ID_META_KEY;
287    use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
288    use parquet::file::properties::WriterProperties;
289    use tempfile::TempDir;
290
291    use super::*;
292    use crate::io::FileIO;
293    use crate::metadata_columns::{
294        RESERVED_COL_NAME_DELETE_FILE_PATH, RESERVED_COL_NAME_DELETE_FILE_POS,
295    };
296    use crate::spec::{
297        DataFileFormat, Literal, NestedField, PartitionSpec, PrimitiveType, Struct, Transform, Type,
298    };
299    use crate::writer::file_writer::ParquetWriterBuilder;
300    use crate::writer::file_writer::location_generator::{
301        DefaultFileNameGenerator, DefaultLocationGenerator,
302    };
303    use crate::writer::file_writer::rolling_writer::RollingFileWriterBuilder;
304
305    #[test]
306    fn test_position_delete_schema_shape() {
307        let schema = position_delete_schema();
308        let fields = schema.as_struct().fields();
309        assert_eq!(fields.len(), 2);
310
311        assert_eq!(fields[0].id, RESERVED_FIELD_ID_DELETE_FILE_PATH);
312        assert_eq!(fields[0].name, RESERVED_COL_NAME_DELETE_FILE_PATH);
313        assert!(fields[0].required);
314        assert_eq!(
315            fields[0].field_type.as_ref(),
316            &Type::Primitive(PrimitiveType::String)
317        );
318
319        assert_eq!(fields[1].id, RESERVED_FIELD_ID_DELETE_FILE_POS);
320        assert_eq!(fields[1].name, RESERVED_COL_NAME_DELETE_FILE_POS);
321        assert!(fields[1].required);
322        assert_eq!(
323            fields[1].field_type.as_ref(),
324            &Type::Primitive(PrimitiveType::Long)
325        );
326
327        // The Arrow projection carries the reserved field ids and non-null flags.
328        let arrow_schema = position_delete_arrow_schema();
329        assert_eq!(arrow_schema.fields().len(), 2);
330        assert_eq!(arrow_schema.field(0).data_type(), &DataType::Utf8);
331        assert_eq!(arrow_schema.field(1).data_type(), &DataType::Int64);
332        assert!(!arrow_schema.field(0).is_nullable());
333        assert!(!arrow_schema.field(1).is_nullable());
334    }
335
336    fn position_delete_batch(paths: Vec<&str>, positions: Vec<i64>) -> RecordBatch {
337        RecordBatch::try_new(position_delete_arrow_schema(), vec![
338            Arc::new(StringArray::from(paths)),
339            Arc::new(Int64Array::from(positions)),
340        ])
341        .unwrap()
342    }
343
344    /// A field carrying an explicit Iceberg field-id metadata entry.
345    fn field_with_id(name: &str, data_type: DataType, field_id: i32) -> Field {
346        Field::new(name, data_type, false).with_metadata(HashMap::from([(
347            PARQUET_FIELD_ID_META_KEY.to_string(),
348            field_id.to_string(),
349        )]))
350    }
351
352    fn writer_setup(
353        temp_dir: &TempDir,
354    ) -> (
355        FileIO,
356        PositionDeleteFileWriterBuilder<
357            ParquetWriterBuilder,
358            DefaultLocationGenerator,
359            DefaultFileNameGenerator,
360        >,
361    ) {
362        let file_io = FileIO::new_with_fs();
363        let location_gen = DefaultLocationGenerator::with_data_location(
364            temp_dir.path().to_str().unwrap().to_string(),
365        );
366        let file_name_gen =
367            DefaultFileNameGenerator::new("test".to_string(), None, DataFileFormat::Parquet);
368
369        let parquet_writer_builder = ParquetWriterBuilder::new(
370            WriterProperties::builder().build(),
371            position_delete_schema(),
372        );
373        let rolling_writer_builder = RollingFileWriterBuilder::new_with_default_file_size(
374            parquet_writer_builder,
375            file_io.clone(),
376            location_gen,
377            file_name_gen,
378        );
379        (
380            file_io,
381            PositionDeleteFileWriterBuilder::new(rolling_writer_builder),
382        )
383    }
384
385    #[tokio::test]
386    async fn test_position_delete_writer_round_trip() -> Result<()> {
387        let temp_dir = TempDir::new().unwrap();
388        let (file_io, builder) = writer_setup(&temp_dir);
389        let mut writer = builder.build(None).await?;
390
391        // Sorted by (file_path, pos): f0/1, f0/4, then f1/2.
392        let batch = position_delete_batch(
393            vec![
394                "s3://bucket/data/f0.parquet",
395                "s3://bucket/data/f0.parquet",
396                "s3://bucket/data/f1.parquet",
397            ],
398            vec![1, 4, 2],
399        );
400        writer.write(batch.clone()).await?;
401        let data_files = writer.close().await?;
402
403        assert_eq!(data_files.len(), 1);
404        let data_file = &data_files[0];
405        assert_eq!(data_file.content_type(), DataContentType::PositionDeletes);
406        assert_eq!(data_file.file_format, DataFileFormat::Parquet);
407        assert_eq!(data_file.record_count, 3);
408        // Unpartitioned writer leaves the default (empty) partition / spec id.
409        assert_eq!(data_file.partition, Struct::empty());
410        assert_eq!(data_file.partition_spec_id, 0);
411        // The rolling writer fills in file statistics.
412        assert!(data_file.file_size_in_bytes > 0);
413
414        // The written Parquet file round-trips back to the exact input rows.
415        let read_back = read_back_single(&file_io, data_file, &batch.schema()).await;
416        assert_eq!(read_back, batch);
417
418        Ok(())
419    }
420
421    #[tokio::test]
422    async fn test_position_delete_writer_parquet_field_ids() -> Result<()> {
423        let temp_dir = TempDir::new().unwrap();
424        let (file_io, builder) = writer_setup(&temp_dir);
425        let mut writer = builder.build(None).await?;
426        writer
427            .write(position_delete_batch(
428                vec!["s3://bucket/data/f0.parquet"],
429                vec![1],
430            ))
431            .await?;
432        let data_files = writer.close().await?;
433
434        // The reserved field ids must survive into the written Parquet schema, not just
435        // the in-memory Arrow schema, so cross-engine readers resolve the columns by id.
436        let content = file_io
437            .new_input(data_files[0].file_path.clone())?
438            .read()
439            .await?;
440        let reader = ParquetRecordBatchReaderBuilder::try_new(content).unwrap();
441        let field_ids: Vec<i32> = reader
442            .parquet_schema()
443            .columns()
444            .iter()
445            .map(|col| col.self_type().get_basic_info().id())
446            .collect();
447        assert_eq!(field_ids, vec![
448            RESERVED_FIELD_ID_DELETE_FILE_PATH,
449            RESERVED_FIELD_ID_DELETE_FILE_POS,
450        ]);
451
452        Ok(())
453    }
454
455    #[tokio::test]
456    async fn test_position_delete_writer_multiple_writes() -> Result<()> {
457        let temp_dir = TempDir::new().unwrap();
458        let (file_io, builder) = writer_setup(&temp_dir);
459        let mut writer = builder.build(None).await?;
460
461        let batch1 = position_delete_batch(
462            vec!["s3://bucket/data/f0.parquet", "s3://bucket/data/f0.parquet"],
463            vec![1, 4],
464        );
465        let batch2 = position_delete_batch(
466            vec!["s3://bucket/data/f1.parquet", "s3://bucket/data/f1.parquet"],
467            vec![2, 7],
468        );
469        writer.write(batch1.clone()).await?;
470        writer.write(batch2.clone()).await?;
471        let data_files = writer.close().await?;
472
473        assert_eq!(data_files.len(), 1);
474        assert_eq!(data_files[0].record_count, 4);
475
476        let expected = concat_batches(&batch1.schema(), [&batch1, &batch2]).unwrap();
477        let read_back = read_back_single(&file_io, &data_files[0], &batch1.schema()).await;
478        assert_eq!(read_back, expected);
479
480        Ok(())
481    }
482
483    #[tokio::test]
484    async fn test_position_delete_writer_sets_partition() -> Result<()> {
485        let temp_dir = TempDir::new().unwrap();
486        let (_file_io, builder) = writer_setup(&temp_dir);
487
488        // A table schema + identity partition spec with a non-default spec id, so the
489        // assertions distinguish real propagation from the DataFileBuilder defaults.
490        let table_schema = Arc::new(
491            Schema::builder()
492                .with_schema_id(1)
493                .with_fields(vec![
494                    NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
495                ])
496                .build()?,
497        );
498        let spec = PartitionSpec::builder(table_schema.clone())
499            .with_spec_id(7)
500            .add_partition_field("id", "id", Transform::Identity)?
501            .build()?;
502        let partition_value = Struct::from_iter([Some(Literal::int(42))]);
503        let partition_key = PartitionKey::new(spec, table_schema.clone(), partition_value.clone());
504
505        let mut writer = builder.build(Some(partition_key)).await?;
506        writer
507            .write(position_delete_batch(
508                vec!["s3://bucket/data/f0.parquet"],
509                vec![1],
510            ))
511            .await?;
512        let data_files = writer.close().await?;
513
514        assert_eq!(data_files.len(), 1);
515        let data_file = &data_files[0];
516        assert_eq!(data_file.content_type(), DataContentType::PositionDeletes);
517        assert_eq!(data_file.partition_spec_id, 7);
518        assert_eq!(data_file.partition, partition_value);
519
520        Ok(())
521    }
522
523    #[tokio::test]
524    async fn test_position_delete_writer_rejects_wrong_column_count() -> Result<()> {
525        let temp_dir = TempDir::new().unwrap();
526        let (_file_io, builder) = writer_setup(&temp_dir);
527        let mut writer = builder.build(None).await?;
528
529        // A batch carrying only the `file_path` column is not a position delete file.
530        let arrow_schema = position_delete_arrow_schema();
531        let path_only = RecordBatch::try_new(Arc::new(arrow_schema.project(&[0]).unwrap()), vec![
532            Arc::new(StringArray::from(vec!["s3://bucket/data/f0.parquet"])),
533        ])
534        .unwrap();
535
536        let err = writer.write(path_only).await.unwrap_err();
537        assert_eq!(err.kind(), ErrorKind::DataInvalid);
538        assert!(
539            err.to_string()
540                .contains("only the two required position delete columns"),
541            "{err}"
542        );
543        Ok(())
544    }
545
546    #[tokio::test]
547    async fn test_position_delete_writer_rejects_missing_field_ids() -> Result<()> {
548        let temp_dir = TempDir::new().unwrap();
549        let (_file_io, builder) = writer_setup(&temp_dir);
550        let mut writer = builder.build(None).await?;
551
552        // Correct shape and types, but plain field names without the reserved
553        // Iceberg field-id metadata: must be rejected.
554        let plain_schema = Arc::new(arrow_schema::Schema::new(vec![
555            Field::new(RESERVED_COL_NAME_DELETE_FILE_PATH, DataType::Utf8, false),
556            Field::new(RESERVED_COL_NAME_DELETE_FILE_POS, DataType::Int64, false),
557        ]));
558        let batch = RecordBatch::try_new(plain_schema, vec![
559            Arc::new(StringArray::from(vec!["s3://bucket/data/f0.parquet"])),
560            Arc::new(Int64Array::from(vec![1_i64])),
561        ])
562        .unwrap();
563
564        let err = writer.write(batch).await.unwrap_err();
565        assert_eq!(err.kind(), ErrorKind::DataInvalid);
566        assert!(err.to_string().contains("field id metadata"), "{err}");
567        Ok(())
568    }
569
570    #[tokio::test]
571    async fn test_position_delete_writer_rejects_unparsable_field_id() -> Result<()> {
572        let temp_dir = TempDir::new().unwrap();
573        let (_file_io, builder) = writer_setup(&temp_dir);
574        let mut writer = builder.build(None).await?;
575
576        // Field-id metadata present but not an integer: must be rejected.
577        let bad_meta = Arc::new(arrow_schema::Schema::new(vec![
578            Field::new(RESERVED_COL_NAME_DELETE_FILE_PATH, DataType::Utf8, false).with_metadata(
579                HashMap::from([(
580                    PARQUET_FIELD_ID_META_KEY.to_string(),
581                    "not_a_number".to_string(),
582                )]),
583            ),
584            field_with_id(
585                RESERVED_COL_NAME_DELETE_FILE_POS,
586                DataType::Int64,
587                RESERVED_FIELD_ID_DELETE_FILE_POS,
588            ),
589        ]));
590        let batch = RecordBatch::try_new(bad_meta, vec![
591            Arc::new(StringArray::from(vec!["s3://bucket/data/f0.parquet"])),
592            Arc::new(Int64Array::from(vec![1_i64])),
593        ])
594        .unwrap();
595
596        let err = writer.write(batch).await.unwrap_err();
597        assert_eq!(err.kind(), ErrorKind::DataInvalid);
598        assert!(err.to_string().contains("invalid field id"), "{err}");
599        Ok(())
600    }
601
602    #[tokio::test]
603    async fn test_position_delete_writer_rejects_wrong_field_id() -> Result<()> {
604        let temp_dir = TempDir::new().unwrap();
605        let (_file_io, builder) = writer_setup(&temp_dir);
606        let mut writer = builder.build(None).await?;
607
608        // Right shape and types, but the two reserved field ids are swapped.
609        let swapped = Arc::new(arrow_schema::Schema::new(vec![
610            field_with_id(
611                RESERVED_COL_NAME_DELETE_FILE_PATH,
612                DataType::Utf8,
613                RESERVED_FIELD_ID_DELETE_FILE_POS,
614            ),
615            field_with_id(
616                RESERVED_COL_NAME_DELETE_FILE_POS,
617                DataType::Int64,
618                RESERVED_FIELD_ID_DELETE_FILE_PATH,
619            ),
620        ]));
621        let batch = RecordBatch::try_new(swapped, vec![
622            Arc::new(StringArray::from(vec!["s3://bucket/data/f0.parquet"])),
623            Arc::new(Int64Array::from(vec![1_i64])),
624        ])
625        .unwrap();
626
627        let err = writer.write(batch).await.unwrap_err();
628        assert_eq!(err.kind(), ErrorKind::DataInvalid);
629        assert!(err.to_string().contains("must be `file_path`"), "{err}");
630        Ok(())
631    }
632
633    #[tokio::test]
634    async fn test_position_delete_writer_rejects_bad_pos_field_id() -> Result<()> {
635        let temp_dir = TempDir::new().unwrap();
636        let (_file_io, builder) = writer_setup(&temp_dir);
637        let mut writer = builder.build(None).await?;
638
639        // Correct `file_path`, but `pos` carries the wrong reserved field id — the
640        // pos-specific branch (not the path branch) must fire.
641        let wrong_pos_id = Arc::new(arrow_schema::Schema::new(vec![
642            field_with_id(
643                RESERVED_COL_NAME_DELETE_FILE_PATH,
644                DataType::Utf8,
645                RESERVED_FIELD_ID_DELETE_FILE_PATH,
646            ),
647            field_with_id(RESERVED_COL_NAME_DELETE_FILE_POS, DataType::Int64, 999),
648        ]));
649        let batch = RecordBatch::try_new(wrong_pos_id, vec![
650            Arc::new(StringArray::from(vec!["s3://bucket/data/f0.parquet"])),
651            Arc::new(Int64Array::from(vec![1_i64])),
652        ])
653        .unwrap();
654        let err = writer.write(batch).await.unwrap_err();
655        assert_eq!(err.kind(), ErrorKind::DataInvalid);
656        assert!(err.to_string().contains("must be `pos`"), "{err}");
657
658        // Correct `file_path`, but `pos` is missing its field-id metadata entirely.
659        let missing_pos_id = Arc::new(arrow_schema::Schema::new(vec![
660            field_with_id(
661                RESERVED_COL_NAME_DELETE_FILE_PATH,
662                DataType::Utf8,
663                RESERVED_FIELD_ID_DELETE_FILE_PATH,
664            ),
665            Field::new(RESERVED_COL_NAME_DELETE_FILE_POS, DataType::Int64, false),
666        ]));
667        let batch = RecordBatch::try_new(missing_pos_id, vec![
668            Arc::new(StringArray::from(vec!["s3://bucket/data/f0.parquet"])),
669            Arc::new(Int64Array::from(vec![1_i64])),
670        ])
671        .unwrap();
672        let err = writer.write(batch).await.unwrap_err();
673        assert_eq!(err.kind(), ErrorKind::DataInvalid);
674        assert!(err.to_string().contains("field id metadata"), "{err}");
675
676        Ok(())
677    }
678
679    #[tokio::test]
680    async fn test_position_delete_writer_rejects_wrong_column_types() -> Result<()> {
681        let temp_dir = TempDir::new().unwrap();
682        let (_file_io, builder) = writer_setup(&temp_dir);
683        let mut writer = builder.build(None).await?;
684
685        // Correct field ids, but `pos` is Int32 rather than Int64.
686        let int32_pos = Arc::new(arrow_schema::Schema::new(vec![
687            field_with_id(
688                RESERVED_COL_NAME_DELETE_FILE_PATH,
689                DataType::Utf8,
690                RESERVED_FIELD_ID_DELETE_FILE_PATH,
691            ),
692            field_with_id(
693                RESERVED_COL_NAME_DELETE_FILE_POS,
694                DataType::Int32,
695                RESERVED_FIELD_ID_DELETE_FILE_POS,
696            ),
697        ]));
698        let batch = RecordBatch::try_new(int32_pos, vec![
699            Arc::new(StringArray::from(vec!["s3://bucket/data/f0.parquet"])),
700            Arc::new(Int32Array::from(vec![1_i32])),
701        ])
702        .unwrap();
703        let err = writer.write(batch).await.unwrap_err();
704        assert_eq!(err.kind(), ErrorKind::DataInvalid);
705        assert!(err.to_string().contains("must be Int64"), "{err}");
706
707        // Correct field ids, but `file_path` is not a string.
708        let int_path = Arc::new(arrow_schema::Schema::new(vec![
709            field_with_id(
710                RESERVED_COL_NAME_DELETE_FILE_PATH,
711                DataType::Int32,
712                RESERVED_FIELD_ID_DELETE_FILE_PATH,
713            ),
714            field_with_id(
715                RESERVED_COL_NAME_DELETE_FILE_POS,
716                DataType::Int64,
717                RESERVED_FIELD_ID_DELETE_FILE_POS,
718            ),
719        ]));
720        let batch = RecordBatch::try_new(int_path, vec![
721            Arc::new(Int32Array::from(vec![1_i32])),
722            Arc::new(Int64Array::from(vec![1_i64])),
723        ])
724        .unwrap();
725        let err = writer.write(batch).await.unwrap_err();
726        assert_eq!(err.kind(), ErrorKind::DataInvalid);
727        assert!(err.to_string().contains("must be Utf8"), "{err}");
728
729        // Correct field ids, but `file_path` is LargeUtf8 — the common shape from
730        // DataFusion/DuckDB/Polars, and the case the validator's comment calls out.
731        let large_path = Arc::new(arrow_schema::Schema::new(vec![
732            field_with_id(
733                RESERVED_COL_NAME_DELETE_FILE_PATH,
734                DataType::LargeUtf8,
735                RESERVED_FIELD_ID_DELETE_FILE_PATH,
736            ),
737            field_with_id(
738                RESERVED_COL_NAME_DELETE_FILE_POS,
739                DataType::Int64,
740                RESERVED_FIELD_ID_DELETE_FILE_POS,
741            ),
742        ]));
743        let batch = RecordBatch::try_new(large_path, vec![
744            Arc::new(LargeStringArray::from(vec!["s3://bucket/data/f0.parquet"])),
745            Arc::new(Int64Array::from(vec![1_i64])),
746        ])
747        .unwrap();
748        let err = writer.write(batch).await.unwrap_err();
749        assert_eq!(err.kind(), ErrorKind::DataInvalid);
750        assert!(err.to_string().contains("must be Utf8"), "{err}");
751
752        Ok(())
753    }
754
755    #[tokio::test]
756    async fn test_position_delete_writer_rejects_nullable_columns() -> Result<()> {
757        let temp_dir = TempDir::new().unwrap();
758        let (_file_io, builder) = writer_setup(&temp_dir);
759        let mut writer = builder.build(None).await?;
760
761        // Correct field ids and types, but the columns are declared nullable. A
762        // required schema with a nullable field could emit nulls -> malformed file.
763        let nullable = |name: &str, data_type: DataType, id: i32| {
764            Field::new(name, data_type, true).with_metadata(HashMap::from([(
765                PARQUET_FIELD_ID_META_KEY.to_string(),
766                id.to_string(),
767            )]))
768        };
769
770        let nullable_path = Arc::new(arrow_schema::Schema::new(vec![
771            nullable(
772                RESERVED_COL_NAME_DELETE_FILE_PATH,
773                DataType::Utf8,
774                RESERVED_FIELD_ID_DELETE_FILE_PATH,
775            ),
776            field_with_id(
777                RESERVED_COL_NAME_DELETE_FILE_POS,
778                DataType::Int64,
779                RESERVED_FIELD_ID_DELETE_FILE_POS,
780            ),
781        ]));
782        let batch = RecordBatch::try_new(nullable_path, vec![
783            Arc::new(StringArray::from(vec!["s3://bucket/data/f0.parquet"])),
784            Arc::new(Int64Array::from(vec![1_i64])),
785        ])
786        .unwrap();
787        let err = writer.write(batch).await.unwrap_err();
788        assert_eq!(err.kind(), ErrorKind::DataInvalid);
789        assert!(
790            err.to_string()
791                .contains("`file_path` column must be required"),
792            "{err}"
793        );
794
795        let nullable_pos = Arc::new(arrow_schema::Schema::new(vec![
796            field_with_id(
797                RESERVED_COL_NAME_DELETE_FILE_PATH,
798                DataType::Utf8,
799                RESERVED_FIELD_ID_DELETE_FILE_PATH,
800            ),
801            nullable(
802                RESERVED_COL_NAME_DELETE_FILE_POS,
803                DataType::Int64,
804                RESERVED_FIELD_ID_DELETE_FILE_POS,
805            ),
806        ]));
807        let batch = RecordBatch::try_new(nullable_pos, vec![
808            Arc::new(StringArray::from(vec!["s3://bucket/data/f0.parquet"])),
809            Arc::new(Int64Array::from(vec![1_i64])),
810        ])
811        .unwrap();
812        let err = writer.write(batch).await.unwrap_err();
813        assert_eq!(err.kind(), ErrorKind::DataInvalid);
814        assert!(
815            err.to_string().contains("`pos` column must be required"),
816            "{err}"
817        );
818
819        Ok(())
820    }
821
822    #[tokio::test]
823    async fn test_position_delete_writer_close_without_writes() -> Result<()> {
824        let temp_dir = TempDir::new().unwrap();
825        let (_file_io, builder) = writer_setup(&temp_dir);
826        let mut writer = builder.build(None).await?;
827
828        // Closing a writer that never received a batch produces no data files.
829        let data_files = writer.close().await?;
830        assert!(data_files.is_empty());
831        Ok(())
832    }
833
834    #[tokio::test]
835    async fn test_position_delete_writer_errors_after_close() -> Result<()> {
836        let temp_dir = TempDir::new().unwrap();
837        let (_file_io, builder) = writer_setup(&temp_dir);
838        let mut writer = builder.build(None).await?;
839        writer.close().await?;
840
841        // Both write() and a second close() report the writer is already closed.
842        let write_err = writer
843            .write(position_delete_batch(
844                vec!["s3://bucket/data/f0.parquet"],
845                vec![1],
846            ))
847            .await
848            .unwrap_err();
849        assert_eq!(write_err.kind(), ErrorKind::Unexpected);
850        assert!(
851            write_err.to_string().contains("cannot write"),
852            "{write_err}"
853        );
854
855        // Even a malformed batch surfaces the closed error, not a validation error:
856        // the closed check runs before validation.
857        let arrow_schema = position_delete_arrow_schema();
858        let invalid = RecordBatch::try_new(Arc::new(arrow_schema.project(&[0]).unwrap()), vec![
859            Arc::new(StringArray::from(vec!["s3://bucket/data/f0.parquet"])),
860        ])
861        .unwrap();
862        let invalid_err = writer.write(invalid).await.unwrap_err();
863        assert_eq!(invalid_err.kind(), ErrorKind::Unexpected);
864
865        let close_err = writer.close().await.unwrap_err();
866        assert_eq!(close_err.kind(), ErrorKind::Unexpected);
867        assert!(
868            close_err.to_string().contains("already closed"),
869            "{close_err}"
870        );
871        Ok(())
872    }
873
874    async fn read_back_single(
875        file_io: &FileIO,
876        data_file: &DataFile,
877        schema: &arrow_schema::SchemaRef,
878    ) -> RecordBatch {
879        let input_content = file_io
880            .new_input(data_file.file_path.clone())
881            .unwrap()
882            .read()
883            .await
884            .unwrap();
885        let reader = ParquetRecordBatchReaderBuilder::try_new(input_content)
886            .unwrap()
887            .build()
888            .unwrap();
889        let batches = reader.map(|b| b.unwrap()).collect::<Vec<_>>();
890        concat_batches(schema, &batches).unwrap()
891    }
892}