diff --git a/crates/iceberg/public-api.txt b/crates/iceberg/public-api.txt index 57edfd6421..ca55804966 100644 --- a/crates/iceberg/public-api.txt +++ b/crates/iceberg/public-api.txt @@ -1306,7 +1306,7 @@ impl core::marker::StructuralPartialEq for iceberg::scan::FileScanTask impl iceberg::scan::FileScanTask pub fn iceberg::scan::FileScanTask::builder() -> FileScanTaskBuilder<((), (), (), (), (), (), (), (), (), (), (), (), (), (), (), (), (), ())> impl serde_core::ser::Serialize for iceberg::scan::FileScanTask -pub fn iceberg::scan::FileScanTask::serialize<__S>(&self, __serializer: __S) -> core::result::Result<<__S as serde_core::ser::Serializer>::Ok, <__S as serde_core::ser::Serializer>::Error> where __S: serde_core::ser::Serializer +pub fn iceberg::scan::FileScanTask::serialize(&self, serializer: S) -> core::result::Result<::Ok, ::Error> where S: serde_core::ser::Serializer impl<'de> serde_core::de::Deserialize<'de> for iceberg::scan::FileScanTask pub fn iceberg::scan::FileScanTask::deserialize<__D>(__deserializer: __D) -> core::result::Result::Error> where __D: serde_core::de::Deserializer<'de> pub struct iceberg::scan::FileScanTaskDeleteFile diff --git a/crates/iceberg/src/scan/mod.rs b/crates/iceberg/src/scan/mod.rs index 2f0a876096..db70806473 100644 --- a/crates/iceberg/src/scan/mod.rs +++ b/crates/iceberg/src/scan/mod.rs @@ -660,12 +660,13 @@ pub mod tests { RESERVED_COL_NAME_POS, RESERVED_COL_NAME_SPEC_ID, RESERVED_FIELD_ID_DELETE_FILE_PATH, RESERVED_FIELD_ID_DELETE_FILE_POS, RESERVED_FIELD_ID_POS, }; - use crate::scan::FileScanTask; + use crate::scan::{FileScanTask, FileScanTaskDeleteFile}; use crate::spec::{ DataContentType, DataFileBuilder, DataFileFormat, Datum, FormatVersion, Literal, MAIN_BRANCH, ManifestEntry, ManifestListWriter, ManifestStatus, ManifestWriterBuilder, - NestedField, Operation, PartitionSpec, PrimitiveType, Schema, Snapshot, Struct, StructType, - Summary, TableMetadata, TableMetadataBuilder, TableProperties, Type, UnboundPartitionSpec, + MappedField, NameMapping, NestedField, Operation, PartitionSpec, PrimitiveType, Schema, + Snapshot, Struct, StructType, Summary, TableMetadata, TableMetadataBuilder, + TableProperties, Transform, Type, UnboundPartitionSpec, }; use crate::table::Table; use crate::test_utils::test_runtime; @@ -2629,33 +2630,64 @@ pub mod tests { assert_eq!(string_arr.value(0), "Apache"); } - #[test] - fn test_file_scan_task_serialize_deserialize() { - let test_fn = |task: FileScanTask| { - let serialized = serde_json::to_string(&task).unwrap(); - let deserialized: FileScanTask = serde_json::from_str(&serialized).unwrap(); - - assert_eq!(task, deserialized); - }; - - // without predicate - let schema = Arc::new( + fn file_scan_task_test_schema(primitive_type: PrimitiveType) -> Arc { + Arc::new( Schema::builder() .with_fields(vec![Arc::new(NestedField::required( 1, "x", - Type::Primitive(PrimitiveType::Binary), + Type::Primitive(primitive_type), ))]) .build() .unwrap(), + ) + } + + fn assert_file_scan_task_serde_round_trip(task: FileScanTask) { + // Regression test for https://github.com/apache/iceberg-rust/issues/3089. + let serialized = serde_json::to_string(&task).unwrap(); + let deserialized: FileScanTask = serde_json::from_str(&serialized).unwrap(); + + assert_eq!(task, deserialized); + } + + fn file_scan_task_with_partition( + primitive_type: PrimitiveType, + transform: Transform, + partition_value: Literal, + ) -> FileScanTask { + let schema = file_scan_task_test_schema(primitive_type); + let partition_spec = Arc::new( + PartitionSpec::builder(schema.clone()) + .add_partition_field("x", "x_partition", transform) + .unwrap() + .build() + .unwrap(), ); + FileScanTask::builder() + .with_data_file_path("data_file_path".to_string()) + .with_file_size_in_bytes(123) + .with_start(10) + .with_length(100) + .with_project_field_ids(vec![1]) + .with_schema(schema) + .with_data_file_format(DataFileFormat::Parquet) + .with_partition(Some(Struct::from_iter([Some(partition_value)]))) + .with_partition_spec(Some(partition_spec)) + .with_case_sensitive(true) + .build() + .unwrap() + } + + #[test] + fn test_file_scan_task_serde_without_predicate() { let task = FileScanTask::builder() .with_data_file_path("data_file_path".to_string()) .with_file_size_in_bytes(0) .with_start(0) .with_length(100) .with_project_field_ids(vec![1, 2, 3]) - .with_schema(schema.clone()) + .with_schema(file_scan_task_test_schema(PrimitiveType::Binary)) .with_record_count(Some(100)) .with_first_row_id(Some(1000)) .with_data_sequence_number(Some(5)) @@ -2663,9 +2695,11 @@ pub mod tests { .with_case_sensitive(false) .build() .unwrap(); - test_fn(task); + assert_file_scan_task_serde_round_trip(task); + } - // with predicate + #[test] + fn test_file_scan_task_serde_with_predicate() { let task = FileScanTask::builder() .with_data_file_path("data_file_path".to_string()) .with_file_size_in_bytes(0) @@ -2673,12 +2707,154 @@ pub mod tests { .with_length(100) .with_project_field_ids(vec![1, 2, 3]) .with_predicate(Some(BoundPredicate::AlwaysTrue)) - .with_schema(schema) + .with_schema(file_scan_task_test_schema(PrimitiveType::Binary)) .with_data_file_format(DataFileFormat::Avro) .with_case_sensitive(false) .build() .unwrap(); - test_fn(task); + + let serialized = serde_json::to_value(&task).unwrap(); + assert!(serialized.get("record_count").is_none()); + assert_file_scan_task_serde_round_trip(task); + } + + #[test] + fn test_unpartitioned_file_scan_task_serde() { + let task = FileScanTask::builder() + .with_data_file_path("data_file_path".to_string()) + .with_file_size_in_bytes(0) + .with_start(0) + .with_length(100) + .with_project_field_ids(vec![1, 2, 3]) + .with_schema(file_scan_task_test_schema(PrimitiveType::Binary)) + .with_data_file_format(DataFileFormat::Parquet) + .with_partition(Some(Struct::empty())) + .with_case_sensitive(false) + .build() + .unwrap(); + assert_file_scan_task_serde_round_trip(task); + } + + #[test] + fn test_file_scan_task_serde_with_all_optional_fields() { + let schema = file_scan_task_test_schema(PrimitiveType::Long); + let partition_spec = Arc::new( + PartitionSpec::builder(schema.clone()) + .add_partition_field("x", "x", Transform::Identity) + .unwrap() + .build() + .unwrap(), + ); + let unified_partition_type = Arc::new(partition_spec.partition_type(&schema).unwrap()); + let task = FileScanTask::builder() + .with_data_file_path("data_file_path".to_string()) + .with_file_size_in_bytes(123) + .with_start(10) + .with_length(100) + .with_project_field_ids(vec![1]) + .with_schema(schema) + .with_data_file_format(DataFileFormat::Parquet) + .with_deletes(vec![ + FileScanTaskDeleteFile::builder() + .with_file_path("delete_file_path".to_string()) + .with_file_size_in_bytes(23) + .with_file_type(DataContentType::EqualityDeletes) + .with_file_format(DataFileFormat::Parquet) + .with_partition_spec_id(0) + .with_equality_ids(Some(vec![1])) + .with_referenced_data_file(Some("data_file_path".to_string())) + .with_content_offset(Some(12)) + .with_content_size_in_bytes(Some(34)) + .with_record_count(Some(5)) + .with_key_metadata(Some(vec![4, 5, 6].into_boxed_slice())) + .build(), + ]) + .with_partition(Some(Struct::from_iter([Some(Literal::long(42))]))) + .with_partition_spec(Some(partition_spec)) + .with_name_mapping(Some(Arc::new(NameMapping::new(vec![MappedField::new( + Some(1), + vec!["x".to_string()], + vec![], + )])))) + .with_unified_partition_type(Some(unified_partition_type)) + .with_case_sensitive(true) + .with_key_metadata(Some(vec![1, 2, 3].into_boxed_slice())) + .build() + .unwrap(); + assert_file_scan_task_serde_round_trip(task); + } + + #[test] + fn test_file_scan_task_serde_with_date_partition() { + let task = file_scan_task_with_partition( + PrimitiveType::Date, + Transform::Identity, + Literal::date(19_000), + ); + assert_file_scan_task_serde_round_trip(task); + } + + #[test] + fn test_file_scan_task_serde_with_timestamp_ns_partition() { + let task = file_scan_task_with_partition( + PrimitiveType::TimestampNs, + Transform::Identity, + Literal::timestamp_nano(1_510_871_468_123_456_789), + ); + assert_file_scan_task_serde_round_trip(task); + } + + #[test] + fn test_file_scan_task_serde_with_timestamptz_ns_partition() { + let task = file_scan_task_with_partition( + PrimitiveType::TimestamptzNs, + Transform::Identity, + Literal::timestamptz_nano(1_510_871_468_123_456_789), + ); + assert_file_scan_task_serde_round_trip(task); + } + + #[test] + fn test_file_scan_task_serde_with_decimal_partition() { + let task = file_scan_task_with_partition( + PrimitiveType::Decimal { + precision: 9, + scale: 2, + }, + Transform::Identity, + Literal::decimal(12_345), + ); + assert_file_scan_task_serde_round_trip(task); + } + + #[test] + fn test_file_scan_task_serde_with_uuid_partition() { + let task = file_scan_task_with_partition( + PrimitiveType::Uuid, + Transform::Identity, + Literal::uuid(Uuid::from_u128(0x12345678_90ab_cdef_1234_567890abcdef)), + ); + assert_file_scan_task_serde_round_trip(task); + } + + #[test] + fn test_file_scan_task_serde_with_fixed_partition() { + let task = file_scan_task_with_partition( + PrimitiveType::Fixed(4), + Transform::Identity, + Literal::fixed([1, 2, 3, 4]), + ); + assert_file_scan_task_serde_round_trip(task); + } + + #[test] + fn test_file_scan_task_serde_with_bucket_partition() { + let task = file_scan_task_with_partition( + PrimitiveType::String, + Transform::Bucket(4), + Literal::int(2), + ); + assert_file_scan_task_serde_round_trip(task); } #[tokio::test] diff --git a/crates/iceberg/src/scan/task.rs b/crates/iceberg/src/scan/task.rs index 3aff0f4ae6..f1a799fcc6 100644 --- a/crates/iceberg/src/scan/task.rs +++ b/crates/iceberg/src/scan/task.rs @@ -18,7 +18,7 @@ use std::sync::Arc; use futures::stream::BoxStream; -use serde::{Deserialize, Serialize, Serializer}; +use serde::{Deserialize, Serialize}; use typed_builder::TypedBuilder; use crate::expr::BoundPredicate; @@ -31,26 +31,9 @@ use crate::{Error, ErrorKind, Result}; /// A stream of [`FileScanTask`]. pub type FileScanTaskStream = BoxStream<'static, Result>; -/// Serialization helper that always returns NotImplementedError. -/// Used for fields that should not be serialized but we want to be explicit about it. -fn serialize_not_implemented(_: &T, _: S) -> std::result::Result -where S: Serializer { - Err(serde::ser::Error::custom( - "Serialization not implemented for this field", - )) -} - -/// Deserialization helper that always returns NotImplementedError. -/// Used for fields that should not be deserialized but we want to be explicit about it. -fn deserialize_not_implemented<'de, D, T>(_: D) -> std::result::Result -where D: serde::Deserializer<'de> { - Err(serde::de::Error::custom( - "Deserialization not implemented for this field", - )) -} - /// A task to scan part of file. -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, TypedBuilder)] +#[derive(Debug, Clone, Deserialize, PartialEq, TypedBuilder)] +#[serde(try_from = "crate::scan::task::_serde::FileScanTaskSerde")] #[builder( field_defaults(setter(prefix = "with_")), build_method(into = Result) @@ -74,7 +57,6 @@ pub struct FileScanTask { /// /// Used to derive the `_row_id` metadata column: for a row without an /// explicit `_row_id`, it is this value plus the row's ordinal position. - #[serde(skip_serializing_if = "Option::is_none")] #[builder(default)] first_row_id: Option, @@ -84,7 +66,6 @@ pub struct FileScanTask { /// manifest that lacks one. /// /// Used to derive the `_last_updated_sequence_number` metadata column. - #[serde(skip_serializing_if = "Option::is_none")] #[builder(default)] data_sequence_number: Option, @@ -99,7 +80,6 @@ pub struct FileScanTask { /// The field ids to project. project_field_ids: Vec, /// The predicate to filter. - #[serde(skip_serializing_if = "Option::is_none")] #[builder(default)] predicate: Option, @@ -110,30 +90,18 @@ pub struct FileScanTask { /// Partition data from the manifest entry, used to identify which columns can use /// constant values from partition metadata vs. reading from the data file. /// Per the Iceberg spec, only identity-transformed partition fields should use constants. - #[serde(default)] - #[serde(skip_serializing_if = "Option::is_none")] - #[serde(serialize_with = "serialize_not_implemented")] - #[serde(deserialize_with = "deserialize_not_implemented")] #[builder(default)] partition: Option, /// The partition spec for this file, used to distinguish identity transforms /// (which use partition metadata constants) from non-identity transforms like /// bucket/truncate (which must read source columns from the data file). - #[serde(default)] - #[serde(skip_serializing_if = "Option::is_none")] - #[serde(serialize_with = "serialize_not_implemented")] - #[serde(deserialize_with = "deserialize_not_implemented")] #[builder(default)] partition_spec: Option>, /// Name mapping from table metadata (property: schema.name-mapping.default), /// used to resolve field IDs from column names when Parquet files lack field IDs /// or have field ID conflicts. - #[serde(default)] - #[serde(skip_serializing_if = "Option::is_none")] - #[serde(serialize_with = "serialize_not_implemented")] - #[serde(deserialize_with = "deserialize_not_implemented")] #[builder(default)] name_mapping: Option>, @@ -145,11 +113,6 @@ pub struct FileScanTask { /// This is a table-level value (same for all tasks in a scan), stored per-task /// so that readers are self-contained without needing back-pointers to table /// metadata. The cost is one Arc clone per task. - /// Serde: not yet implemented (same pattern as partition, partition_spec, name_mapping). - #[serde(default)] - #[serde(skip_serializing_if = "Option::is_none")] - #[serde(serialize_with = "serialize_not_implemented")] - #[serde(deserialize_with = "deserialize_not_implemented")] #[builder(default)] unified_partition_type: Option>, @@ -161,11 +124,9 @@ pub struct FileScanTask { /// /// Note on the trust boundary: for the standard encryption scheme this /// carries `StandardKeyMetadata`, whose payload is the *plaintext* DEK. - /// Because `FileScanTask` derives `Serialize`, that plaintext DEK is part + /// Because `FileScanTask` implements [`Serialize`], that plaintext DEK is part /// of the serialized scan plan should these tasks ever be serialized and sent /// over the network. - #[serde(default)] - #[serde(skip_serializing_if = "Option::is_none")] #[builder(default)] key_metadata: Option>, } @@ -403,6 +364,177 @@ pub struct FileScanTaskDeleteFile { pub key_metadata: Option>, } +mod _serde { + use std::sync::Arc; + + use serde::{Deserialize, Serialize}; + + use super::{FileScanTask, FileScanTaskDeleteFile}; + use crate::expr::BoundPredicate; + use crate::spec::{ + DataFileFormat, Literal, NameMapping, PartitionSpec, RawLiteral, SchemaRef, StructType, + Type, + }; + use crate::{Error, ErrorKind, Result}; + + #[derive(Deserialize)] + pub(super) struct FileScanTaskSerde { + file_size_in_bytes: u64, + start: u64, + length: u64, + record_count: Option, + first_row_id: Option, + data_sequence_number: Option, + data_file_path: String, + data_file_format: DataFileFormat, + schema: SchemaRef, + project_field_ids: Vec, + predicate: Option, + deletes: Vec, + #[serde(default)] + partition: Option, + #[serde(default)] + partition_spec: Option>, + #[serde(default)] + name_mapping: Option>, + #[serde(default)] + unified_partition_type: Option>, + case_sensitive: bool, + #[serde(default)] + key_metadata: Option>, + } + + #[derive(Serialize)] + struct FileScanTaskRefSerde<'a> { + file_size_in_bytes: u64, + start: u64, + length: u64, + #[serde(skip_serializing_if = "Option::is_none")] + record_count: Option, + #[serde(skip_serializing_if = "Option::is_none")] + first_row_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + data_sequence_number: Option, + data_file_path: &'a str, + data_file_format: DataFileFormat, + schema: &'a SchemaRef, + project_field_ids: &'a [i32], + #[serde(skip_serializing_if = "Option::is_none")] + predicate: Option<&'a BoundPredicate>, + deletes: &'a [FileScanTaskDeleteFile], + #[serde(skip_serializing_if = "Option::is_none")] + partition: Option, + #[serde(skip_serializing_if = "Option::is_none")] + partition_spec: Option<&'a Arc>, + #[serde(skip_serializing_if = "Option::is_none")] + name_mapping: Option<&'a Arc>, + #[serde(skip_serializing_if = "Option::is_none")] + unified_partition_type: Option<&'a Arc>, + case_sensitive: bool, + #[serde(skip_serializing_if = "Option::is_none")] + key_metadata: Option<&'a [u8]>, + } + + fn partition_type( + partition_spec: Option<&PartitionSpec>, + schema: &crate::spec::Schema, + ) -> Result { + let partition_type = match partition_spec { + Some(partition_spec) => partition_spec.partition_type(schema)?, + None => PartitionSpec::unpartition_spec().partition_type(schema)?, + }; + Ok(Type::Struct(partition_type)) + } + + impl<'a> TryFrom<&'a FileScanTask> for FileScanTaskRefSerde<'a> { + type Error = Error; + + fn try_from(value: &'a FileScanTask) -> Result { + let partition = value + .partition + .as_ref() + .map(|partition| { + let partition_type = + partition_type(value.partition_spec.as_deref(), &value.schema)?; + RawLiteral::try_from(Literal::Struct(partition.clone()), &partition_type) + }) + .transpose()?; + + Ok(Self { + file_size_in_bytes: value.file_size_in_bytes, + start: value.start, + length: value.length, + record_count: value.record_count, + first_row_id: value.first_row_id, + data_sequence_number: value.data_sequence_number, + data_file_path: &value.data_file_path, + data_file_format: value.data_file_format, + schema: &value.schema, + project_field_ids: &value.project_field_ids, + predicate: value.predicate.as_ref(), + deletes: &value.deletes, + partition, + partition_spec: value.partition_spec.as_ref(), + name_mapping: value.name_mapping.as_ref(), + unified_partition_type: value.unified_partition_type.as_ref(), + case_sensitive: value.case_sensitive, + key_metadata: value.key_metadata.as_deref(), + }) + } + } + + impl Serialize for FileScanTask { + fn serialize(&self, serializer: S) -> std::result::Result + where S: serde::Serializer { + FileScanTaskRefSerde::try_from(self) + .map_err(serde::ser::Error::custom)? + .serialize(serializer) + } + } + + impl TryFrom for FileScanTask { + type Error = Error; + + fn try_from(value: FileScanTaskSerde) -> Result { + let partition = value + .partition + .map(|partition| { + let partition_type = + partition_type(value.partition_spec.as_deref(), &value.schema)?; + match partition.try_into(&partition_type)? { + Some(Literal::Struct(partition)) => Ok(partition), + _ => Err(Error::new( + ErrorKind::DataInvalid, + "FileScanTask partition must be a struct", + )), + } + }) + .transpose()?; + + Self::builder() + .with_file_size_in_bytes(value.file_size_in_bytes) + .with_start(value.start) + .with_length(value.length) + .with_record_count(value.record_count) + .with_first_row_id(value.first_row_id) + .with_data_sequence_number(value.data_sequence_number) + .with_data_file_path(value.data_file_path) + .with_data_file_format(value.data_file_format) + .with_schema(value.schema) + .with_project_field_ids(value.project_field_ids) + .with_predicate(value.predicate) + .with_deletes(value.deletes) + .with_partition(partition) + .with_partition_spec(value.partition_spec) + .with_name_mapping(value.name_mapping) + .with_unified_partition_type(value.unified_partition_type) + .with_case_sensitive(value.case_sensitive) + .with_key_metadata(value.key_metadata) + .build() + } + } +} + #[cfg(test)] mod tests { use super::*; diff --git a/crates/iceberg/src/spec/values/serde.rs b/crates/iceberg/src/spec/values/serde.rs index 10a815aa45..8a00a79dee 100644 --- a/crates/iceberg/src/spec/values/serde.rs +++ b/crates/iceberg/src/spec/values/serde.rs @@ -418,6 +418,12 @@ pub(crate) mod _serde { Type::Primitive(PrimitiveType::Timestamptz) => { Ok(Some(Literal::timestamptz(v))) } + Type::Primitive(PrimitiveType::TimestampNs) => { + Ok(Some(Literal::timestamp_nano(v))) + } + Type::Primitive(PrimitiveType::TimestamptzNs) => { + Ok(Some(Literal::timestamptz_nano(v))) + } _ => Err(invalid_err("long")), }, RawLiteralEnum::Float(v) => match ty { diff --git a/crates/iceberg/src/spec/values/tests.rs b/crates/iceberg/src/spec/values/tests.rs index 8c72332f13..1d716a48da 100644 --- a/crates/iceberg/src/spec/values/tests.rs +++ b/crates/iceberg/src/spec/values/tests.rs @@ -45,6 +45,17 @@ fn check_json_serde(json: &str, expected_literal: Literal, expected_type: &Type) assert_eq!(parsed_json_value, raw_json_value); } +fn check_raw_literal_json_serde(expected_literal: Literal, expected_type: &Type) { + let raw_literal = RawLiteral::try_from(expected_literal.clone(), expected_type).unwrap(); + let serialized = serde_json::to_string(&raw_literal).unwrap(); + let deserialized: RawLiteral = serde_json::from_str(&serialized).unwrap(); + + assert_eq!( + deserialized.try_into(expected_type).unwrap(), + Some(expected_literal) + ); +} + fn check_avro_bytes_serde(input: Vec, expected_datum: Datum, expected_type: &PrimitiveType) { let raw_schema = r#""bytes""#; let schema = apache_avro::Schema::parse_str(raw_schema).unwrap(); @@ -235,6 +246,18 @@ fn json_timestamptz_ns() { ); } +#[test] +fn raw_literal_json_serde_nanosecond_timestamps() { + check_raw_literal_json_serde( + Literal::timestamp_nano(1510871468123456789), + &Primitive(PrimitiveType::TimestampNs), + ); + check_raw_literal_json_serde( + Literal::timestamptz_nano(1510871468123456789), + &Primitive(PrimitiveType::TimestamptzNs), + ); +} + #[test] fn json_timestamptz_ns_rejects_non_utc_offset() { // Per the spec, timestamptz_ns single-value serialization must use offset "+00:00"; Java's