Skip to content
Open
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
6 changes: 6 additions & 0 deletions crates/catalog/glue/src/schema.rs
Original file line number Diff line number Diff line change
Expand Up @@ -149,6 +149,12 @@ impl SchemaVisitor for GlueSchemaBuilder {

fn primitive(&mut self, p: &PrimitiveType) -> Result<Self::T> {
let glue_type = match p {
PrimitiveType::Unknown => {
return Err(Error::new(
ErrorKind::FeatureUnsupported,
format!("Conversion from {p:?} is not supported"),
));
}
PrimitiveType::Boolean => "boolean".to_string(),
PrimitiveType::Int => "int".to_string(),
PrimitiveType::Long => "bigint".to_string(),
Expand Down
6 changes: 6 additions & 0 deletions crates/catalog/hms/src/schema.rs
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,12 @@ impl SchemaVisitor for HiveSchemaBuilder {

fn primitive(&mut self, p: &PrimitiveType) -> Result<String> {
let hive_type = match p {
PrimitiveType::Unknown => {
return Err(Error::new(
ErrorKind::FeatureUnsupported,
format!("Conversion from {p:?} is not supported"),
));
}
PrimitiveType::Boolean => "boolean".to_string(),
PrimitiveType::Int => "int".to_string(),
PrimitiveType::Long => "bigint".to_string(),
Expand Down
1 change: 1 addition & 0 deletions crates/iceberg/public-api.txt
Original file line number Diff line number Diff line change
Expand Up @@ -1659,6 +1659,7 @@ pub iceberg::spec::PrimitiveType::Timestamp
pub iceberg::spec::PrimitiveType::TimestampNs
pub iceberg::spec::PrimitiveType::Timestamptz
pub iceberg::spec::PrimitiveType::TimestamptzNs
pub iceberg::spec::PrimitiveType::Unknown
pub iceberg::spec::PrimitiveType::Uuid
impl iceberg::spec::PrimitiveType
pub fn iceberg::spec::PrimitiveType::compatible(&self, literal: &iceberg::spec::PrimitiveLiteral) -> bool
Expand Down
90 changes: 89 additions & 1 deletion crates/iceberg/src/arrow/reader/pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -320,7 +320,10 @@ impl FileScanTaskReader {
.schema()
.fields()
.iter()
.any(|f| f.name() == RESERVED_COL_NAME_ROW_ID);
.any(|f| {
f.name() == RESERVED_COL_NAME_ROW_ID
&& f.metadata().get(PARQUET_FIELD_ID_META_KEY).is_none()
});

// Read the physical column only when first_row_id is set. A null first_row_id
// nulls the whole column below (matching Java `ValueReaders.rowIds`).
Expand Down Expand Up @@ -2291,6 +2294,91 @@ mod tests {
assert_row_id_column(&batches, &[Some(100), Some(101), Some(102)]);
}

#[tokio::test]
async fn test_row_id_name_collision_with_user_field_id() {
let tmp_dir = TempDir::new().unwrap();
let dir = tmp_dir.path().to_str().unwrap();

// A user column may share the reserved metadata column's name. Its non-reserved
// field id distinguishes it from a physically-stored metadata column, so preserve
// the user values and synthesize the projected metadata column independently.
let user_field_id = 2;
let user_row_id_field = Field::new(RESERVED_COL_NAME_ROW_ID, DataType::Int64, true)
.with_metadata(HashMap::from([(
PARQUET_FIELD_ID_META_KEY.to_string(),
user_field_id.to_string(),
)]));
let user_values = Arc::new(Int64Array::from(vec![7i64, 8, 9])) as ArrayRef;
let file_path = write_plain_parquet(
dir,
"user_row_id_with_field_id.parquet",
vec![user_row_id_field],
vec![user_values],
);

let schema = Arc::new(
Schema::builder()
.with_schema_id(1)
.with_fields(vec![
NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
NestedField::optional(
user_field_id,
RESERVED_COL_NAME_ROW_ID,
Type::Primitive(PrimitiveType::Long),
)
.into(),
])
.build()
.unwrap(),
);
let task = FileScanTask::builder()
.with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
.with_start(0)
.with_length(0)
.with_data_file_path(file_path)
.with_data_file_format(DataFileFormat::Parquet)
.with_schema(schema)
.with_project_field_ids(vec![user_field_id, RESERVED_FIELD_ID_ROW_ID])
.with_first_row_id(Some(100))
.with_case_sensitive(false)
.build();

let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
let batches: Vec<RecordBatch> = reader
.read(tasks)
.unwrap()
.stream()
.try_collect()
.await
.unwrap();

let batch = &batches[0];
let values_by_field_id = |field_id: i32| {
let index = batch
.schema()
.fields()
.iter()
.position(|field| {
field
.metadata()
.get(PARQUET_FIELD_ID_META_KEY)
.is_some_and(|id| id == &field_id.to_string())
})
.unwrap();
batch
.column(index)
.as_primitive::<arrow_array::types::Int64Type>()
.values()
.to_vec()
};

assert_eq!(values_by_field_id(user_field_id), vec![7, 8, 9]);
assert_eq!(values_by_field_id(RESERVED_FIELD_ID_ROW_ID), vec![
100, 101, 102
]);
}

#[tokio::test]
async fn test_read_encrypted_parquet_with_wrong_key_fails() {
let encryption_key = b"0123456789abcdef";
Expand Down
14 changes: 14 additions & 0 deletions crates/iceberg/src/arrow/reader/row_lineage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,13 @@ pub(crate) fn synthesize_row_id_column(
let row_id: ArrayRef = match first_row_id {
None => Arc::new(Int64Array::new_null(batch.num_rows())),
Some(base) => {
if base < 0 {
return Err(Error::new(
ErrorKind::DataInvalid,
format!("first_row_id must be non-negative, got {base}"),
));
}

let pos = column_by_field_id(&batch, RESERVED_FIELD_ID_POS).ok_or_else(|| {
Error::new(
ErrorKind::Unexpected,
Expand Down Expand Up @@ -204,6 +211,13 @@ mod tests {
assert_eq!(row_ids(&out), vec![Some(110), Some(111), Some(112)]);
}

#[test]
fn negative_first_row_id_is_rejected() {
let err = synthesize_row_id_column(batch(vec![0, 1, 2], None), Some(-1)).unwrap_err();
assert_eq!(err.kind(), ErrorKind::DataInvalid);
assert!(format!("{err}").contains("first_row_id must be non-negative"));
}

#[test]
fn coalesce_with_physical_column() {
let out = synthesize_row_id_column(
Expand Down
11 changes: 11 additions & 0 deletions crates/iceberg/src/arrow/schema.rs
Original file line number Diff line number Diff line change
Expand Up @@ -191,6 +191,7 @@ fn visit_type<V: ArrowSchemaVisitor>(r#type: &DataType, visitor: &mut V) -> Resu
| DataType::Utf8
| DataType::LargeUtf8
| DataType::Utf8View
| DataType::Null
| DataType::Binary
| DataType::LargeBinary
| DataType::BinaryView
Expand Down Expand Up @@ -509,6 +510,7 @@ impl ArrowSchemaVisitor for ArrowSchemaConverter {

fn primitive(&mut self, p: &DataType) -> Result<Self::T> {
match p {
DataType::Null => Ok(Type::Primitive(PrimitiveType::Unknown)),
DataType::Boolean => Ok(Type::Primitive(PrimitiveType::Boolean)),
DataType::Int8 | DataType::Int16 | DataType::Int32 => {
Ok(Type::Primitive(PrimitiveType::Int))
Expand Down Expand Up @@ -697,6 +699,7 @@ impl SchemaVisitor for ToArrowSchemaConverter {

fn primitive(&mut self, p: &PrimitiveType) -> Result<ArrowSchemaOrFieldOrType> {
match p {
PrimitiveType::Unknown => Ok(ArrowSchemaOrFieldOrType::Type(DataType::Null)),
PrimitiveType::Boolean => Ok(ArrowSchemaOrFieldOrType::Type(DataType::Boolean)),
PrimitiveType::Int => Ok(ArrowSchemaOrFieldOrType::Type(DataType::Int32)),
PrimitiveType::Long => Ok(ArrowSchemaOrFieldOrType::Type(DataType::Int64)),
Expand Down Expand Up @@ -1209,6 +1212,7 @@ pub(crate) fn primitive_type_to_arrow_type_with_ree(primitive_type: &PrimitiveTy
};

match primitive_type {
PrimitiveType::Unknown => make_ree(DataType::Null),
PrimitiveType::Boolean => make_ree(DataType::Boolean),
PrimitiveType::Int => make_ree(DataType::Int32),
PrimitiveType::Long => make_ree(DataType::Int64),
Expand Down Expand Up @@ -2227,6 +2231,13 @@ mod tests {
assert_eq!(iceberg_type, arrow_type_to_type(&arrow_type).unwrap());
}

{
let arrow_type = DataType::Null;
let iceberg_type = Type::Primitive(PrimitiveType::Unknown);
assert_eq!(arrow_type, type_to_arrow_type(&iceberg_type).unwrap());
assert_eq!(iceberg_type, arrow_type_to_type(&arrow_type).unwrap());
}

// test struct type
{
// no metadata will cause error
Expand Down
24 changes: 23 additions & 1 deletion crates/iceberg/src/arrow/value.rs
Original file line number Diff line number Diff line change
Expand Up @@ -206,6 +206,7 @@ impl SchemaWithPartnerVisitor<ArrayRef> for ArrowArrayToIcebergStructConverter {

fn primitive(&mut self, p: &PrimitiveType, partner: &ArrayRef) -> Result<Vec<Option<Literal>>> {
match p {
PrimitiveType::Unknown => Ok(vec![None; partner.len()]),
PrimitiveType::Boolean => {
let array = partner
.as_any()
Expand Down Expand Up @@ -636,6 +637,7 @@ pub(crate) fn create_primitive_array_single_element(
prim_lit: &Option<PrimitiveLiteral>,
) -> Result<ArrayRef> {
match (data_type, prim_lit) {
(DataType::Null, None) => Ok(Arc::new(arrow_array::NullArray::new(1))),
(DataType::Boolean, Some(PrimitiveLiteral::Boolean(v))) => {
Ok(Arc::new(BooleanArray::from(vec![*v])))
}
Expand Down Expand Up @@ -958,7 +960,7 @@ pub(crate) fn create_primitive_array_repeated(
Some(NullBuffer::new_null(num_rows)),
))
}
(DataType::Null, _) => Arc::new(arrow_array::NullArray::new(num_rows)),
(DataType::Null, None) => Arc::new(arrow_array::NullArray::new(num_rows)),

// --- Catch-all null arm: use arrow-rs new_null_array for any remaining DataType ---
(dt, None) => new_null_array(dt, num_rows),
Expand Down Expand Up @@ -1881,6 +1883,26 @@ mod test {
}
}

#[test]
fn test_create_null_array_rejects_non_null_literal() {
let literal = Some(PrimitiveLiteral::Int(1));

assert!(create_primitive_array_single_element(&DataType::Null, &literal).is_err());
assert!(create_primitive_array_repeated(&DataType::Null, &literal, 2).is_err());
assert_eq!(
create_primitive_array_single_element(&DataType::Null, &None)
.unwrap()
.len(),
1
);
assert_eq!(
create_primitive_array_repeated(&DataType::Null, &None, 2)
.unwrap()
.len(),
2
);
}

#[test]
fn test_create_decimal_array_repeated_respects_precision() {
// Ensure repeated arrays also respect target precision, not Arrow's default.
Expand Down
Loading
Loading