1use std::sync::Arc;
22
23use itertools::Itertools;
24use serde::{Deserialize, Serialize};
25use typed_builder::TypedBuilder;
26
27use super::transform::Transform;
28use super::{NestedField, PrimitiveType, Schema, SchemaRef, StructType};
29use crate::Result;
30use crate::error::invalid_data;
31use crate::spec::Struct;
32
33pub(crate) const UNPARTITIONED_LAST_ASSIGNED_ID: i32 = 999;
34pub(crate) const DEFAULT_PARTITION_SPEC_ID: i32 = 0;
35
36#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone, TypedBuilder)]
38#[serde(rename_all = "kebab-case")]
39pub struct PartitionField {
40 pub source_id: i32,
42 pub field_id: i32,
45 pub name: String,
47 pub transform: Transform,
49}
50
51impl PartitionField {
52 pub fn into_unbound(self) -> UnboundPartitionField {
54 self.into()
55 }
56}
57
58pub type PartitionSpecRef = Arc<PartitionSpec>;
60#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone)]
67#[serde(rename_all = "kebab-case")]
68pub struct PartitionSpec {
69 spec_id: i32,
71 fields: Vec<PartitionField>,
73}
74
75impl PartitionSpec {
76 pub fn builder(schema: impl Into<SchemaRef>) -> PartitionSpecBuilder {
78 PartitionSpecBuilder::new(schema)
79 }
80
81 pub fn fields(&self) -> &[PartitionField] {
83 &self.fields
84 }
85
86 pub fn spec_id(&self) -> i32 {
88 self.spec_id
89 }
90
91 pub fn unpartition_spec() -> Self {
93 Self {
94 spec_id: DEFAULT_PARTITION_SPEC_ID,
95 fields: vec![],
96 }
97 }
98
99 pub fn is_unpartitioned(&self) -> bool {
103 self.fields.is_empty() || self.fields.iter().all(|f| f.transform == Transform::Void)
104 }
105
106 pub fn partition_type(&self, schema: &Schema) -> Result<StructType> {
110 let mut struct_fields = Vec::with_capacity(self.fields.len());
111 for partition_field in &self.fields {
112 let res_type = match schema.field_by_id(partition_field.source_id) {
113 Some(field) => partition_field.transform.result_type(&field.field_type)?,
114 None => match partition_field.transform {
117 Transform::Bucket(_) | Transform::Year | Transform::Month | Transform::Hour => {
118 PrimitiveType::Int.into()
119 }
120 Transform::Day => PrimitiveType::Date.into(),
121 Transform::Unknown => PrimitiveType::String.into(),
122 Transform::Identity | Transform::Truncate(_) | Transform::Void => {
123 PrimitiveType::Unknown.into()
124 }
125 },
126 };
127 struct_fields.push(
128 NestedField::optional(partition_field.field_id, &partition_field.name, res_type)
129 .into(),
130 );
131 }
132 Ok(StructType::new(struct_fields))
133 }
134
135 pub fn into_unbound(self) -> UnboundPartitionSpec {
137 self.into()
138 }
139
140 pub fn with_spec_id(self, spec_id: i32) -> Self {
142 Self { spec_id, ..self }
143 }
144
145 pub fn has_sequential_ids(&self) -> bool {
149 has_sequential_ids(self.fields.iter().map(|f| f.field_id))
150 }
151
152 pub fn highest_field_id(&self) -> Option<i32> {
154 self.fields.iter().map(|f| f.field_id).max()
155 }
156
157 pub fn is_compatible_with(&self, other: &PartitionSpec) -> bool {
167 if self.fields.len() != other.fields.len() {
168 return false;
169 }
170
171 for (this_field, other_field) in self.fields.iter().zip(other.fields.iter()) {
172 if this_field.source_id != other_field.source_id
173 || this_field.name != other_field.name
174 || this_field.transform != other_field.transform
175 {
176 return false;
177 }
178 }
179
180 true
181 }
182
183 pub fn partition_to_path(&self, data: &Struct, schema: SchemaRef) -> String {
186 let partition_type = self.partition_type(&schema).unwrap();
187 let field_types = partition_type.fields();
188
189 self.fields
190 .iter()
191 .enumerate()
192 .map(|(i, field)| {
193 let value = data[i].as_ref();
194 form_urlencoded::Serializer::new(String::new())
195 .append_pair(
196 &field.name,
197 &field
198 .transform
199 .to_human_string(&field_types[i].field_type, value),
200 )
201 .finish()
202 })
203 .join("/")
204 }
205}
206
207#[derive(Clone, Debug)]
210pub struct PartitionKey {
211 spec: PartitionSpec,
213 schema: SchemaRef,
215 data: Struct,
217}
218
219impl PartitionKey {
220 pub fn new(spec: PartitionSpec, schema: SchemaRef, data: Struct) -> Self {
222 Self { spec, schema, data }
223 }
224
225 pub fn copy_with_data(&self, data: Struct) -> Self {
227 Self {
228 spec: self.spec.clone(),
229 schema: self.schema.clone(),
230 data,
231 }
232 }
233
234 pub fn to_path(&self) -> String {
236 self.spec.partition_to_path(&self.data, self.schema.clone())
237 }
238
239 pub fn is_effectively_none(partition_key: Option<&PartitionKey>) -> bool {
242 match partition_key {
243 None => true,
244 Some(pk) => pk.spec.is_unpartitioned(),
245 }
246 }
247
248 pub fn spec(&self) -> &PartitionSpec {
250 &self.spec
251 }
252
253 pub fn schema(&self) -> &SchemaRef {
255 &self.schema
256 }
257
258 pub fn data(&self) -> &Struct {
260 &self.data
261 }
262}
263
264pub type UnboundPartitionSpecRef = Arc<UnboundPartitionSpec>;
266#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone, TypedBuilder)]
272#[serde(
273 try_from = "self::_serde::UnboundPartitionFieldSerde",
274 into = "self::_serde::UnboundPartitionFieldSerde"
275)]
276#[builder(build_method(into = Result<UnboundPartitionField>))]
277pub struct UnboundPartitionField {
278 source_ids: Vec<i32>,
281 #[builder(default, setter(strip_option(fallback = field_id_opt)))]
284 field_id: Option<i32>,
285 #[builder(setter(into))]
287 name: String,
288 transform: Transform,
290}
291
292impl UnboundPartitionField {
293 pub fn source_id(&self) -> Result<i32> {
298 match self.source_ids.as_slice() {
299 [source_id] => Ok(*source_id),
300 source_ids => Err(invalid_data!(
301 "Partition field '{}' reads {} source columns and has no single source id",
302 self.name,
303 source_ids.len()
304 )),
305 }
306 }
307
308 pub fn source_ids(&self) -> &[i32] {
310 &self.source_ids
311 }
312
313 pub fn field_id(&self) -> Option<i32> {
315 self.field_id
316 }
317
318 pub fn name(&self) -> &str {
320 &self.name
321 }
322
323 pub fn transform(&self) -> Transform {
325 self.transform
326 }
327
328 pub fn with_field_id(self, field_id: i32) -> Self {
330 Self {
331 field_id: Some(field_id),
332 ..self
333 }
334 }
335
336 fn validate(&self) -> Result<()> {
337 if self.source_ids.is_empty() {
338 return Err(invalid_data!("Empty source-ids is not allowed"));
339 }
340 Ok(())
341 }
342}
343
344impl From<UnboundPartitionField> for Result<UnboundPartitionField> {
345 fn from(field: UnboundPartitionField) -> Self {
346 field.validate()?;
347 Ok(field)
348 }
349}
350
351mod _serde {
352 use serde::{Deserialize, Serialize};
353
354 use super::UnboundPartitionField;
355 use crate::Error;
356 use crate::error::invalid_data;
357 use crate::spec::Transform;
358
359 #[derive(Serialize, Deserialize)]
363 #[serde(rename_all = "kebab-case")]
364 pub(super) struct UnboundPartitionFieldSerde {
365 #[serde(default, skip_serializing_if = "Option::is_none")]
366 source_id: Option<i32>,
367 #[serde(default, skip_serializing_if = "Option::is_none")]
368 source_ids: Option<Vec<i32>>,
369 #[serde(default, skip_serializing_if = "Option::is_none")]
370 field_id: Option<i32>,
371 name: String,
372 transform: Transform,
373 }
374
375 impl TryFrom<UnboundPartitionFieldSerde> for UnboundPartitionField {
376 type Error = Error;
377
378 fn try_from(value: UnboundPartitionFieldSerde) -> Result<Self, Error> {
379 let source_ids = match (value.source_id, value.source_ids) {
380 (Some(source_id), None) => vec![source_id],
381 (None, Some(source_ids)) => source_ids,
382 (Some(_), Some(_)) => {
383 return Err(invalid_data!(
384 "source-id and source-ids are mutually exclusive"
385 ));
386 }
387 (None, None) => {
388 return Err(invalid_data!(
389 "Either `source-id` or `source-ids` must be present"
390 ));
391 }
392 };
393
394 UnboundPartitionField::builder()
395 .source_ids(source_ids)
396 .field_id_opt(value.field_id)
397 .name(value.name)
398 .transform(value.transform)
399 .build()
400 }
401 }
402
403 impl From<UnboundPartitionField> for UnboundPartitionFieldSerde {
404 fn from(value: UnboundPartitionField) -> Self {
405 let (source_id, source_ids) = match value.source_ids.as_slice() {
406 [source_id] => (Some(*source_id), None),
407 _ => (None, Some(value.source_ids)),
408 };
409 Self {
410 source_id,
411 source_ids,
412 field_id: value.field_id,
413 name: value.name,
414 transform: value.transform,
415 }
416 }
417 }
418}
419
420#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone, Default)]
424#[serde(rename_all = "kebab-case")]
425pub struct UnboundPartitionSpec {
426 #[serde(skip_serializing_if = "Option::is_none")]
428 pub(crate) spec_id: Option<i32>,
429 pub(crate) fields: Vec<UnboundPartitionField>,
431}
432
433impl UnboundPartitionSpec {
434 pub fn builder() -> UnboundPartitionSpecBuilder {
436 UnboundPartitionSpecBuilder::default()
437 }
438
439 pub fn bind(self, schema: impl Into<SchemaRef>) -> Result<PartitionSpec> {
441 PartitionSpecBuilder::new_from_unbound(self, schema)?.build()
442 }
443
444 pub fn spec_id(&self) -> Option<i32> {
446 self.spec_id
447 }
448
449 pub fn fields(&self) -> &[UnboundPartitionField] {
451 &self.fields
452 }
453
454 pub fn with_spec_id(self, spec_id: i32) -> Self {
456 Self {
457 spec_id: Some(spec_id),
458 ..self
459 }
460 }
461}
462
463fn has_sequential_ids(field_ids: impl Iterator<Item = i32>) -> bool {
464 for (index, field_id) in field_ids.enumerate() {
465 let expected_id = (UNPARTITIONED_LAST_ASSIGNED_ID as i64)
466 .checked_add(1)
467 .and_then(|id| id.checked_add(index as i64))
468 .unwrap_or(i64::MAX);
469
470 if field_id as i64 != expected_id {
471 return false;
472 }
473 }
474
475 true
476}
477
478impl From<PartitionField> for UnboundPartitionField {
479 fn from(field: PartitionField) -> Self {
480 UnboundPartitionField {
481 source_ids: vec![field.source_id],
482 field_id: Some(field.field_id),
483 name: field.name,
484 transform: field.transform,
485 }
486 }
487}
488
489impl From<PartitionSpec> for UnboundPartitionSpec {
490 fn from(spec: PartitionSpec) -> Self {
491 UnboundPartitionSpec {
492 spec_id: Some(spec.spec_id),
493 fields: spec.fields.into_iter().map(Into::into).collect(),
494 }
495 }
496}
497
498#[derive(Debug, Default)]
500pub struct UnboundPartitionSpecBuilder {
501 spec_id: Option<i32>,
502 fields: Vec<UnboundPartitionField>,
503}
504
505impl UnboundPartitionSpecBuilder {
506 pub fn new() -> Self {
508 Self {
509 spec_id: None,
510 fields: vec![],
511 }
512 }
513
514 pub fn with_spec_id(mut self, spec_id: i32) -> Self {
516 self.spec_id = Some(spec_id);
517 self
518 }
519
520 pub fn add_partition_field(mut self, field: UnboundPartitionField) -> Result<Self> {
522 self.check_name_set_and_unique(&field.name)?;
523 self.check_for_redundant_partitions(&field.source_ids, &field.transform)?;
524 if let Some(partition_field_id) = field.field_id {
525 self.check_partition_id_unique(partition_field_id)?;
526 }
527 self.fields.push(field);
528 Ok(self)
529 }
530
531 pub fn add_partition_fields(
533 self,
534 fields: impl IntoIterator<Item = UnboundPartitionField>,
535 ) -> Result<Self> {
536 let mut builder = self;
537 for field in fields {
538 builder = builder.add_partition_field(field)?;
539 }
540 Ok(builder)
541 }
542
543 pub fn build(self) -> UnboundPartitionSpec {
545 UnboundPartitionSpec {
546 spec_id: self.spec_id,
547 fields: self.fields,
548 }
549 }
550}
551
552#[derive(Debug)]
554pub struct PartitionSpecBuilder {
555 spec_id: Option<i32>,
556 last_assigned_field_id: i32,
557 fields: Vec<UnboundPartitionField>,
558 schema: SchemaRef,
559}
560
561impl PartitionSpecBuilder {
562 pub fn new(schema: impl Into<SchemaRef>) -> Self {
564 Self {
565 spec_id: None,
566 fields: vec![],
567 last_assigned_field_id: UNPARTITIONED_LAST_ASSIGNED_ID,
568 schema: schema.into(),
569 }
570 }
571
572 pub fn new_from_unbound(
574 unbound: UnboundPartitionSpec,
575 schema: impl Into<SchemaRef>,
576 ) -> Result<Self> {
577 let mut builder =
578 Self::new(schema).with_spec_id(unbound.spec_id.unwrap_or(DEFAULT_PARTITION_SPEC_ID));
579
580 for field in unbound.fields {
581 builder = builder.add_unbound_field(field)?;
582 }
583 Ok(builder)
584 }
585
586 pub fn with_last_assigned_field_id(mut self, last_assigned_field_id: i32) -> Self {
592 self.last_assigned_field_id = last_assigned_field_id;
593 self
594 }
595
596 pub fn with_spec_id(mut self, spec_id: i32) -> Self {
598 self.spec_id = Some(spec_id);
599 self
600 }
601
602 pub fn add_partition_field(
604 self,
605 source_name: impl AsRef<str>,
606 target_name: impl Into<String>,
607 transform: Transform,
608 ) -> Result<Self> {
609 let source_id = self
610 .schema
611 .field_by_name(source_name.as_ref())
612 .ok_or_else(|| {
613 invalid_data!(
614 "Cannot find source column with name: {} in schema",
615 source_name.as_ref()
616 )
617 })?
618 .id;
619 let field = UnboundPartitionField {
620 source_ids: vec![source_id],
621 field_id: None,
622 name: target_name.into(),
623 transform,
624 };
625
626 self.add_unbound_field(field)
627 }
628
629 pub fn add_unbound_field(mut self, field: UnboundPartitionField) -> Result<Self> {
634 self.check_name_set_and_unique(&field.name)?;
635 self.check_for_redundant_partitions(&field.source_ids, &field.transform)?;
636 Self::check_name_does_not_collide_with_schema(&field, &self.schema)?;
637 Self::check_transform_compatibility(&field, &self.schema)?;
638 if let Some(partition_field_id) = field.field_id {
639 self.check_partition_id_unique(partition_field_id)?;
640 }
641
642 self.fields.push(field);
644 Ok(self)
645 }
646
647 pub fn add_unbound_fields(
649 self,
650 fields: impl IntoIterator<Item = UnboundPartitionField>,
651 ) -> Result<Self> {
652 let mut builder = self;
653 for field in fields {
654 builder = builder.add_unbound_field(field)?;
655 }
656 Ok(builder)
657 }
658
659 pub fn build(self) -> Result<PartitionSpec> {
661 let fields = Self::set_field_ids(self.fields, self.last_assigned_field_id)?;
662 Ok(PartitionSpec {
663 spec_id: self.spec_id.unwrap_or(DEFAULT_PARTITION_SPEC_ID),
664 fields,
665 })
666 }
667
668 fn set_field_ids(
669 fields: Vec<UnboundPartitionField>,
670 last_assigned_field_id: i32,
671 ) -> Result<Vec<PartitionField>> {
672 let mut last_assigned_field_id = last_assigned_field_id;
673 let assigned_ids = fields
676 .iter()
677 .filter_map(|f| f.field_id)
678 .collect::<std::collections::HashSet<_>>();
679
680 fn _check_add_1(prev: i32) -> Result<i32> {
681 prev.checked_add(1)
682 .ok_or_else(|| invalid_data!("Cannot assign more partition ids. Overflow."))
683 }
684
685 let mut bound_fields = Vec::with_capacity(fields.len());
686 for field in fields.into_iter() {
687 let partition_field_id = if let Some(partition_field_id) = field.field_id {
688 last_assigned_field_id = std::cmp::max(last_assigned_field_id, partition_field_id);
689 partition_field_id
690 } else {
691 last_assigned_field_id = _check_add_1(last_assigned_field_id)?;
692 while assigned_ids.contains(&last_assigned_field_id) {
693 last_assigned_field_id = _check_add_1(last_assigned_field_id)?;
694 }
695 last_assigned_field_id
696 };
697
698 bound_fields.push(PartitionField {
699 source_id: field.source_id()?,
700 field_id: partition_field_id,
701 name: field.name,
702 transform: field.transform,
703 })
704 }
705
706 Ok(bound_fields)
707 }
708
709 fn check_name_does_not_collide_with_schema(
714 field: &UnboundPartitionField,
715 schema: &Schema,
716 ) -> Result<()> {
717 match schema.field_by_name(field.name.as_str()) {
718 Some(schema_collision) => {
719 if field.transform == Transform::Identity {
720 if schema_collision.id == field.source_id()? {
721 Ok(())
722 } else {
723 Err(invalid_data!(
724 "Cannot create identity partition sourced from different field in schema. Field name '{}' has id `{}` in schema but partition source id is `{}`",
725 field.name,
726 schema_collision.id,
727 field.source_id()?
728 ))
729 }
730 } else {
731 Err(invalid_data!(
732 "Cannot create partition with name: '{}' that conflicts with schema field and is not an identity transform.",
733 field.name
734 ))
735 }
736 }
737 None => Ok(()),
738 }
739 }
740
741 fn check_transform_compatibility(field: &UnboundPartitionField, schema: &Schema) -> Result<()> {
744 let source_id = field.source_id()?;
747 let schema_field = schema.field_by_id(source_id).ok_or_else(|| {
748 invalid_data!("Cannot find partition source field with id `{source_id}` in schema")
749 })?;
750
751 if field.transform != Transform::Void {
752 if !schema_field.field_type.is_primitive() {
753 return Err(invalid_data!(
754 "Cannot partition by non-primitive source field: '{}'.",
755 schema_field.field_type
756 ));
757 }
758
759 if field
760 .transform
761 .result_type(&schema_field.field_type)
762 .is_err()
763 {
764 return Err(invalid_data!(
765 "Invalid source type: '{}' for transform: '{}'.",
766 schema_field.field_type,
767 field.transform.dedup_name()
768 ));
769 }
770 }
771
772 Ok(())
773 }
774}
775
776trait CorePartitionSpecValidator {
778 fn check_name_set_and_unique(&self, name: &str) -> Result<()> {
780 if name.is_empty() {
781 return Err(invalid_data!("Cannot use empty partition name"));
782 }
783
784 if self.fields().iter().any(|f| f.name == name) {
785 return Err(invalid_data!(
786 "Cannot use partition name more than once: {name}"
787 ));
788 }
789 Ok(())
790 }
791
792 fn check_for_redundant_partitions(
794 &self,
795 source_ids: &[i32],
796 transform: &Transform,
797 ) -> Result<()> {
798 let collision = self.fields().iter().find(|f| {
799 f.source_ids == source_ids && f.transform.dedup_name() == transform.dedup_name()
800 });
801
802 if let Some(collision) = collision {
803 Err(invalid_data!(
804 "Cannot add redundant partition with source ids `{:?}` and transform `{}`. A partition with the same source ids and transform already exists with name `{}`",
805 source_ids,
806 transform.dedup_name(),
807 collision.name
808 ))
809 } else {
810 Ok(())
811 }
812 }
813
814 fn check_partition_id_unique(&self, field_id: i32) -> Result<()> {
816 if self.fields().iter().any(|f| f.field_id == Some(field_id)) {
817 return Err(invalid_data!(
818 "Cannot use field id more than once in one PartitionSpec: {field_id}"
819 ));
820 }
821
822 Ok(())
823 }
824
825 fn fields(&self) -> &Vec<UnboundPartitionField>;
826}
827
828impl CorePartitionSpecValidator for PartitionSpecBuilder {
829 fn fields(&self) -> &Vec<UnboundPartitionField> {
830 &self.fields
831 }
832}
833
834impl CorePartitionSpecValidator for UnboundPartitionSpecBuilder {
835 fn fields(&self) -> &Vec<UnboundPartitionField> {
836 &self.fields
837 }
838}
839
840#[cfg(test)]
841mod tests {
842 use super::*;
843 use crate::ErrorKind;
844 use crate::spec::{Literal, PrimitiveType, Type};
845
846 #[test]
847 fn test_partition_spec() {
848 let spec = r#"
849 {
850 "spec-id": 1,
851 "fields": [ {
852 "source-id": 4,
853 "field-id": 1000,
854 "name": "ts_day",
855 "transform": "day"
856 }, {
857 "source-id": 1,
858 "field-id": 1001,
859 "name": "id_bucket",
860 "transform": "bucket[16]"
861 }, {
862 "source-id": 2,
863 "field-id": 1002,
864 "name": "id_truncate",
865 "transform": "truncate[4]"
866 } ]
867 }
868 "#;
869
870 let partition_spec: PartitionSpec = serde_json::from_str(spec).unwrap();
871 assert_eq!(4, partition_spec.fields[0].source_id);
872 assert_eq!(1000, partition_spec.fields[0].field_id);
873 assert_eq!("ts_day", partition_spec.fields[0].name);
874 assert_eq!(Transform::Day, partition_spec.fields[0].transform);
875
876 assert_eq!(1, partition_spec.fields[1].source_id);
877 assert_eq!(1001, partition_spec.fields[1].field_id);
878 assert_eq!("id_bucket", partition_spec.fields[1].name);
879 assert_eq!(Transform::Bucket(16), partition_spec.fields[1].transform);
880
881 assert_eq!(2, partition_spec.fields[2].source_id);
882 assert_eq!(1002, partition_spec.fields[2].field_id);
883 assert_eq!("id_truncate", partition_spec.fields[2].name);
884 assert_eq!(Transform::Truncate(4), partition_spec.fields[2].transform);
885 }
886
887 #[test]
888 fn test_is_unpartitioned() {
889 let schema = Schema::builder()
890 .with_fields(vec![
891 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
892 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
893 ])
894 .build()
895 .unwrap();
896 let partition_spec = PartitionSpec::builder(schema.clone())
897 .with_spec_id(1)
898 .build()
899 .unwrap();
900 assert!(
901 partition_spec.is_unpartitioned(),
902 "Empty partition spec should be unpartitioned"
903 );
904
905 let partition_spec = PartitionSpec::builder(schema.clone())
906 .add_unbound_fields(vec![
907 UnboundPartitionField::builder()
908 .source_ids(vec![1])
909 .name("id".to_string())
910 .transform(Transform::Identity)
911 .build()
912 .unwrap(),
913 UnboundPartitionField::builder()
914 .source_ids(vec![2])
915 .name("name_string".to_string())
916 .transform(Transform::Void)
917 .build()
918 .unwrap(),
919 ])
920 .unwrap()
921 .with_spec_id(1)
922 .build()
923 .unwrap();
924 assert!(
925 !partition_spec.is_unpartitioned(),
926 "Partition spec with one non void transform should not be unpartitioned"
927 );
928
929 let partition_spec = PartitionSpec::builder(schema.clone())
930 .with_spec_id(1)
931 .add_unbound_fields(vec![
932 UnboundPartitionField::builder()
933 .source_ids(vec![1])
934 .name("id_void".to_string())
935 .transform(Transform::Void)
936 .build()
937 .unwrap(),
938 UnboundPartitionField::builder()
939 .source_ids(vec![2])
940 .name("name_void".to_string())
941 .transform(Transform::Void)
942 .build()
943 .unwrap(),
944 ])
945 .unwrap()
946 .build()
947 .unwrap();
948 assert!(
949 partition_spec.is_unpartitioned(),
950 "Partition spec with all void field should be unpartitioned"
951 );
952 }
953
954 #[test]
955 fn test_unbound_partition_spec() {
956 let spec = r#"
957 {
958 "spec-id": 1,
959 "fields": [ {
960 "source-id": 4,
961 "field-id": 1000,
962 "name": "ts_day",
963 "transform": "day"
964 }, {
965 "source-id": 1,
966 "field-id": 1001,
967 "name": "id_bucket",
968 "transform": "bucket[16]"
969 }, {
970 "source-id": 2,
971 "field-id": 1002,
972 "name": "id_truncate",
973 "transform": "truncate[4]"
974 } ]
975 }
976 "#;
977
978 let partition_spec: UnboundPartitionSpec = serde_json::from_str(spec).unwrap();
979 assert_eq!(Some(1), partition_spec.spec_id);
980
981 assert_eq!(4, partition_spec.fields[0].source_id().unwrap());
982 assert_eq!(Some(1000), partition_spec.fields[0].field_id);
983 assert_eq!("ts_day", partition_spec.fields[0].name);
984 assert_eq!(Transform::Day, partition_spec.fields[0].transform);
985
986 assert_eq!(1, partition_spec.fields[1].source_id().unwrap());
987 assert_eq!(Some(1001), partition_spec.fields[1].field_id);
988 assert_eq!("id_bucket", partition_spec.fields[1].name);
989 assert_eq!(Transform::Bucket(16), partition_spec.fields[1].transform);
990
991 assert_eq!(2, partition_spec.fields[2].source_id().unwrap());
992 assert_eq!(Some(1002), partition_spec.fields[2].field_id);
993 assert_eq!("id_truncate", partition_spec.fields[2].name);
994 assert_eq!(Transform::Truncate(4), partition_spec.fields[2].transform);
995
996 let spec = r#"
997 {
998 "fields": [ {
999 "source-id": 4,
1000 "name": "ts_day",
1001 "transform": "day"
1002 } ]
1003 }
1004 "#;
1005 let partition_spec: UnboundPartitionSpec = serde_json::from_str(spec).unwrap();
1006 assert_eq!(None, partition_spec.spec_id);
1007
1008 assert_eq!(4, partition_spec.fields[0].source_id().unwrap());
1009 assert_eq!(None, partition_spec.fields[0].field_id);
1010 assert_eq!("ts_day", partition_spec.fields[0].name);
1011 assert_eq!(Transform::Day, partition_spec.fields[0].transform);
1012 }
1013
1014 #[test]
1015 fn test_unbound_partition_spec_serialization_skips_none_fields() {
1016 let spec = UnboundPartitionSpec::builder()
1017 .add_partition_field(
1018 UnboundPartitionField::builder()
1019 .source_ids(vec![4])
1020 .name("ts_day")
1021 .transform(Transform::Day)
1022 .build()
1023 .unwrap(),
1024 )
1025 .unwrap()
1026 .build();
1027
1028 let value = serde_json::to_value(&spec).unwrap();
1029 let object = value.as_object().unwrap();
1030 assert!(!object.contains_key("spec-id"));
1031 let field = object["fields"][0].as_object().unwrap();
1032 assert!(!field.contains_key("field-id"));
1033
1034 let value = serde_json::to_value(spec.with_spec_id(1)).unwrap();
1035 let object = value.as_object().unwrap();
1036 assert_eq!(Some(&serde_json::json!(1)), object.get("spec-id"));
1037
1038 let spec: UnboundPartitionSpec = serde_json::from_str(
1041 r#"{
1042 "spec-id": null,
1043 "fields": [
1044 {"source-id": 4, "name": "ts_day", "transform": "day", "field-id": null}
1045 ]
1046 }"#,
1047 )
1048 .unwrap();
1049 assert_eq!(None, spec.spec_id);
1050 assert_eq!(None, spec.fields[0].field_id);
1051 }
1052
1053 #[test]
1054 fn test_new_unpartition() {
1055 let schema = Schema::builder()
1056 .with_fields(vec![
1057 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1058 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1059 ])
1060 .build()
1061 .unwrap();
1062 let partition_spec = PartitionSpec::builder(schema.clone())
1063 .with_spec_id(0)
1064 .build()
1065 .unwrap();
1066 let partition_type = partition_spec.partition_type(&schema).unwrap();
1067 assert_eq!(0, partition_type.fields().len());
1068
1069 let unpartition_spec = PartitionSpec::unpartition_spec();
1070 assert_eq!(partition_spec, unpartition_spec);
1071 }
1072
1073 #[test]
1074 fn test_partition_type() {
1075 let spec = r#"
1076 {
1077 "spec-id": 1,
1078 "fields": [ {
1079 "source-id": 4,
1080 "field-id": 1000,
1081 "name": "ts_day",
1082 "transform": "day"
1083 }, {
1084 "source-id": 1,
1085 "field-id": 1001,
1086 "name": "id_bucket",
1087 "transform": "bucket[16]"
1088 }, {
1089 "source-id": 2,
1090 "field-id": 1002,
1091 "name": "id_truncate",
1092 "transform": "truncate[4]"
1093 } ]
1094 }
1095 "#;
1096
1097 let partition_spec: PartitionSpec = serde_json::from_str(spec).unwrap();
1098 let schema = Schema::builder()
1099 .with_fields(vec![
1100 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1101 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1102 NestedField::required(3, "ts", Type::Primitive(PrimitiveType::Timestamp)).into(),
1103 NestedField::required(4, "ts_day", Type::Primitive(PrimitiveType::Timestamp))
1104 .into(),
1105 NestedField::required(5, "id_bucket", Type::Primitive(PrimitiveType::Int)).into(),
1106 NestedField::required(6, "id_truncate", Type::Primitive(PrimitiveType::Int)).into(),
1107 ])
1108 .build()
1109 .unwrap();
1110
1111 let partition_type = partition_spec.partition_type(&schema).unwrap();
1112 assert_eq!(3, partition_type.fields().len());
1113 assert_eq!(
1114 *partition_type.fields()[0],
1115 NestedField::optional(
1116 partition_spec.fields[0].field_id,
1117 &partition_spec.fields[0].name,
1118 Type::Primitive(PrimitiveType::Date)
1119 )
1120 );
1121 assert_eq!(
1122 *partition_type.fields()[1],
1123 NestedField::optional(
1124 partition_spec.fields[1].field_id,
1125 &partition_spec.fields[1].name,
1126 Type::Primitive(PrimitiveType::Int)
1127 )
1128 );
1129 assert_eq!(
1130 *partition_type.fields()[2],
1131 NestedField::optional(
1132 partition_spec.fields[2].field_id,
1133 &partition_spec.fields[2].name,
1134 Type::Primitive(PrimitiveType::String)
1135 )
1136 );
1137 }
1138
1139 #[test]
1140 fn test_partition_empty() {
1141 let spec = r#"
1142 {
1143 "spec-id": 1,
1144 "fields": []
1145 }
1146 "#;
1147
1148 let partition_spec: PartitionSpec = serde_json::from_str(spec).unwrap();
1149 let schema = Schema::builder()
1150 .with_fields(vec![
1151 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1152 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1153 NestedField::required(3, "ts", Type::Primitive(PrimitiveType::Timestamp)).into(),
1154 NestedField::required(4, "ts_day", Type::Primitive(PrimitiveType::Timestamp))
1155 .into(),
1156 NestedField::required(5, "id_bucket", Type::Primitive(PrimitiveType::Int)).into(),
1157 NestedField::required(6, "id_truncate", Type::Primitive(PrimitiveType::Int)).into(),
1158 ])
1159 .build()
1160 .unwrap();
1161
1162 let partition_type = partition_spec.partition_type(&schema).unwrap();
1163 assert_eq!(0, partition_type.fields().len());
1164 }
1165
1166 #[test]
1167 fn test_partition_type_with_dropped_source_column() {
1168 let spec = r#"
1169 {
1170 "spec-id": 1,
1171 "fields": [ {
1172 "source-id": 4,
1173 "field-id": 1000,
1174 "name": "ts_day",
1175 "transform": "day"
1176 }, {
1177 "source-id": 1,
1178 "field-id": 1001,
1179 "name": "id_bucket",
1180 "transform": "bucket[16]"
1181 }, {
1182 "source-id": 2,
1183 "field-id": 1002,
1184 "name": "id_truncate",
1185 "transform": "truncate[4]"
1186 } ]
1187 }
1188 "#;
1189
1190 let partition_spec: PartitionSpec = serde_json::from_str(spec).unwrap();
1191 let schema = Schema::builder()
1192 .with_fields(vec![
1193 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1194 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1195 ])
1196 .build()
1197 .unwrap();
1198
1199 assert_eq!(
1200 partition_spec.partition_type(&schema).unwrap(),
1201 StructType::new(vec![
1202 NestedField::optional(1000, "ts_day", PrimitiveType::Date.into()).into(),
1203 NestedField::optional(1001, "id_bucket", PrimitiveType::Int.into()).into(),
1204 NestedField::optional(1002, "id_truncate", PrimitiveType::String.into()).into(),
1205 ])
1206 );
1207
1208 let empty_schema = Schema::builder().build().unwrap();
1210 assert_eq!(
1211 partition_spec.partition_type(&empty_schema).unwrap(),
1212 StructType::new(vec![
1213 NestedField::optional(1000, "ts_day", PrimitiveType::Date.into()).into(),
1214 NestedField::optional(1001, "id_bucket", PrimitiveType::Int.into()).into(),
1215 NestedField::optional(1002, "id_truncate", PrimitiveType::Unknown.into()).into(),
1216 ])
1217 );
1218
1219 let incompatible_schema = Schema::builder()
1221 .with_fields(vec![
1222 NestedField::required(1, "id", PrimitiveType::Boolean.into()).into(),
1223 ])
1224 .build()
1225 .unwrap();
1226 let err = partition_spec
1227 .partition_type(&incompatible_schema)
1228 .unwrap_err();
1229 assert_eq!(err.kind(), ErrorKind::DataInvalid);
1230 assert_eq!(
1231 err.message(),
1232 "boolean is not a valid input type of bucket transform"
1233 );
1234 }
1235
1236 #[test]
1237 fn test_partition_type_without_source_type() {
1238 let schema = Schema::builder().build().unwrap();
1239 for (transform, expected_type) in [
1240 (Transform::Identity, PrimitiveType::Unknown),
1241 (Transform::Truncate(4), PrimitiveType::Unknown),
1242 (Transform::Void, PrimitiveType::Unknown),
1243 (Transform::Bucket(16), PrimitiveType::Int),
1244 (Transform::Year, PrimitiveType::Int),
1245 (Transform::Month, PrimitiveType::Int),
1246 (Transform::Day, PrimitiveType::Date),
1247 (Transform::Hour, PrimitiveType::Int),
1248 (Transform::Unknown, PrimitiveType::String),
1249 ] {
1250 let spec = PartitionSpec {
1251 spec_id: 0,
1252 fields: vec![PartitionField {
1253 source_id: 1,
1254 field_id: 1000,
1255 name: "partition".to_string(),
1256 transform,
1257 }],
1258 };
1259 assert_eq!(
1260 spec.partition_type(&schema).unwrap(),
1261 StructType::new(vec![
1262 NestedField::optional(1000, "partition", expected_type.into()).into(),
1263 ])
1264 );
1265 }
1266 }
1267
1268 #[test]
1269 fn test_builder_disallow_duplicate_names() {
1270 UnboundPartitionSpec::builder()
1271 .add_partition_field(
1272 UnboundPartitionField::builder()
1273 .source_ids(vec![1])
1274 .name("ts_day")
1275 .transform(Transform::Day)
1276 .build()
1277 .unwrap(),
1278 )
1279 .unwrap()
1280 .add_partition_field(
1281 UnboundPartitionField::builder()
1282 .source_ids(vec![2])
1283 .name("ts_day")
1284 .transform(Transform::Day)
1285 .build()
1286 .unwrap(),
1287 )
1288 .unwrap_err();
1289 }
1290
1291 #[test]
1292 fn test_builder_disallow_duplicate_field_ids() {
1293 let schema = Schema::builder()
1294 .with_fields(vec![
1295 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1296 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1297 ])
1298 .build()
1299 .unwrap();
1300 PartitionSpec::builder(schema.clone())
1301 .add_unbound_field(UnboundPartitionField {
1302 source_ids: vec![1],
1303 field_id: Some(1000),
1304 name: "id".to_string(),
1305 transform: Transform::Identity,
1306 })
1307 .unwrap()
1308 .add_unbound_field(UnboundPartitionField {
1309 source_ids: vec![2],
1310 field_id: Some(1000),
1311 name: "id_bucket".to_string(),
1312 transform: Transform::Bucket(16),
1313 })
1314 .unwrap_err();
1315 }
1316
1317 #[test]
1318 fn test_builder_auto_assign_field_ids() {
1319 let schema = Schema::builder()
1320 .with_fields(vec![
1321 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1322 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1323 NestedField::required(3, "ts", Type::Primitive(PrimitiveType::Timestamp)).into(),
1324 ])
1325 .build()
1326 .unwrap();
1327 let spec = PartitionSpec::builder(schema.clone())
1328 .with_spec_id(1)
1329 .add_unbound_field(UnboundPartitionField {
1330 source_ids: vec![1],
1331 name: "id".to_string(),
1332 transform: Transform::Identity,
1333 field_id: Some(1012),
1334 })
1335 .unwrap()
1336 .add_unbound_field(UnboundPartitionField {
1337 source_ids: vec![2],
1338 name: "name_void".to_string(),
1339 transform: Transform::Void,
1340 field_id: None,
1341 })
1342 .unwrap()
1343 .add_unbound_field(UnboundPartitionField {
1345 source_ids: vec![3],
1346 name: "year".to_string(),
1347 transform: Transform::Year,
1348 field_id: Some(1),
1349 })
1350 .unwrap()
1351 .build()
1352 .unwrap();
1353
1354 assert_eq!(1012, spec.fields[0].field_id);
1355 assert_eq!(1013, spec.fields[1].field_id);
1356 assert_eq!(1, spec.fields[2].field_id);
1357 }
1358
1359 #[test]
1360 fn test_builder_valid_schema() {
1361 let schema = Schema::builder()
1362 .with_fields(vec![
1363 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1364 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1365 ])
1366 .build()
1367 .unwrap();
1368
1369 PartitionSpec::builder(schema.clone())
1370 .with_spec_id(1)
1371 .build()
1372 .unwrap();
1373
1374 let spec = PartitionSpec::builder(schema.clone())
1375 .with_spec_id(1)
1376 .add_partition_field("id", "id_bucket[16]", Transform::Bucket(16))
1377 .unwrap()
1378 .build()
1379 .unwrap();
1380
1381 assert_eq!(spec, PartitionSpec {
1382 spec_id: 1,
1383 fields: vec![PartitionField {
1384 source_id: 1,
1385 field_id: 1000,
1386 name: "id_bucket[16]".to_string(),
1387 transform: Transform::Bucket(16),
1388 }],
1389 });
1390 assert_eq!(
1391 spec.partition_type(&schema).unwrap(),
1392 StructType::new(vec![
1393 NestedField::optional(1000, "id_bucket[16]", Type::Primitive(PrimitiveType::Int))
1394 .into()
1395 ])
1396 )
1397 }
1398
1399 #[test]
1400 fn test_collision_with_schema_name() {
1401 let schema = Schema::builder()
1402 .with_fields(vec![
1403 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1404 ])
1405 .build()
1406 .unwrap();
1407
1408 PartitionSpec::builder(schema.clone())
1409 .with_spec_id(1)
1410 .build()
1411 .unwrap();
1412
1413 let err = PartitionSpec::builder(schema)
1414 .with_spec_id(1)
1415 .add_unbound_field(UnboundPartitionField {
1416 source_ids: vec![1],
1417 field_id: None,
1418 name: "id".to_string(),
1419 transform: Transform::Bucket(16),
1420 })
1421 .unwrap_err();
1422 assert!(err.message().contains("conflicts with schema"))
1423 }
1424
1425 #[test]
1426 fn test_builder_collision_is_ok_for_identity_transforms() {
1427 let schema = Schema::builder()
1428 .with_fields(vec![
1429 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1430 NestedField::required(2, "number", Type::Primitive(PrimitiveType::Int)).into(),
1431 ])
1432 .build()
1433 .unwrap();
1434
1435 PartitionSpec::builder(schema.clone())
1436 .with_spec_id(1)
1437 .build()
1438 .unwrap();
1439
1440 PartitionSpec::builder(schema.clone())
1441 .with_spec_id(1)
1442 .add_unbound_field(UnboundPartitionField {
1443 source_ids: vec![1],
1444 field_id: None,
1445 name: "id".to_string(),
1446 transform: Transform::Identity,
1447 })
1448 .unwrap()
1449 .build()
1450 .unwrap();
1451
1452 PartitionSpec::builder(schema)
1454 .with_spec_id(1)
1455 .add_unbound_field(UnboundPartitionField {
1456 source_ids: vec![2],
1457 field_id: None,
1458 name: "id".to_string(),
1459 transform: Transform::Identity,
1460 })
1461 .unwrap_err();
1462 }
1463
1464 #[test]
1465 fn test_builder_all_source_ids_must_exist() {
1466 let schema = Schema::builder()
1467 .with_fields(vec![
1468 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1469 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1470 NestedField::required(3, "ts", Type::Primitive(PrimitiveType::Timestamp)).into(),
1471 ])
1472 .build()
1473 .unwrap();
1474
1475 PartitionSpec::builder(schema.clone())
1477 .with_spec_id(1)
1478 .add_unbound_fields(vec![
1479 UnboundPartitionField {
1480 source_ids: vec![1],
1481 field_id: None,
1482 name: "id_bucket".to_string(),
1483 transform: Transform::Bucket(16),
1484 },
1485 UnboundPartitionField {
1486 source_ids: vec![2],
1487 field_id: None,
1488 name: "name".to_string(),
1489 transform: Transform::Identity,
1490 },
1491 ])
1492 .unwrap()
1493 .build()
1494 .unwrap();
1495
1496 PartitionSpec::builder(schema)
1498 .with_spec_id(1)
1499 .add_unbound_fields(vec![
1500 UnboundPartitionField {
1501 source_ids: vec![1],
1502 field_id: None,
1503 name: "id_bucket".to_string(),
1504 transform: Transform::Bucket(16),
1505 },
1506 UnboundPartitionField {
1507 source_ids: vec![4],
1508 field_id: None,
1509 name: "name".to_string(),
1510 transform: Transform::Identity,
1511 },
1512 ])
1513 .unwrap_err();
1514 }
1515
1516 #[test]
1517 fn test_builder_disallows_variant_source() {
1518 let schema = Schema::builder()
1519 .with_fields(vec![
1520 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1521 NestedField::optional(2, "v", Type::Variant(crate::spec::VariantType)).into(),
1522 ])
1523 .build()
1524 .unwrap();
1525
1526 let err = PartitionSpec::builder(schema)
1527 .with_spec_id(1)
1528 .add_unbound_fields(vec![UnboundPartitionField {
1529 source_ids: vec![2],
1530 field_id: None,
1531 name: "v_part".to_string(),
1532 transform: Transform::Identity,
1533 }])
1534 .expect_err("variant must not be allowed as a partition source");
1535
1536 assert_eq!(
1537 err.message(),
1538 "Cannot partition by non-primitive source field: 'variant'."
1539 );
1540 }
1541
1542 #[test]
1543 fn test_builder_disallows_redundant() {
1544 let err = UnboundPartitionSpec::builder()
1545 .with_spec_id(1)
1546 .add_partition_field(
1547 UnboundPartitionField::builder()
1548 .source_ids(vec![1])
1549 .name("id_bucket[16]")
1550 .transform(Transform::Bucket(16))
1551 .build()
1552 .unwrap(),
1553 )
1554 .unwrap()
1555 .add_partition_field(
1556 UnboundPartitionField::builder()
1557 .source_ids(vec![1])
1558 .name("id_bucket_with_other_name")
1559 .transform(Transform::Bucket(16))
1560 .build()
1561 .unwrap(),
1562 )
1563 .unwrap_err();
1564 assert!(err.message().contains("redundant partition"));
1565 }
1566
1567 #[test]
1568 fn test_builder_incompatible_transforms_disallowed() {
1569 let schema = Schema::builder()
1570 .with_fields(vec![
1571 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1572 ])
1573 .build()
1574 .unwrap();
1575
1576 PartitionSpec::builder(schema)
1577 .with_spec_id(1)
1578 .add_unbound_field(UnboundPartitionField {
1579 source_ids: vec![1],
1580 field_id: None,
1581 name: "id_year".to_string(),
1582 transform: Transform::Year,
1583 })
1584 .unwrap_err();
1585 }
1586
1587 #[test]
1588 fn test_build_unbound_specs_without_partition_id() {
1589 let spec = UnboundPartitionSpec::builder()
1590 .with_spec_id(1)
1591 .add_partition_fields(vec![UnboundPartitionField {
1592 source_ids: vec![1],
1593 field_id: None,
1594 name: "id_bucket[16]".to_string(),
1595 transform: Transform::Bucket(16),
1596 }])
1597 .unwrap()
1598 .build();
1599
1600 assert_eq!(spec, UnboundPartitionSpec {
1601 spec_id: Some(1),
1602 fields: vec![UnboundPartitionField {
1603 source_ids: vec![1],
1604 field_id: None,
1605 name: "id_bucket[16]".to_string(),
1606 transform: Transform::Bucket(16),
1607 }]
1608 });
1609 }
1610
1611 #[test]
1612 fn test_is_compatible_with() {
1613 let schema = Schema::builder()
1614 .with_fields(vec![
1615 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1616 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1617 ])
1618 .build()
1619 .unwrap();
1620
1621 let partition_spec_1 = PartitionSpec::builder(schema.clone())
1622 .with_spec_id(1)
1623 .add_unbound_field(UnboundPartitionField {
1624 source_ids: vec![1],
1625 field_id: None,
1626 name: "id_bucket".to_string(),
1627 transform: Transform::Bucket(16),
1628 })
1629 .unwrap()
1630 .build()
1631 .unwrap();
1632
1633 let partition_spec_2 = PartitionSpec::builder(schema)
1634 .with_spec_id(1)
1635 .add_unbound_field(UnboundPartitionField {
1636 source_ids: vec![1],
1637 field_id: None,
1638 name: "id_bucket".to_string(),
1639 transform: Transform::Bucket(16),
1640 })
1641 .unwrap()
1642 .build()
1643 .unwrap();
1644
1645 assert!(partition_spec_1.is_compatible_with(&partition_spec_2));
1646 }
1647
1648 #[test]
1649 fn test_not_compatible_with_transform_different() {
1650 let schema = Schema::builder()
1651 .with_fields(vec![
1652 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1653 ])
1654 .build()
1655 .unwrap();
1656
1657 let partition_spec_1 = PartitionSpec::builder(schema.clone())
1658 .with_spec_id(1)
1659 .add_unbound_field(UnboundPartitionField {
1660 source_ids: vec![1],
1661 field_id: None,
1662 name: "id_bucket".to_string(),
1663 transform: Transform::Bucket(16),
1664 })
1665 .unwrap()
1666 .build()
1667 .unwrap();
1668
1669 let partition_spec_2 = PartitionSpec::builder(schema)
1670 .with_spec_id(1)
1671 .add_unbound_field(UnboundPartitionField {
1672 source_ids: vec![1],
1673 field_id: None,
1674 name: "id_bucket".to_string(),
1675 transform: Transform::Bucket(32),
1676 })
1677 .unwrap()
1678 .build()
1679 .unwrap();
1680
1681 assert!(!partition_spec_1.is_compatible_with(&partition_spec_2));
1682 }
1683
1684 #[test]
1685 fn test_not_compatible_with_source_id_different() {
1686 let schema = Schema::builder()
1687 .with_fields(vec![
1688 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1689 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1690 ])
1691 .build()
1692 .unwrap();
1693
1694 let partition_spec_1 = PartitionSpec::builder(schema.clone())
1695 .with_spec_id(1)
1696 .add_unbound_field(UnboundPartitionField {
1697 source_ids: vec![1],
1698 field_id: None,
1699 name: "id_bucket".to_string(),
1700 transform: Transform::Bucket(16),
1701 })
1702 .unwrap()
1703 .build()
1704 .unwrap();
1705
1706 let partition_spec_2 = PartitionSpec::builder(schema)
1707 .with_spec_id(1)
1708 .add_unbound_field(UnboundPartitionField {
1709 source_ids: vec![2],
1710 field_id: None,
1711 name: "id_bucket".to_string(),
1712 transform: Transform::Bucket(16),
1713 })
1714 .unwrap()
1715 .build()
1716 .unwrap();
1717
1718 assert!(!partition_spec_1.is_compatible_with(&partition_spec_2));
1719 }
1720
1721 #[test]
1722 fn test_not_compatible_with_order_different() {
1723 let schema = Schema::builder()
1724 .with_fields(vec![
1725 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1726 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1727 ])
1728 .build()
1729 .unwrap();
1730
1731 let partition_spec_1 = PartitionSpec::builder(schema.clone())
1732 .with_spec_id(1)
1733 .add_unbound_field(UnboundPartitionField {
1734 source_ids: vec![1],
1735 field_id: None,
1736 name: "id_bucket".to_string(),
1737 transform: Transform::Bucket(16),
1738 })
1739 .unwrap()
1740 .add_unbound_field(UnboundPartitionField {
1741 source_ids: vec![2],
1742 field_id: None,
1743 name: "name".to_string(),
1744 transform: Transform::Identity,
1745 })
1746 .unwrap()
1747 .build()
1748 .unwrap();
1749
1750 let partition_spec_2 = PartitionSpec::builder(schema)
1751 .with_spec_id(1)
1752 .add_unbound_field(UnboundPartitionField {
1753 source_ids: vec![2],
1754 field_id: None,
1755 name: "name".to_string(),
1756 transform: Transform::Identity,
1757 })
1758 .unwrap()
1759 .add_unbound_field(UnboundPartitionField {
1760 source_ids: vec![1],
1761 field_id: None,
1762 name: "id_bucket".to_string(),
1763 transform: Transform::Bucket(16),
1764 })
1765 .unwrap()
1766 .build()
1767 .unwrap();
1768
1769 assert!(!partition_spec_1.is_compatible_with(&partition_spec_2));
1770 }
1771
1772 #[test]
1773 fn test_highest_field_id_unpartitioned() {
1774 let spec = PartitionSpec::builder(Schema::builder().with_fields(vec![]).build().unwrap())
1775 .with_spec_id(1)
1776 .build()
1777 .unwrap();
1778
1779 assert!(spec.highest_field_id().is_none());
1780 }
1781
1782 #[test]
1783 fn test_highest_field_id() {
1784 let schema = Schema::builder()
1785 .with_fields(vec![
1786 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1787 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1788 ])
1789 .build()
1790 .unwrap();
1791
1792 let spec = PartitionSpec::builder(schema)
1793 .with_spec_id(1)
1794 .add_unbound_field(UnboundPartitionField {
1795 source_ids: vec![1],
1796 field_id: Some(1001),
1797 name: "id".to_string(),
1798 transform: Transform::Identity,
1799 })
1800 .unwrap()
1801 .add_unbound_field(UnboundPartitionField {
1802 source_ids: vec![2],
1803 field_id: Some(1000),
1804 name: "name".to_string(),
1805 transform: Transform::Identity,
1806 })
1807 .unwrap()
1808 .build()
1809 .unwrap();
1810
1811 assert_eq!(Some(1001), spec.highest_field_id());
1812 }
1813
1814 #[test]
1815 fn test_has_sequential_ids() {
1816 let schema = Schema::builder()
1817 .with_fields(vec![
1818 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1819 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1820 ])
1821 .build()
1822 .unwrap();
1823
1824 let spec = PartitionSpec::builder(schema)
1825 .with_spec_id(1)
1826 .add_unbound_field(UnboundPartitionField {
1827 source_ids: vec![1],
1828 field_id: Some(1000),
1829 name: "id".to_string(),
1830 transform: Transform::Identity,
1831 })
1832 .unwrap()
1833 .add_unbound_field(UnboundPartitionField {
1834 source_ids: vec![2],
1835 field_id: Some(1001),
1836 name: "name".to_string(),
1837 transform: Transform::Identity,
1838 })
1839 .unwrap()
1840 .build()
1841 .unwrap();
1842
1843 assert_eq!(1000, spec.fields[0].field_id);
1844 assert_eq!(1001, spec.fields[1].field_id);
1845 assert!(spec.has_sequential_ids());
1846 }
1847
1848 #[test]
1849 fn test_sequential_ids_must_start_at_1000() {
1850 let schema = Schema::builder()
1851 .with_fields(vec![
1852 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1853 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1854 ])
1855 .build()
1856 .unwrap();
1857
1858 let spec = PartitionSpec::builder(schema)
1859 .with_spec_id(1)
1860 .add_unbound_field(UnboundPartitionField {
1861 source_ids: vec![1],
1862 field_id: Some(999),
1863 name: "id".to_string(),
1864 transform: Transform::Identity,
1865 })
1866 .unwrap()
1867 .add_unbound_field(UnboundPartitionField {
1868 source_ids: vec![2],
1869 field_id: Some(1000),
1870 name: "name".to_string(),
1871 transform: Transform::Identity,
1872 })
1873 .unwrap()
1874 .build()
1875 .unwrap();
1876
1877 assert_eq!(999, spec.fields[0].field_id);
1878 assert_eq!(1000, spec.fields[1].field_id);
1879 assert!(!spec.has_sequential_ids());
1880 }
1881
1882 #[test]
1883 fn test_sequential_ids_must_have_no_gaps() {
1884 let schema = Schema::builder()
1885 .with_fields(vec![
1886 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1887 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1888 ])
1889 .build()
1890 .unwrap();
1891
1892 let spec = PartitionSpec::builder(schema)
1893 .with_spec_id(1)
1894 .add_unbound_field(UnboundPartitionField {
1895 source_ids: vec![1],
1896 field_id: Some(1000),
1897 name: "id".to_string(),
1898 transform: Transform::Identity,
1899 })
1900 .unwrap()
1901 .add_unbound_field(UnboundPartitionField {
1902 source_ids: vec![2],
1903 field_id: Some(1002),
1904 name: "name".to_string(),
1905 transform: Transform::Identity,
1906 })
1907 .unwrap()
1908 .build()
1909 .unwrap();
1910
1911 assert_eq!(1000, spec.fields[0].field_id);
1912 assert_eq!(1002, spec.fields[1].field_id);
1913 assert!(!spec.has_sequential_ids());
1914 }
1915
1916 #[test]
1917 fn test_partition_to_path() {
1918 let schema = Schema::builder()
1919 .with_fields(vec![
1920 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1921 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1922 NestedField::required(3, "timestamp", Type::Primitive(PrimitiveType::Timestamp))
1923 .into(),
1924 NestedField::required(4, "empty", Type::Primitive(PrimitiveType::String)).into(),
1925 ])
1926 .build()
1927 .unwrap();
1928
1929 let spec = PartitionSpec::builder(schema.clone())
1930 .add_partition_field("id", "id", Transform::Identity)
1931 .unwrap()
1932 .add_partition_field("name", "name", Transform::Identity)
1933 .unwrap()
1934 .add_partition_field("timestamp", "ts_hour", Transform::Hour)
1935 .unwrap()
1936 .add_partition_field("empty", "empty_void", Transform::Void)
1937 .unwrap()
1938 .build()
1939 .unwrap();
1940
1941 let data = Struct::from_iter([
1942 Some(Literal::int(42)),
1943 Some(Literal::string("alice")),
1944 Some(Literal::int(1000)),
1945 Some(Literal::string("empty")),
1946 ]);
1947
1948 assert_eq!(
1949 spec.partition_to_path(&data, schema.into()),
1950 "id=42/name=alice/ts_hour=1970-02-11-16/empty_void=null"
1951 );
1952 }
1953
1954 #[test]
1955 fn test_partition_to_path_escaped_strings() {
1956 let schema = Schema::builder()
1957 .with_fields(vec![
1958 NestedField::required(1, "\"esc\"#1", Type::Primitive(PrimitiveType::String))
1959 .into(),
1960 NestedField::required(2, "data", Type::Primitive(PrimitiveType::String)).into(),
1961 ])
1962 .build()
1963 .unwrap();
1964
1965 let spec = PartitionSpec::builder(schema.clone())
1966 .add_partition_field("\"esc\"#1", "\"esc\"#1", Transform::Identity)
1967 .unwrap()
1968 .build()
1969 .unwrap();
1970
1971 let data = Struct::from_iter([
1972 Some(Literal::string("a/b/c/d")),
1973 Some(Literal::string("val#1")),
1974 ]);
1975
1976 assert_eq!(
1977 spec.partition_to_path(&data, schema.into()),
1978 "%22esc%22%231=a%2Fb%2Fc%2Fd"
1979 );
1980 }
1981
1982 #[test]
1983 fn test_partition_to_path_escaped_field_name() {
1984 let schema = Schema::builder()
1985 .with_fields(vec![
1986 NestedField::required(1, "\"esc\"#1", Type::Primitive(PrimitiveType::String))
1987 .into(),
1988 NestedField::required(2, "data", Type::Primitive(PrimitiveType::String)).into(),
1989 ])
1990 .build()
1991 .unwrap();
1992
1993 let spec = PartitionSpec::builder(schema.clone())
1994 .add_partition_field("data", "data", Transform::Identity)
1995 .unwrap()
1996 .add_partition_field("data", "data_truc_10", Transform::Truncate(10))
1997 .unwrap()
1998 .build()
1999 .unwrap();
2000
2001 let data = Struct::from_iter([
2002 Some(Literal::string("a/b/c/d")),
2003 Some(Literal::string("a/b/c/d")),
2004 ]);
2005
2006 assert_eq!(
2007 spec.partition_to_path(&data, schema.into()),
2008 "data=a%2Fb%2Fc%2Fd/data_truc_10=a%2Fb%2Fc%2Fd"
2009 );
2010 }
2011
2012 #[test]
2013 fn test_unbound_partition_field_reads_source_ids_only() {
2014 let field: UnboundPartitionField = serde_json::from_str(
2015 r#"{"source-ids": [1, 2], "name": "m", "transform": "bucket[4]"}"#,
2016 )
2017 .unwrap();
2018
2019 assert_eq!([1, 2], field.source_ids());
2020 assert!(
2022 field
2023 .source_id()
2024 .unwrap_err()
2025 .to_string()
2026 .contains("has no single source id")
2027 );
2028
2029 let serialized = serde_json::to_value(&field).unwrap();
2030 assert_eq!(
2031 Some(&serde_json::json!([1, 2])),
2032 serialized.get("source-ids")
2033 );
2034 assert!(serialized.get("source-id").is_none());
2035 }
2036
2037 #[test]
2038 fn test_unbound_partition_field_single_source_id_round_trip() {
2039 let field: UnboundPartitionField =
2040 serde_json::from_str(r#"{"source-id": 1, "name": "m", "transform": "identity"}"#)
2041 .unwrap();
2042
2043 assert_eq!(1, field.source_id().unwrap());
2044 assert_eq!([1], field.source_ids());
2045
2046 let serialized = serde_json::to_value(&field).unwrap();
2048 assert_eq!(Some(&serde_json::json!(1)), serialized.get("source-id"));
2049 assert!(serialized.get("source-ids").is_none());
2050 }
2051
2052 #[test]
2053 fn test_unbound_partition_field_rejects_malformed_source_ids() {
2054 for (input, expected) in [
2055 (
2056 r#"{"source-ids": [], "name": "m", "transform": "identity"}"#,
2057 "Empty source-ids is not allowed",
2058 ),
2059 (
2060 r#"{"name": "m", "transform": "identity"}"#,
2061 "Either `source-id` or `source-ids` must be present",
2062 ),
2063 (
2064 r#"{"source-id": 1, "source-ids": [1], "name": "m", "transform": "identity"}"#,
2065 "mutually exclusive",
2066 ),
2067 ] {
2068 let err = serde_json::from_str::<UnboundPartitionField>(input).unwrap_err();
2069 assert!(
2070 err.to_string().contains(expected),
2071 "unexpected error for {input}: {err}"
2072 );
2073 }
2074 }
2075
2076 #[test]
2077 fn test_binding_a_multi_argument_field_fails_loudly() {
2078 let schema = Schema::builder()
2079 .with_fields(vec![
2080 NestedField::required(1, "a", Type::Primitive(PrimitiveType::Int)).into(),
2081 NestedField::required(2, "b", Type::Primitive(PrimitiveType::Int)).into(),
2082 ])
2083 .build()
2084 .unwrap();
2085
2086 let spec: UnboundPartitionSpec = serde_json::from_str(
2087 r#"{"spec-id": 1, "fields": [{"source-ids": [1, 2], "field-id": 1000, "name": "m", "transform": "bucket[4]"}]}"#,
2088 )
2089 .unwrap();
2090
2091 let err = spec.bind(schema).unwrap_err();
2094 assert!(
2095 err.to_string().contains("has no single source id"),
2096 "unexpected error: {err}"
2097 );
2098 }
2099
2100 #[test]
2101 fn test_unbound_partition_field_builder_rejects_empty_source_ids() {
2102 let err = UnboundPartitionField::builder()
2103 .source_ids(vec![])
2104 .name("m")
2105 .transform(Transform::Identity)
2106 .build()
2107 .unwrap_err();
2108 assert!(
2109 err.to_string().contains("Empty source-ids is not allowed"),
2110 "unexpected error: {err}"
2111 );
2112 }
2113}