Copilot commented on code in PR #3291:
URL: https://github.com/apache/iceberg-rust/pull/3291#discussion_r4126066631
##########
crates/iceberg/src/arrow/delete_file_loader.rs:
##########
@@ -141,14 +144,785 @@ impl DeleteFileLoader for BasicDeleteFileLoader {
}
}
+/// In-memory positions from a file-scoped position-delete file.
+///
+/// Iteration is sorted and de-duplicated.
+pub struct PositionDeleteIndex {
+ positions: DeleteVector,
+}
+
+impl PositionDeleteIndex {
+ fn new() -> Self {
+ Self {
+ positions: DeleteVector::default(),
+ }
+ }
+
+ /// Iterates deleted row positions in ascending order.
+ pub fn iter(&self) -> impl Iterator<Item = i64> + '_ {
+ self.positions.iter().map(|position| position as i64)
+ }
+}
+
+/// Loads a file-scoped V2 Parquet position-delete file into an index.
+///
+/// Validates `file_path`, `pos`, and the expected data-file target. Snapshot
applicability and
+/// replacement policy remain the caller's responsibility.
+pub struct PositionDeleteIndexLoader {
+ basic_loader: BasicDeleteFileLoader,
+}
+
+impl PositionDeleteIndexLoader {
+ /// Creates a loader for the given Iceberg `FileIO`.
+ pub fn new(file_io: FileIO) -> Self {
+ Self {
+ basic_loader: BasicDeleteFileLoader::new(file_io,
ScanMetrics::new()),
Review Comment:
`PositionDeleteIndexLoader` is tightly coupled to `BasicDeleteFileLoader`
just to reuse `file_io()` and a bytes counter, and it unconditionally creates a
fresh `ScanMetrics`. This makes it hard to integrate loader IO/metrics into a
caller’s existing scan metrics and increases coupling between otherwise
separate responsibilities; consider storing `FileIO` (and optionally a
`ScanMetrics`/counter reference passed from the caller) directly in
`PositionDeleteIndexLoader` instead of embedding `BasicDeleteFileLoader`.
##########
crates/iceberg/src/arrow/delete_file_loader.rs:
##########
@@ -141,14 +144,785 @@ impl DeleteFileLoader for BasicDeleteFileLoader {
}
}
+/// In-memory positions from a file-scoped position-delete file.
+///
+/// Iteration is sorted and de-duplicated.
+pub struct PositionDeleteIndex {
+ positions: DeleteVector,
+}
+
+impl PositionDeleteIndex {
+ fn new() -> Self {
+ Self {
+ positions: DeleteVector::default(),
+ }
+ }
+
+ /// Iterates deleted row positions in ascending order.
+ pub fn iter(&self) -> impl Iterator<Item = i64> + '_ {
+ self.positions.iter().map(|position| position as i64)
+ }
+}
+
+/// Loads a file-scoped V2 Parquet position-delete file into an index.
+///
+/// Validates `file_path`, `pos`, and the expected data-file target. Snapshot
applicability and
+/// replacement policy remain the caller's responsibility.
+pub struct PositionDeleteIndexLoader {
+ basic_loader: BasicDeleteFileLoader,
+}
+
+impl PositionDeleteIndexLoader {
+ /// Creates a loader for the given Iceberg `FileIO`.
+ pub fn new(file_io: FileIO) -> Self {
+ Self {
+ basic_loader: BasicDeleteFileLoader::new(file_io,
ScanMetrics::new()),
+ }
+ }
+
+ fn reserved_field_index(
+ schema: &arrow_schema::Schema,
+ field_id: i32,
+ logical_name: &str,
+ delete_file_path: &str,
+ ) -> Result<usize> {
+ let mut matches = schema
+ .fields()
+ .iter()
+ .enumerate()
+ .filter_map(|(index, field)| {
+ field
+ .metadata()
+ .get(PARQUET_FIELD_ID_META_KEY)
+ .and_then(|id| id.parse::<i32>().ok())
+ .filter(|id| *id == field_id)
+ .map(|_| index)
+ });
+
+ let Some(index) = matches.next() else {
+ return Err(invalid_data!(
+ "Position-delete file {delete_file_path} has no
`{logical_name}` column with reserved field id {field_id}"
+ ));
+ };
+ if matches.next().is_some() {
+ return Err(invalid_data!(
+ "Position-delete file {delete_file_path} has multiple columns
with reserved field id {field_id}"
+ ));
+ }
+
+ Ok(index)
+ }
+
+ fn position_delete_columns(
+ schema: &arrow_schema::Schema,
+ delete_file_path: &str,
+ ) -> Result<(usize, usize)> {
+ let path_index = Self::reserved_field_index(
+ schema,
+ crate::metadata_columns::RESERVED_FIELD_ID_DELETE_FILE_PATH,
+ "file_path",
+ delete_file_path,
+ )?;
+ let position_index = Self::reserved_field_index(
+ schema,
+ crate::metadata_columns::RESERVED_FIELD_ID_DELETE_FILE_POS,
+ "pos",
+ delete_file_path,
+ )?;
+
+ let path_field = schema.field(path_index);
+ if path_field.data_type() != &arrow_schema::DataType::Utf8 ||
path_field.is_nullable() {
+ return Err(invalid_data!(
+ "Position-delete file {delete_file_path} requires non-nullable
Utf8 file_path"
+ ));
+ }
+
+ let position_field = schema.field(position_index);
+ if position_field.data_type() != &arrow_schema::DataType::Int64
+ || position_field.is_nullable()
+ {
+ return Err(invalid_data!(
+ "Position-delete file {delete_file_path} requires non-nullable
Int64 pos"
+ ));
+ }
+
+ Ok((path_index, position_index))
+ }
+
+ /// Loads positions targeting `expected_data_file`.
+ ///
+ /// Positions are sorted and de-duplicated after validating the physical
row count.
+ pub async fn load_file_scoped_positions(
+ &self,
+ delete_file: &FileScanTaskDeleteFile,
+ expected_data_file: &str,
+ ) -> Result<PositionDeleteIndex> {
+ if delete_file.file_type != DataContentType::PositionDeletes
+ || delete_file.file_format != DataFileFormat::Parquet
+ {
+ return Err(Error::new(
+ ErrorKind::FeatureUnsupported,
+ format!(
+ "Expected a V2 Parquet position-delete file, got {:?}/{:?}
at {}",
+ delete_file.file_type, delete_file.file_format,
delete_file.file_path
+ ),
+ ));
+ }
+
+ if delete_file.equality_ids.is_some() {
+ return Err(invalid_data!(
+ "Position-delete file {} must not carry equality_ids",
+ delete_file.file_path
+ ));
+ }
+
+ if delete_file.content_offset.is_some() ||
delete_file.content_size_in_bytes.is_some() {
+ return Err(invalid_data!(
+ "V2 Parquet position-delete file {} must not carry
deletion-vector content coordinates",
+ delete_file.file_path
+ ));
+ }
+
+ if let Some(referenced_data_file) =
delete_file.referenced_data_file.as_deref()
+ && referenced_data_file != expected_data_file
+ {
+ return Err(invalid_data!(
+ "Position-delete file {} references {referenced_data_file},
expected {expected_data_file}",
+ delete_file.file_path
+ ));
+ }
+
+ let parquet_read_options = ParquetReadOptions::builder().build();
+ let (parquet_file_reader, arrow_metadata) =
ArrowReader::open_parquet_file(
+ &delete_file.file_path,
+ self.basic_loader.file_io(),
+ delete_file.file_size_in_bytes,
+ parquet_read_options,
+ self.basic_loader.scan_metrics.bytes_read_counter(),
+ delete_file.key_metadata.as_deref(),
+ )
+ .await?;
+
+ let mut stream_builder =
+
ParquetRecordBatchStreamBuilder::new_with_metadata(parquet_file_reader,
arrow_metadata);
+
+ // Validate before reading so malformed empty files are rejected.
+ let (path_root_index, position_root_index) =
Self::position_delete_columns(
+ stream_builder.schema().as_ref(),
+ &delete_file.file_path,
+ )?;
+
+ // Read only file_path and pos; skip optional deleted-row payloads.
+ let projection =
ProjectionMask::roots(stream_builder.parquet_schema(), vec![
+ path_root_index,
+ position_root_index,
+ ]);
+ stream_builder = stream_builder.with_projection(projection);
+
+ // Projection compacts ordinals, so resolve columns from the projected
schema.
+ let projected_stream = stream_builder.build()?;
+ let (path_index, position_index) = Self::position_delete_columns(
+ projected_stream.schema().as_ref(),
+ &delete_file.file_path,
+ )?;
Review Comment:
This re-resolves columns after projection by looking for
`PARQUET_FIELD_ID_META_KEY` metadata in the *projected* Arrow schema. That
assumes the projection path preserves field-id metadata on `Field` objects; if
a future parquet/arrow change drops or rewrites metadata during projection,
this will start failing on valid files. A more robust approach is to derive the
projected column ordinals deterministically from the projection you applied (or
otherwise track the mapping once) so you don’t depend on metadata being
preserved through projection.
##########
crates/iceberg/src/arrow/delete_file_loader.rs:
##########
@@ -141,14 +144,785 @@ impl DeleteFileLoader for BasicDeleteFileLoader {
}
}
+/// In-memory positions from a file-scoped position-delete file.
+///
+/// Iteration is sorted and de-duplicated.
+pub struct PositionDeleteIndex {
+ positions: DeleteVector,
+}
+
+impl PositionDeleteIndex {
+ fn new() -> Self {
+ Self {
+ positions: DeleteVector::default(),
+ }
+ }
+
+ /// Iterates deleted row positions in ascending order.
+ pub fn iter(&self) -> impl Iterator<Item = i64> + '_ {
+ self.positions.iter().map(|position| position as i64)
+ }
+}
+
+/// Loads a file-scoped V2 Parquet position-delete file into an index.
+///
+/// Validates `file_path`, `pos`, and the expected data-file target. Snapshot
applicability and
+/// replacement policy remain the caller's responsibility.
+pub struct PositionDeleteIndexLoader {
+ basic_loader: BasicDeleteFileLoader,
+}
+
+impl PositionDeleteIndexLoader {
+ /// Creates a loader for the given Iceberg `FileIO`.
+ pub fn new(file_io: FileIO) -> Self {
+ Self {
+ basic_loader: BasicDeleteFileLoader::new(file_io,
ScanMetrics::new()),
+ }
+ }
+
+ fn reserved_field_index(
+ schema: &arrow_schema::Schema,
+ field_id: i32,
+ logical_name: &str,
+ delete_file_path: &str,
+ ) -> Result<usize> {
+ let mut matches = schema
+ .fields()
+ .iter()
+ .enumerate()
+ .filter_map(|(index, field)| {
+ field
+ .metadata()
+ .get(PARQUET_FIELD_ID_META_KEY)
+ .and_then(|id| id.parse::<i32>().ok())
+ .filter(|id| *id == field_id)
+ .map(|_| index)
+ });
+
+ let Some(index) = matches.next() else {
+ return Err(invalid_data!(
+ "Position-delete file {delete_file_path} has no
`{logical_name}` column with reserved field id {field_id}"
+ ));
+ };
+ if matches.next().is_some() {
+ return Err(invalid_data!(
+ "Position-delete file {delete_file_path} has multiple columns
with reserved field id {field_id}"
+ ));
+ }
+
+ Ok(index)
+ }
+
+ fn position_delete_columns(
+ schema: &arrow_schema::Schema,
+ delete_file_path: &str,
+ ) -> Result<(usize, usize)> {
+ let path_index = Self::reserved_field_index(
+ schema,
+ crate::metadata_columns::RESERVED_FIELD_ID_DELETE_FILE_PATH,
+ "file_path",
+ delete_file_path,
+ )?;
+ let position_index = Self::reserved_field_index(
+ schema,
+ crate::metadata_columns::RESERVED_FIELD_ID_DELETE_FILE_POS,
+ "pos",
+ delete_file_path,
+ )?;
+
+ let path_field = schema.field(path_index);
+ if path_field.data_type() != &arrow_schema::DataType::Utf8 ||
path_field.is_nullable() {
+ return Err(invalid_data!(
+ "Position-delete file {delete_file_path} requires non-nullable
Utf8 file_path"
+ ));
+ }
+
+ let position_field = schema.field(position_index);
+ if position_field.data_type() != &arrow_schema::DataType::Int64
+ || position_field.is_nullable()
+ {
+ return Err(invalid_data!(
+ "Position-delete file {delete_file_path} requires non-nullable
Int64 pos"
+ ));
+ }
+
+ Ok((path_index, position_index))
+ }
+
+ /// Loads positions targeting `expected_data_file`.
+ ///
+ /// Positions are sorted and de-duplicated after validating the physical
row count.
+ pub async fn load_file_scoped_positions(
+ &self,
+ delete_file: &FileScanTaskDeleteFile,
+ expected_data_file: &str,
+ ) -> Result<PositionDeleteIndex> {
+ if delete_file.file_type != DataContentType::PositionDeletes
+ || delete_file.file_format != DataFileFormat::Parquet
+ {
+ return Err(Error::new(
+ ErrorKind::FeatureUnsupported,
+ format!(
+ "Expected a V2 Parquet position-delete file, got {:?}/{:?}
at {}",
+ delete_file.file_type, delete_file.file_format,
delete_file.file_path
+ ),
+ ));
+ }
+
+ if delete_file.equality_ids.is_some() {
+ return Err(invalid_data!(
+ "Position-delete file {} must not carry equality_ids",
+ delete_file.file_path
+ ));
+ }
+
+ if delete_file.content_offset.is_some() ||
delete_file.content_size_in_bytes.is_some() {
+ return Err(invalid_data!(
+ "V2 Parquet position-delete file {} must not carry
deletion-vector content coordinates",
+ delete_file.file_path
+ ));
+ }
+
+ if let Some(referenced_data_file) =
delete_file.referenced_data_file.as_deref()
+ && referenced_data_file != expected_data_file
+ {
+ return Err(invalid_data!(
+ "Position-delete file {} references {referenced_data_file},
expected {expected_data_file}",
+ delete_file.file_path
+ ));
+ }
+
+ let parquet_read_options = ParquetReadOptions::builder().build();
+ let (parquet_file_reader, arrow_metadata) =
ArrowReader::open_parquet_file(
+ &delete_file.file_path,
+ self.basic_loader.file_io(),
+ delete_file.file_size_in_bytes,
+ parquet_read_options,
+ self.basic_loader.scan_metrics.bytes_read_counter(),
+ delete_file.key_metadata.as_deref(),
+ )
+ .await?;
Review Comment:
The PR description says this “covers … encrypted files”, but the added tests
don’t appear to exercise `PositionDeleteIndexLoader` behavior with an encrypted
Parquet delete file (e.g., via `write_encrypted_parquet` + `key_metadata`).
Either add a targeted test for encrypted position-delete files (success or
expected failure, depending on intended support), or adjust the description to
match what’s currently covered.
##########
crates/iceberg/src/arrow/delete_file_loader.rs:
##########
@@ -141,14 +144,785 @@ impl DeleteFileLoader for BasicDeleteFileLoader {
}
}
+/// In-memory positions from a file-scoped position-delete file.
+///
+/// Iteration is sorted and de-duplicated.
+pub struct PositionDeleteIndex {
+ positions: DeleteVector,
+}
+
+impl PositionDeleteIndex {
+ fn new() -> Self {
+ Self {
+ positions: DeleteVector::default(),
+ }
+ }
+
+ /// Iterates deleted row positions in ascending order.
+ pub fn iter(&self) -> impl Iterator<Item = i64> + '_ {
+ self.positions.iter().map(|position| position as i64)
+ }
+}
+
+/// Loads a file-scoped V2 Parquet position-delete file into an index.
+///
+/// Validates `file_path`, `pos`, and the expected data-file target. Snapshot
applicability and
+/// replacement policy remain the caller's responsibility.
+pub struct PositionDeleteIndexLoader {
+ basic_loader: BasicDeleteFileLoader,
+}
+
+impl PositionDeleteIndexLoader {
+ /// Creates a loader for the given Iceberg `FileIO`.
+ pub fn new(file_io: FileIO) -> Self {
+ Self {
+ basic_loader: BasicDeleteFileLoader::new(file_io,
ScanMetrics::new()),
+ }
+ }
+
+ fn reserved_field_index(
+ schema: &arrow_schema::Schema,
+ field_id: i32,
+ logical_name: &str,
+ delete_file_path: &str,
+ ) -> Result<usize> {
+ let mut matches = schema
+ .fields()
+ .iter()
+ .enumerate()
+ .filter_map(|(index, field)| {
+ field
+ .metadata()
+ .get(PARQUET_FIELD_ID_META_KEY)
+ .and_then(|id| id.parse::<i32>().ok())
+ .filter(|id| *id == field_id)
+ .map(|_| index)
+ });
+
+ let Some(index) = matches.next() else {
+ return Err(invalid_data!(
+ "Position-delete file {delete_file_path} has no
`{logical_name}` column with reserved field id {field_id}"
+ ));
+ };
+ if matches.next().is_some() {
+ return Err(invalid_data!(
+ "Position-delete file {delete_file_path} has multiple columns
with reserved field id {field_id}"
+ ));
+ }
+
+ Ok(index)
+ }
+
+ fn position_delete_columns(
+ schema: &arrow_schema::Schema,
+ delete_file_path: &str,
+ ) -> Result<(usize, usize)> {
+ let path_index = Self::reserved_field_index(
+ schema,
+ crate::metadata_columns::RESERVED_FIELD_ID_DELETE_FILE_PATH,
+ "file_path",
+ delete_file_path,
+ )?;
+ let position_index = Self::reserved_field_index(
+ schema,
+ crate::metadata_columns::RESERVED_FIELD_ID_DELETE_FILE_POS,
+ "pos",
+ delete_file_path,
+ )?;
+
+ let path_field = schema.field(path_index);
+ if path_field.data_type() != &arrow_schema::DataType::Utf8 ||
path_field.is_nullable() {
+ return Err(invalid_data!(
+ "Position-delete file {delete_file_path} requires non-nullable
Utf8 file_path"
+ ));
+ }
+
+ let position_field = schema.field(position_index);
+ if position_field.data_type() != &arrow_schema::DataType::Int64
+ || position_field.is_nullable()
+ {
+ return Err(invalid_data!(
+ "Position-delete file {delete_file_path} requires non-nullable
Int64 pos"
+ ));
+ }
+
+ Ok((path_index, position_index))
+ }
+
+ /// Loads positions targeting `expected_data_file`.
+ ///
+ /// Positions are sorted and de-duplicated after validating the physical
row count.
+ pub async fn load_file_scoped_positions(
+ &self,
+ delete_file: &FileScanTaskDeleteFile,
+ expected_data_file: &str,
+ ) -> Result<PositionDeleteIndex> {
+ if delete_file.file_type != DataContentType::PositionDeletes
+ || delete_file.file_format != DataFileFormat::Parquet
+ {
+ return Err(Error::new(
+ ErrorKind::FeatureUnsupported,
+ format!(
+ "Expected a V2 Parquet position-delete file, got {:?}/{:?}
at {}",
+ delete_file.file_type, delete_file.file_format,
delete_file.file_path
+ ),
+ ));
+ }
+
+ if delete_file.equality_ids.is_some() {
+ return Err(invalid_data!(
+ "Position-delete file {} must not carry equality_ids",
+ delete_file.file_path
+ ));
+ }
+
+ if delete_file.content_offset.is_some() ||
delete_file.content_size_in_bytes.is_some() {
+ return Err(invalid_data!(
+ "V2 Parquet position-delete file {} must not carry
deletion-vector content coordinates",
+ delete_file.file_path
+ ));
+ }
+
+ if let Some(referenced_data_file) =
delete_file.referenced_data_file.as_deref()
+ && referenced_data_file != expected_data_file
+ {
+ return Err(invalid_data!(
+ "Position-delete file {} references {referenced_data_file},
expected {expected_data_file}",
+ delete_file.file_path
+ ));
+ }
+
+ let parquet_read_options = ParquetReadOptions::builder().build();
+ let (parquet_file_reader, arrow_metadata) =
ArrowReader::open_parquet_file(
+ &delete_file.file_path,
+ self.basic_loader.file_io(),
+ delete_file.file_size_in_bytes,
+ parquet_read_options,
+ self.basic_loader.scan_metrics.bytes_read_counter(),
+ delete_file.key_metadata.as_deref(),
+ )
+ .await?;
+
+ let mut stream_builder =
+
ParquetRecordBatchStreamBuilder::new_with_metadata(parquet_file_reader,
arrow_metadata);
+
+ // Validate before reading so malformed empty files are rejected.
+ let (path_root_index, position_root_index) =
Self::position_delete_columns(
+ stream_builder.schema().as_ref(),
+ &delete_file.file_path,
+ )?;
+
+ // Read only file_path and pos; skip optional deleted-row payloads.
+ let projection =
ProjectionMask::roots(stream_builder.parquet_schema(), vec![
+ path_root_index,
+ position_root_index,
+ ]);
+ stream_builder = stream_builder.with_projection(projection);
+
+ // Projection compacts ordinals, so resolve columns from the projected
schema.
+ let projected_stream = stream_builder.build()?;
+ let (path_index, position_index) = Self::position_delete_columns(
+ projected_stream.schema().as_ref(),
+ &delete_file.file_path,
+ )?;
+ let mut batches = projected_stream.map_err(|e| {
+ Error::new(
+ ErrorKind::Unexpected,
+ format!(
+ "Failed to read position-delete file {}",
+ delete_file.file_path
+ ),
+ )
+ .with_source(e)
+ });
+
+ let mut index = PositionDeleteIndex::new();
+ let mut rows_read = 0u64;
+
+ while let Some(batch) = batches.try_next().await? {
+ let paths = batch
+ .column(path_index)
+ .as_any()
+ .downcast_ref::<StringArray>()
+ .ok_or_else(|| {
+ invalid_data!(
+ "Position-delete file {} has a non-Utf8 file_path
column",
+ delete_file.file_path
+ )
+ })?;
+ let row_positions = batch
+ .column(position_index)
+ .as_any()
+ .downcast_ref::<Int64Array>()
+ .ok_or_else(|| {
+ invalid_data!(
+ "Position-delete file {} has a non-Int64 pos column",
+ delete_file.file_path
+ )
+ })?;
+ if paths.null_count() != 0 || row_positions.null_count() != 0 {
+ return Err(invalid_data!(
+ "Position-delete file {} contains nulls",
+ delete_file.file_path
+ ));
+ }
+
+ rows_read += batch.num_rows() as u64;
+ let mut batch_positions = Vec::with_capacity(batch.num_rows());
+ for row in 0..batch.num_rows() {
+ let data_file = paths.value(row);
+ if data_file != expected_data_file {
+ return Err(invalid_data!(
+ "File-scoped position-delete file {} contains target
{data_file}, expected {expected_data_file}",
+ delete_file.file_path
+ ));
+ }
+
+ let position = row_positions.value(row);
+ if position < 0 {
+ return Err(invalid_data!(
+ "Position-delete file {} contains a negative row
position {position}",
+ delete_file.file_path
+ ));
+ }
+ batch_positions.push(position as u64);
+ }
+
+ // Fast path ordered batches; fall back for duplicates or
out-of-order positions.
+ if !batch_positions.is_empty()
+ && let Err(err) =
index.positions.insert_positions(&batch_positions)
+ {
+ tracing::debug!(
+ delete_file = %delete_file.file_path,
+ batch_len = batch_positions.len(),
+ error = %err,
+ "position-delete batch fell back to per-position insert"
+ );
+ for position in batch_positions {
+ index.positions.insert(position);
+ }
+ }
+ }
+
+ if let Some(expected_count) = delete_file.record_count
+ && rows_read != expected_count
+ {
+ return Err(invalid_data!(
+ "Position-delete file {} contains {rows_read} rows, expected
{expected_count} from record_count",
+ delete_file.file_path
+ ));
+ }
+
+ Ok(index)
+ }
+}
+
#[cfg(test)]
mod tests {
+ use std::collections::HashMap;
+ use std::fs::File;
+ use std::sync::Arc;
+
+ use arrow_array::{
+ Array, ArrayRef, Int32Array, Int64Array, RecordBatch, StringArray,
StructArray,
+ };
+ use arrow_schema::{DataType, Field, Schema as ArrowSchema};
+ use parquet::arrow::ArrowWriter;
+ use parquet::file::properties::WriterProperties;
use tempfile::TempDir;
use super::*;
use crate::arrow::delete_filter::tests::setup;
use crate::arrow::test_utils::write_encrypted_parquet;
+ fn write_plain_parquet(path: &str, batch: &RecordBatch) {
+ let file = File::create(path).unwrap();
+ let mut writer = ArrowWriter::try_new(file, batch.schema(),
None).unwrap();
+ writer.write(batch).unwrap();
+ writer.close().unwrap();
+ }
+
+ fn write_plain_parquet_batches(path: &str, batches: &[RecordBatch],
row_group_size: usize) {
+ let file = File::create(path).unwrap();
+ let properties = WriterProperties::builder()
+ .set_max_row_group_row_count(Some(row_group_size))
+ .build();
+ let mut writer = ArrowWriter::try_new(file, batches[0].schema(),
Some(properties)).unwrap();
+ for batch in batches {
+ writer.write(batch).unwrap();
+ }
+ writer.close().unwrap();
+ }
Review Comment:
This helper will panic on `batches[0]` if it’s ever called with an empty
slice. Even though current call sites pass non-empty arrays, adding an explicit
`assert!(!batches.is_empty())` (or returning early) makes the test helper safer
and keeps failures clearer if future tests reuse it differently.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]