1use 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
53static 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
67static 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
78pub fn position_delete_schema() -> SchemaRef {
84 POSITION_DELETE_SCHEMA.clone()
85}
86
87pub(crate) fn position_delete_arrow_schema() -> arrow_schema::SchemaRef {
92 POSITION_DELETE_ARROW_SCHEMA.clone()
93}
94
95fn 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
115fn 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 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 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#[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 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#[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 async fn write(&mut self, batch: RecordBatch) -> Result<()> {
241 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 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 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 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 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 assert_eq!(data_file.partition, Struct::empty());
410 assert_eq!(data_file.partition_spec_id, 0);
411 assert!(data_file.file_size_in_bytes > 0);
413
414 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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}