Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
78 changes: 26 additions & 52 deletions crates/iceberg/src/arrow/caching_delete_file_loader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ use crate::arrow::scan_metrics::ScanMetrics;
use crate::arrow::{arrow_primitive_to_literal, arrow_schema_to_schema};
use crate::delete_vector::DeleteVector;
use crate::encryption::{EncryptedInputFile, StandardKeyMetadata};
use crate::error::invalid_data;
use crate::expr::Predicate::AlwaysTrue;
use crate::expr::{Predicate, Reference};
use crate::io::FileIO;
Expand Down Expand Up @@ -348,55 +349,37 @@ impl CachingDeleteFileLoader {
task: &FileScanTaskDeleteFile,
) -> Result<(u64, u64, String, u64)> {
let content_offset = task.content_offset.ok_or_else(|| {
Error::new(
ErrorKind::DataInvalid,
format!(
"deletion vector {} is missing content_offset",
task.file_path
),
invalid_data!(
"deletion vector {} is missing content_offset",
task.file_path
)
})?;
let content_size = task.content_size_in_bytes.ok_or_else(|| {
Error::new(
ErrorKind::DataInvalid,
format!(
"deletion vector {} is missing content_size_in_bytes",
task.file_path
),
invalid_data!(
"deletion vector {} is missing content_size_in_bytes",
task.file_path
)
})?;
let data_file_path = task.referenced_data_file.clone().ok_or_else(|| {
Error::new(
ErrorKind::DataInvalid,
format!(
"deletion vector {} is missing referenced_data_file",
task.file_path
),
invalid_data!(
"deletion vector {} is missing referenced_data_file",
task.file_path
)
})?;
let record_count = task.record_count.ok_or_else(|| {
Error::new(
ErrorKind::DataInvalid,
format!("deletion vector {} is missing record_count", task.file_path),
)
invalid_data!("deletion vector {} is missing record_count", task.file_path)
})?;

let start = u64::try_from(content_offset).map_err(|_| {
Error::new(
ErrorKind::DataInvalid,
format!(
"deletion vector {} has negative content_offset {content_offset}",
task.file_path
),
invalid_data!(
"deletion vector {} has negative content_offset {content_offset}",
task.file_path
)
})?;
let len = u64::try_from(content_size).map_err(|_| {
Error::new(
ErrorKind::DataInvalid,
format!(
"deletion vector {} has negative content_size_in_bytes {content_size}",
task.file_path
),
invalid_data!(
"deletion vector {} has negative content_size_in_bytes {content_size}",
task.file_path
)
})?;

Expand All @@ -412,11 +395,8 @@ impl CachingDeleteFileLoader {
) -> Result<()> {
let actual = delete_vector.len();
if actual != expected {
return Err(Error::new(
ErrorKind::DataInvalid,
format!(
"deletion vector {dv_path} decoded to {actual} positions, expected {expected} from record_count"
),
return Err(invalid_data!(
"deletion vector {dv_path} decoded to {actual} positions, expected {expected} from record_count"
));
}
Ok(())
Expand Down Expand Up @@ -528,15 +508,13 @@ impl CachingDeleteFileLoader {
let columns = batch.columns();

let Some(file_paths) = columns[0].as_any().downcast_ref::<StringArray>() else {
return Err(Error::new(
ErrorKind::DataInvalid,
"Could not downcast file paths array to StringArray",
return Err(invalid_data!(
"Could not downcast file paths array to StringArray"
));
};
let Some(positions) = columns[1].as_any().downcast_ref::<Int64Array>() else {
return Err(Error::new(
ErrorKind::DataInvalid,
"Could not downcast positions array to Int64Array",
return Err(invalid_data!(
"Could not downcast positions array to Int64Array"
));
};

Expand All @@ -551,15 +529,11 @@ impl CachingDeleteFileLoader {

for (file_path, pos) in file_paths.iter().zip(positions.iter()) {
let (Some(file_path), Some(pos)) = (file_path, pos) else {
return Err(Error::new(
ErrorKind::DataInvalid,
"null values in delete file",
));
return Err(invalid_data!("null values in delete file"));
};
if pos < 0 {
return Err(Error::new(
ErrorKind::DataInvalid,
format!("negative position in delete file {file_path}: {pos}"),
return Err(invalid_data!(
"negative position in delete file {file_path}: {pos}"
));
}

Expand Down
20 changes: 6 additions & 14 deletions crates/iceberg/src/arrow/partition_value_calculator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,9 +27,10 @@ use arrow_schema::DataType;

use super::record_batch_projector::RecordBatchProjector;
use super::type_to_arrow_type;
use crate::Result;
use crate::error::invalid_data;
use crate::spec::{PartitionSpec, Schema, StructType, Type};
use crate::transform::{BoxedTransformFunction, create_transform_function};
use crate::{Error, ErrorKind, Result};

/// Calculator for partition values in Iceberg tables.
///
Expand Down Expand Up @@ -63,9 +64,8 @@ impl PartitionValueCalculator {
/// - Projector initialization fails
pub fn try_new(partition_spec: &PartitionSpec, table_schema: &Schema) -> Result<Self> {
if partition_spec.is_unpartitioned() {
return Err(Error::new(
ErrorKind::DataInvalid,
"Cannot create partition calculator for unpartitioned table",
return Err(invalid_data!(
"Cannot create partition calculator for unpartitioned table"
));
}

Expand Down Expand Up @@ -140,10 +140,7 @@ impl PartitionValueCalculator {
let expected_struct_fields = match &self.partition_arrow_type {
DataType::Struct(fields) => fields.clone(),
_ => {
return Err(Error::new(
ErrorKind::DataInvalid,
"Expected partition type must be a struct",
));
return Err(invalid_data!("Expected partition type must be a struct"));
}
};

Expand All @@ -156,12 +153,7 @@ impl PartitionValueCalculator {

// Construct the StructArray
let struct_array = StructArray::try_new(expected_struct_fields, partition_values, None)
.map_err(|e| {
Error::new(
ErrorKind::DataInvalid,
format!("Failed to create partition struct array: {e}"),
)
})?;
.map_err(|e| invalid_data!("Failed to create partition struct array: {e}"))?;

Ok(Arc::new(struct_array))
}
Expand Down
11 changes: 4 additions & 7 deletions crates/iceberg/src/arrow/reader/pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ use crate::arrow::int96::coerce_int96_timestamps;
use crate::arrow::record_batch_transformer::RecordBatchTransformerBuilder;
use crate::arrow::scan_metrics::{CountingFileRead, ScanMetrics, ScanResult};
use crate::encryption::StandardKeyMetadata;
use crate::error::Result;
use crate::error::{Result, invalid_data};
use crate::expr::BoundPredicate;
use crate::expr::visitors::bloom_filter_evaluator::{
BloomFilterEvaluator, ColumnBloomFilter, collect_bloom_filter_field_ids,
Expand Down Expand Up @@ -475,12 +475,9 @@ impl FileScanTaskReader {
// inheritance a committed entry always has one, so this is a malformed
// manifest rather than a legitimate null.
(Some(_), None) => {
return Err(Error::new(
ErrorKind::DataInvalid,
format!(
"Data file {} has a first_row_id but no data sequence number",
task.data_file_path()
),
return Err(invalid_data!(
"Data file {} has a first_row_id but no data sequence number",
task.data_file_path()
));
}
};
Expand Down
19 changes: 6 additions & 13 deletions crates/iceberg/src/arrow/reader/predicate_visitor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -35,11 +35,10 @@ use fnv::FnvHashSet;
use parquet::schema::types::SchemaDescriptor;

use crate::arrow::get_arrow_datum;
use crate::error::Result;
use crate::error::{Result, invalid_data};
use crate::expr::visitors::bound_predicate_visitor::BoundPredicateVisitor;
use crate::expr::{BoundPredicate, BoundReference};
use crate::spec::Datum;
use crate::{Error, ErrorKind};

/// A visitor to collect field ids from bound predicates.
pub(super) struct CollectFieldIdVisitor {
Expand Down Expand Up @@ -215,12 +214,9 @@ impl PredicateConverter<'_> {
// The leaf column's index in Parquet schema.
if let Some(column_idx) = self.column_map.get(&reference.field().id) {
if self.parquet_schema.get_column_root(*column_idx).is_group() {
return Err(Error::new(
ErrorKind::DataInvalid,
format!(
"Leaf column `{}` in predicates isn't a root column in Parquet schema.",
reference.field().name
),
return Err(invalid_data!(
"Leaf column `{}` in predicates isn't a root column in Parquet schema.",
reference.field().name
));
}

Expand All @@ -229,13 +225,10 @@ impl PredicateConverter<'_> {
.column_indices
.iter()
.position(|&idx| idx == *column_idx)
.ok_or(Error::new(
ErrorKind::DataInvalid,
format!(
.ok_or(invalid_data!(
"Leaf column `{}` in predicates cannot be found in the required column indices.",
reference.field().name
),
))?;
))?;

Ok(Some(index))
} else {
Expand Down
9 changes: 3 additions & 6 deletions crates/iceberg/src/arrow/reader/projection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@ use parquet::schema::types::{SchemaDescriptor, Type as ParquetType};

use super::{ArrowReader, CollectFieldIdVisitor};
use crate::arrow::arrow_schema_to_schema;
use crate::error::Result;
use crate::error::{Result, invalid_data};
use crate::expr::BoundPredicate;
use crate::expr::visitors::bound_predicate_visitor::visit;
use crate::spec::{NameMapping, NestedField, PrimitiveType, Schema, Type};
Expand Down Expand Up @@ -299,11 +299,8 @@ pub(super) fn build_field_id_map(
column_map.insert(basic_info.id(), idx);
}
ParquetType::GroupType { .. } => {
return Err(Error::new(
ErrorKind::DataInvalid,
format!(
"Leaf column in schema should be primitive type but got {field_type:?}"
),
return Err(invalid_data!(
"Leaf column in schema should be primitive type but got {field_type:?}"
));
}
};
Expand Down
7 changes: 4 additions & 3 deletions crates/iceberg/src/arrow/reader/row_lineage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ use arrow_schema::{DataType, Field, Schema};
use arrow_select::zip::zip;
use parquet::arrow::PARQUET_FIELD_ID_META_KEY;

use crate::error::invalid_data;
use crate::metadata_columns::{
RESERVED_COL_NAME_ROW_ID, RESERVED_FIELD_ID_POS, RESERVED_FIELD_ID_ROW_ID,
};
Expand Down Expand Up @@ -75,9 +76,9 @@ pub(crate) fn synthesize_row_id_column(
match column_by_field_id(&batch, RESERVED_FIELD_ID_ROW_ID) {
Some(id) => {
if id.data_type() != &DataType::Int64 {
return Err(Error::new(
ErrorKind::DataInvalid,
format!("_row_id source must be Int64, got {}", id.data_type()),
return Err(invalid_data!(
"_row_id source must be Int64, got {}",
id.data_type()
));
}

Expand Down
27 changes: 9 additions & 18 deletions crates/iceberg/src/arrow/record_batch_partition_splitter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,8 +24,9 @@ use arrow_select::filter::filter_record_batch;

use super::arrow_struct_to_literal;
use super::partition_value_calculator::PartitionValueCalculator;
use crate::Result;
use crate::error::invalid_data;
use crate::spec::{Literal, PartitionKey, PartitionSpecRef, SchemaRef, StructType};
use crate::{Error, ErrorKind, Result};

/// Column name for the projected partition values struct
pub const PROJECTED_PARTITION_VALUE_COLUMN: &str = "_partition";
Expand Down Expand Up @@ -129,9 +130,8 @@ impl RecordBatchPartitionSplitter {
if let Some(Literal::Struct(s)) = s {
Ok(s)
} else {
Err(Error::new(
ErrorKind::DataInvalid,
"Partition value is not a struct literal or is null",
Err(invalid_data!(
"Partition value is not a struct literal or is null"
))
}
})
Expand All @@ -141,23 +141,15 @@ impl RecordBatchPartitionSplitter {
let partition_column = batch
.column_by_name(PROJECTED_PARTITION_VALUE_COLUMN)
.ok_or_else(|| {
Error::new(
ErrorKind::DataInvalid,
format!(
"Partition column '{PROJECTED_PARTITION_VALUE_COLUMN}' not found in batch"
),
invalid_data!(
"Partition column '{PROJECTED_PARTITION_VALUE_COLUMN}' not found in batch"
)
})?;

let partition_struct_array = partition_column
.as_any()
.downcast_ref::<StructArray>()
.ok_or_else(|| {
Error::new(
ErrorKind::DataInvalid,
"Partition column is not a StructArray",
)
})?;
.ok_or_else(|| invalid_data!("Partition column is not a StructArray"))?;

let arrow_struct_array = Arc::new(partition_struct_array.clone()) as ArrayRef;
let struct_array = arrow_struct_to_literal(&arrow_struct_array, &self.partition_type)?;
Expand All @@ -168,9 +160,8 @@ impl RecordBatchPartitionSplitter {
if let Some(Literal::Struct(s)) = s {
Ok(s)
} else {
Err(Error::new(
ErrorKind::DataInvalid,
"Partition value is not a struct literal or is null",
Err(invalid_data!(
"Partition value is not a struct literal or is null"
))
}
})
Expand Down
13 changes: 5 additions & 8 deletions crates/iceberg/src/arrow/record_batch_projector.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ use arrow_schema::{DataType, Field, FieldRef, Fields, Schema, SchemaRef};
use parquet::arrow::PARQUET_FIELD_ID_META_KEY;

use crate::arrow::schema::schema_to_arrow_schema;
use crate::error::Result;
use crate::error::{Result, invalid_data};
use crate::spec::Schema as IcebergSchema;
use crate::{Error, ErrorKind};

Expand Down Expand Up @@ -97,12 +97,9 @@ impl RecordBatchProjector {
let field_id_fetch_func = |field: &Field| -> Result<Option<i64>> {
if let Some(value) = field.metadata().get(PARQUET_FIELD_ID_META_KEY) {
let field_id = value.parse::<i32>().map_err(|e| {
Error::new(
ErrorKind::DataInvalid,
"Failed to parse field id".to_string(),
)
.with_context("value", value)
.with_source(e)
invalid_data!("Failed to parse field id")
.with_context("value", value)
.with_source(e)
})?;
Ok(Some(field_id as i64))
} else {
Expand Down Expand Up @@ -167,7 +164,7 @@ impl RecordBatchProjector {
self.projected_schema.clone(),
self.project_column(batch.columns())?,
)
.map_err(|err| Error::new(ErrorKind::DataInvalid, format!("{err}")))
.map_err(|err| invalid_data!("{err}"))
}

/// Do projection with columns
Expand Down
Loading
Loading