1#![doc(
157 html_logo_url = "https://arrow.apache.org/img/arrow-logo_chevrons_black-txt_white-bg.svg",
158 html_favicon_url = "https://arrow.apache.org/img/arrow-logo_chevrons_black-txt_transparent-bg.svg"
159)]
160#![cfg_attr(docsrs, feature(doc_cfg))]
161#![warn(missing_docs)]
162use std::cmp::Ordering;
163use std::hash::{Hash, Hasher};
164use std::iter::Map;
165use std::slice::Windows;
166use std::sync::Arc;
167
168use arrow_array::cast::*;
169use arrow_array::types::{ArrowDictionaryKeyType, ByteArrayType, ByteViewType};
170use arrow_array::*;
171use arrow_buffer::{ArrowNativeType, Buffer, OffsetBuffer, ScalarBuffer};
172use arrow_schema::*;
173use variable::{decode_binary_view, decode_string_view};
174
175use crate::fixed::{decode_bool, decode_fixed_size_binary, decode_primitive};
176use crate::list::{compute_lengths_fixed_size_list, encode_fixed_size_list};
177use crate::variable::{decode_binary, decode_string};
178use arrow_array::types::{Int16Type, Int32Type, Int64Type};
179
180mod fixed;
181mod list;
182mod run;
183mod variable;
184
185#[derive(Debug)]
568pub struct RowConverter {
569 fields: Arc<[SortField]>,
570 codecs: Vec<Codec>,
572}
573
574#[derive(Debug)]
575enum Codec {
576 Stateless,
578 Dictionary(RowConverter, OwnedRow),
581 Struct(RowConverter, OwnedRow),
584 List(RowConverter),
586 Map(RowConverter),
588 RunEndEncoded(RowConverter),
590 Union(Vec<RowConverter>, Vec<i8>, Vec<OwnedRow>),
593}
594
595fn compute_list_view_bounds<O: OffsetSizeTrait>(array: &GenericListViewArray<O>) -> (usize, usize) {
598 if array.is_empty() {
599 return (0, 0);
600 }
601
602 let offsets = array.value_offsets();
603 let sizes = array.value_sizes();
604 let values_len = array.values().len();
605
606 let mut min_offset = usize::MAX;
607 let mut max_end = 0usize;
608
609 for i in 0..array.len() {
610 let offset = offsets[i].as_usize();
611 let size = sizes[i].as_usize();
612 let end = offset + size;
613
614 if size > 0 {
615 min_offset = min_offset.min(offset);
616 max_end = max_end.max(end);
617 }
618
619 if min_offset == 0 && max_end == values_len {
623 break;
624 }
625 }
626
627 if min_offset == usize::MAX {
628 (0, 0)
630 } else {
631 (min_offset, max_end)
632 }
633}
634
635impl Codec {
636 fn new(sort_field: &SortField) -> Result<Self, ArrowError> {
637 match &sort_field.data_type {
638 DataType::Dictionary(_, values) => {
639 let sort_field =
640 SortField::new_with_options(values.as_ref().clone(), sort_field.options);
641
642 let converter = RowConverter::new(vec![sort_field])?;
643 let null_array = new_null_array(values.as_ref(), 1);
644 let nulls = converter.convert_columns(&[null_array])?;
645
646 let owned = OwnedRow {
647 data: nulls.buffer.into(),
648 config: nulls.config,
649 };
650 Ok(Self::Dictionary(converter, owned))
651 }
652 DataType::RunEndEncoded(_, values) => {
653 let options = SortOptions {
655 descending: false,
656 nulls_first: sort_field.options.nulls_first != sort_field.options.descending,
657 };
658
659 let field = SortField::new_with_options(values.data_type().clone(), options);
660 let converter = RowConverter::new(vec![field])?;
661 Ok(Self::RunEndEncoded(converter))
662 }
663 d if !d.is_nested() => Ok(Self::Stateless),
664 DataType::List(f)
665 | DataType::LargeList(f)
666 | DataType::ListView(f)
667 | DataType::LargeListView(f) => {
668 let options = SortOptions {
672 descending: false,
673 nulls_first: sort_field.options.nulls_first != sort_field.options.descending,
674 };
675
676 let field = SortField::new_with_options(f.data_type().clone(), options);
677 let converter = RowConverter::new(vec![field])?;
678 Ok(Self::List(converter))
679 }
680 DataType::Map(f, _) => {
681 let options = SortOptions {
685 descending: false,
686 nulls_first: sort_field.options.nulls_first != sort_field.options.descending,
687 };
688
689 let DataType::Struct(fields) = f.data_type() else {
690 return Err(ArrowError::InvalidArgumentError(format!(
691 "expected struct field in map, got {:?}",
692 f.data_type()
693 )));
694 };
695
696 let fields = fields
698 .iter()
699 .map(|struct_field| {
700 SortField::new_with_options(struct_field.data_type().clone(), options)
701 })
702 .collect::<Vec<_>>();
703 assert_eq!(fields.len(), 2);
704 let converter = RowConverter::new(fields)?;
705 Ok(Self::Map(converter))
706 }
707 DataType::FixedSizeList(f, _) => {
708 let field = SortField::new_with_options(f.data_type().clone(), sort_field.options);
709 let converter = RowConverter::new(vec![field])?;
710 Ok(Self::List(converter))
711 }
712 DataType::Struct(f) => {
713 let sort_fields = f
714 .iter()
715 .map(|x| SortField::new_with_options(x.data_type().clone(), sort_field.options))
716 .collect();
717
718 let converter = RowConverter::new(sort_fields)?;
719 let nulls: Vec<_> = f.iter().map(|x| new_null_array(x.data_type(), 1)).collect();
720
721 let nulls = converter.convert_columns(&nulls)?;
722 let owned = OwnedRow {
723 data: nulls.buffer.into(),
724 config: nulls.config,
725 };
726
727 Ok(Self::Struct(converter, owned))
728 }
729 DataType::Union(fields, _mode) => {
730 let options = SortOptions {
733 descending: false,
734 nulls_first: sort_field.options.nulls_first != sort_field.options.descending,
735 };
736
737 let mut converters = Vec::with_capacity(fields.len());
738 let mut type_ids = Vec::with_capacity(fields.len());
739 let mut null_rows = Vec::with_capacity(fields.len());
740
741 for (type_id, field) in fields.iter() {
742 let sort_field =
743 SortField::new_with_options(field.data_type().clone(), options);
744 let converter = RowConverter::new(vec![sort_field])?;
745
746 let null_array = new_null_array(field.data_type(), 1);
747 let nulls = converter.convert_columns(&[null_array])?;
748 let owned = OwnedRow {
749 data: nulls.buffer.into(),
750 config: nulls.config,
751 };
752
753 converters.push(converter);
754 type_ids.push(type_id);
755 null_rows.push(owned);
756 }
757
758 Ok(Self::Union(converters, type_ids, null_rows))
759 }
760 _ => Err(ArrowError::NotYetImplemented(format!(
761 "not yet implemented: {:?}",
762 sort_field.data_type
763 ))),
764 }
765 }
766
767 fn encoder(&self, array: &dyn Array) -> Result<Encoder<'_>, ArrowError> {
768 match self {
769 Codec::Stateless => Ok(Encoder::Stateless),
770 Codec::Dictionary(converter, nulls) => {
771 let values = array.as_any_dictionary().values().clone();
772 let rows = converter.convert_columns(&[values])?;
773 Ok(Encoder::Dictionary(rows, nulls.row()))
774 }
775 Codec::Struct(converter, null) => {
776 let v = as_struct_array(array);
777 let rows = converter.convert_columns(v.columns())?;
778 Ok(Encoder::Struct(rows, null.row()))
779 }
780 Codec::List(converter) => {
781 let values = match array.data_type() {
782 DataType::List(_) => {
783 let list_array = as_list_array(array);
784 let first_offset = list_array.offsets()[0] as usize;
785 let last_offset =
786 list_array.offsets()[list_array.offsets().len() - 1] as usize;
787
788 list_array
791 .values()
792 .slice(first_offset, last_offset - first_offset)
793 }
794 DataType::LargeList(_) => {
795 let list_array = as_large_list_array(array);
796
797 let first_offset = list_array.offsets()[0] as usize;
798 let last_offset =
799 list_array.offsets()[list_array.offsets().len() - 1] as usize;
800
801 list_array
804 .values()
805 .slice(first_offset, last_offset - first_offset)
806 }
807 DataType::ListView(_) => {
808 let list_view_array = array.as_list_view::<i32>();
809 let (min_offset, max_end) = compute_list_view_bounds(list_view_array);
810 list_view_array
811 .values()
812 .slice(min_offset, max_end - min_offset)
813 }
814 DataType::LargeListView(_) => {
815 let list_view_array = array.as_list_view::<i64>();
816 let (min_offset, max_end) = compute_list_view_bounds(list_view_array);
817 list_view_array
818 .values()
819 .slice(min_offset, max_end - min_offset)
820 }
821 DataType::FixedSizeList(_, _) => {
822 as_fixed_size_list_array(array).values().clone()
823 }
824 _ => unreachable!(),
825 };
826 let rows = converter.convert_columns(&[values])?;
827 Ok(Encoder::List(rows))
828 }
829 Codec::Map(converter) => {
830 let map_array = as_map_array(array);
831
832 let first_offset = map_array.offsets()[0] as usize;
833 let last_offset = map_array.offsets()[map_array.offsets().len() - 1] as usize;
834
835 let sliced_entries = map_array
838 .entries()
839 .slice(first_offset, last_offset - first_offset);
840
841 let rows = converter.convert_columns(sliced_entries.columns())?;
843 Ok(Encoder::Map(rows))
844 }
845 Codec::RunEndEncoded(converter) => {
846 let values = match array.data_type() {
847 DataType::RunEndEncoded(r, _) => match r.data_type() {
848 DataType::Int16 => array.as_run::<Int16Type>().values_slice(),
849 DataType::Int32 => array.as_run::<Int32Type>().values_slice(),
850 DataType::Int64 => array.as_run::<Int64Type>().values_slice(),
851 _ => unreachable!("Unsupported run end index type: {r:?}"),
852 },
853 _ => unreachable!(),
854 };
855 let rows = converter.convert_columns(std::slice::from_ref(&values))?;
856 Ok(Encoder::RunEndEncoded(rows))
857 }
858 Codec::Union(converters, field_to_type_ids, _) => {
859 let union_array = array
860 .as_any()
861 .downcast_ref::<UnionArray>()
862 .expect("expected Union array");
863
864 let type_ids = union_array.type_ids().clone();
865 let offsets = union_array.offsets().cloned();
866
867 let mut child_rows = Vec::with_capacity(converters.len());
868 for (field_idx, converter) in converters.iter().enumerate() {
869 let type_id = field_to_type_ids[field_idx];
870 let child_array = union_array.child(type_id);
871 let rows = converter.convert_columns(std::slice::from_ref(child_array))?;
872 child_rows.push(rows);
873 }
874
875 Ok(Encoder::Union {
876 child_rows,
877 field_to_type_ids: field_to_type_ids.clone(),
878 type_ids,
879 offsets,
880 })
881 }
882 }
883 }
884
885 fn size(&self) -> usize {
886 match self {
887 Codec::Stateless => 0,
888 Codec::Dictionary(converter, nulls) => converter.size() + nulls.data.len(),
889 Codec::Struct(converter, nulls) => converter.size() + nulls.data.len(),
890 Codec::List(converter) => converter.size(),
891 Codec::Map(converter) => converter.size(),
892 Codec::RunEndEncoded(converter) => converter.size(),
893 Codec::Union(converters, _, null_rows) => {
894 converters.iter().map(|c| c.size()).sum::<usize>()
895 + null_rows.iter().map(|n| n.data.len()).sum::<usize>()
896 }
897 }
898 }
899}
900
901#[derive(Debug)]
902enum Encoder<'a> {
903 Stateless,
905 Dictionary(Rows, Row<'a>),
907 Struct(Rows, Row<'a>),
913 List(Rows),
915 Map(Rows),
917 RunEndEncoded(Rows),
919 Union {
921 child_rows: Vec<Rows>,
922 field_to_type_ids: Vec<i8>,
923 type_ids: ScalarBuffer<i8>,
924 offsets: Option<ScalarBuffer<i32>>,
925 },
926}
927
928#[derive(Debug, Clone, PartialEq, Eq)]
930pub struct SortField {
931 options: SortOptions,
933 data_type: DataType,
935}
936
937impl SortField {
938 pub fn new(data_type: DataType) -> Self {
940 Self::new_with_options(data_type, Default::default())
941 }
942
943 pub fn new_with_options(data_type: DataType, options: SortOptions) -> Self {
945 Self { options, data_type }
946 }
947
948 pub fn size(&self) -> usize {
952 self.data_type.size() + std::mem::size_of::<Self>() - std::mem::size_of::<DataType>()
953 }
954}
955
956impl RowConverter {
957 pub fn new(fields: Vec<SortField>) -> Result<Self, ArrowError> {
959 if !Self::supports_fields(&fields) {
960 return Err(ArrowError::NotYetImplemented(format!(
961 "Row format support not yet implemented for: {fields:?}"
962 )));
963 }
964
965 let codecs = fields.iter().map(Codec::new).collect::<Result<_, _>>()?;
966 Ok(Self {
967 fields: fields.into(),
968 codecs,
969 })
970 }
971
972 pub fn supports_fields(fields: &[SortField]) -> bool {
974 fields.iter().all(|x| Self::supports_datatype(&x.data_type))
975 }
976
977 fn supports_datatype(d: &DataType) -> bool {
978 match d {
979 _ if !d.is_nested() => true,
980 DataType::List(f)
981 | DataType::LargeList(f)
982 | DataType::ListView(f)
983 | DataType::LargeListView(f)
984 | DataType::FixedSizeList(f, _)
985 | DataType::Map(f, _) => Self::supports_datatype(f.data_type()),
986 DataType::Struct(f) => f.iter().all(|x| Self::supports_datatype(x.data_type())),
987 DataType::RunEndEncoded(_, values) => Self::supports_datatype(values.data_type()),
988 DataType::Union(fs, _mode) => fs
989 .iter()
990 .all(|(_, f)| Self::supports_datatype(f.data_type())),
991 _ => false,
992 }
993 }
994
995 pub fn convert_columns(&self, columns: &[ArrayRef]) -> Result<Rows, ArrowError> {
1005 let num_rows = columns.first().map(|x| x.len()).unwrap_or(0);
1006 let mut rows = self.empty_rows(num_rows, 0);
1007 self.append(&mut rows, columns)?;
1008 Ok(rows)
1009 }
1010
1011 pub fn append(&self, rows: &mut Rows, columns: &[ArrayRef]) -> Result<(), ArrowError> {
1042 assert!(
1043 Arc::ptr_eq(&rows.config.fields, &self.fields),
1044 "rows were not produced by this RowConverter"
1045 );
1046
1047 if columns.len() != self.fields.len() {
1048 return Err(ArrowError::InvalidArgumentError(format!(
1049 "Incorrect number of arrays provided to RowConverter, expected {} got {}",
1050 self.fields.len(),
1051 columns.len()
1052 )));
1053 }
1054 for colum in columns.iter().skip(1) {
1055 if colum.len() != columns[0].len() {
1056 return Err(ArrowError::InvalidArgumentError(format!(
1057 "RowConverter columns must all have the same length, expected {} got {}",
1058 columns[0].len(),
1059 colum.len()
1060 )));
1061 }
1062 }
1063
1064 let encoders = columns
1065 .iter()
1066 .zip(&self.codecs)
1067 .zip(self.fields.iter())
1068 .map(|((column, codec), field)| {
1069 if !column.data_type().equals_datatype(&field.data_type) {
1070 return Err(ArrowError::InvalidArgumentError(format!(
1071 "RowConverter column schema mismatch, expected {} got {}",
1072 field.data_type,
1073 column.data_type()
1074 )));
1075 }
1076 codec.encoder(column.as_ref())
1077 })
1078 .collect::<Result<Vec<_>, _>>()?;
1079
1080 let write_offset = rows.num_rows();
1081 let lengths = row_lengths(columns, &encoders);
1082 let total = lengths.extend_offsets(rows.offsets[write_offset], &mut rows.offsets);
1083 rows.buffer.resize(total, 0);
1084
1085 for ((column, field), encoder) in columns.iter().zip(self.fields.iter()).zip(encoders) {
1086 encode_column(
1088 &mut rows.buffer,
1089 &mut rows.offsets[write_offset..],
1090 column.as_ref(),
1091 field.options,
1092 &encoder,
1093 )
1094 }
1095
1096 if cfg!(debug_assertions) {
1097 assert_eq!(*rows.offsets.last().unwrap(), rows.buffer.len());
1098 rows.offsets
1099 .windows(2)
1100 .for_each(|w| assert!(w[0] <= w[1], "offsets should be monotonic"));
1101 }
1102
1103 Ok(())
1104 }
1105
1106 pub fn convert_rows<'a, I>(&self, rows: I) -> Result<Vec<ArrayRef>, ArrowError>
1114 where
1115 I: IntoIterator<Item = Row<'a>>,
1116 {
1117 let mut validate_utf8 = false;
1118 let mut rows: Vec<_> = rows
1119 .into_iter()
1120 .map(|row| {
1121 assert!(
1122 Arc::ptr_eq(&row.config.fields, &self.fields),
1123 "rows were not produced by this RowConverter"
1124 );
1125 validate_utf8 |= row.config.validate_utf8;
1126 row.data
1127 })
1128 .collect();
1129
1130 let result = unsafe { self.convert_raw(&mut rows, validate_utf8) }?;
1134
1135 if cfg!(debug_assertions) {
1136 for (i, row) in rows.iter().enumerate() {
1137 if !row.is_empty() {
1138 return Err(ArrowError::InvalidArgumentError(format!(
1139 "Codecs {codecs:?} did not consume all bytes for row {i}, remaining bytes: {row:?}",
1140 codecs = self.codecs
1141 )));
1142 }
1143 }
1144 }
1145
1146 Ok(result)
1147 }
1148
1149 pub fn empty_rows(&self, row_capacity: usize, data_capacity: usize) -> Rows {
1178 let mut offsets = Vec::with_capacity(row_capacity.saturating_add(1));
1179 offsets.push(0);
1180
1181 Rows {
1182 offsets,
1183 buffer: Vec::with_capacity(data_capacity),
1184 config: RowConfig {
1185 fields: self.fields.clone(),
1186 validate_utf8: false,
1187 },
1188 }
1189 }
1190
1191 pub fn from_binary(&self, array: BinaryArray) -> Rows {
1218 assert_eq!(
1219 array.null_count(),
1220 0,
1221 "can't construct Rows instance from array with nulls"
1222 );
1223 let (offsets, values, _) = array.into_parts();
1224 let offsets = offsets.iter().map(|&i| i.as_usize()).collect();
1225 let buffer = values.into_vec().unwrap_or_else(|values| values.to_vec());
1227 Rows {
1228 buffer,
1229 offsets,
1230 config: RowConfig {
1231 fields: Arc::clone(&self.fields),
1232 validate_utf8: true,
1233 },
1234 }
1235 }
1236
1237 unsafe fn convert_raw(
1243 &self,
1244 rows: &mut [&[u8]],
1245 validate_utf8: bool,
1246 ) -> Result<Vec<ArrayRef>, ArrowError> {
1247 self.fields
1248 .iter()
1249 .zip(&self.codecs)
1250 .map(|(field, codec)| unsafe { decode_column(field, rows, codec, validate_utf8) })
1251 .collect()
1252 }
1253
1254 pub fn parser(&self) -> RowParser {
1256 RowParser::new(Arc::clone(&self.fields))
1257 }
1258
1259 pub unsafe fn parser_skip_utf8_validation(&self) -> RowParser {
1264 unsafe { RowParser::with_skip_utf8_validate(Arc::clone(&self.fields)) }
1265 }
1266
1267 pub fn size(&self) -> usize {
1271 std::mem::size_of::<Self>()
1272 + self.fields.iter().map(|x| x.size()).sum::<usize>()
1273 + self.codecs.capacity() * std::mem::size_of::<Codec>()
1274 + self.codecs.iter().map(Codec::size).sum::<usize>()
1275 }
1276}
1277
1278#[derive(Debug)]
1280pub struct RowParser {
1281 config: RowConfig,
1282}
1283
1284impl RowParser {
1285 fn new(fields: Arc<[SortField]>) -> Self {
1286 Self {
1287 config: RowConfig {
1288 fields,
1289 validate_utf8: true,
1290 },
1291 }
1292 }
1293 unsafe fn with_skip_utf8_validate(fields: Arc<[SortField]>) -> Self {
1298 Self {
1299 config: RowConfig {
1300 fields,
1301 validate_utf8: false,
1302 },
1303 }
1304 }
1305
1306 pub fn parse<'a>(&'a self, bytes: &'a [u8]) -> Row<'a> {
1311 Row {
1312 data: bytes,
1313 config: &self.config,
1314 }
1315 }
1316}
1317
1318#[derive(Debug, Clone)]
1320struct RowConfig {
1321 fields: Arc<[SortField]>,
1323 validate_utf8: bool,
1325}
1326
1327#[derive(Debug, Clone)]
1331pub struct Rows {
1332 buffer: Vec<u8>,
1334 offsets: Vec<usize>,
1336 config: RowConfig,
1338}
1339
1340pub type RowLengthIter<'a> = Map<Windows<'a, usize>, fn(&'a [usize]) -> usize>;
1342
1343impl Rows {
1344 pub fn push(&mut self, row: Row<'_>) {
1346 assert!(
1347 Arc::ptr_eq(&row.config.fields, &self.config.fields),
1348 "row was not produced by this RowConverter"
1349 );
1350 self.config.validate_utf8 |= row.config.validate_utf8;
1351 self.buffer.extend_from_slice(row.data);
1352 self.offsets.push(self.buffer.len())
1353 }
1354
1355 pub fn reserve(&mut self, row_capacity: usize, data_capacity: usize) {
1357 self.buffer.reserve(data_capacity);
1358 self.offsets.reserve(row_capacity);
1359 }
1360
1361 pub fn row(&self, row: usize) -> Row<'_> {
1363 self.checked_row_end(row);
1364 unsafe { self.row_unchecked(row) }
1365 }
1366
1367 fn checked_row_end(&self, row: usize) -> usize {
1368 row.checked_add(1)
1369 .filter(|end| *end < self.offsets.len())
1370 .expect("row index out of bounds")
1371 }
1372
1373 pub unsafe fn row_unchecked(&self, index: usize) -> Row<'_> {
1378 let end = unsafe { self.offsets.get_unchecked(index + 1) };
1379 let start = unsafe { self.offsets.get_unchecked(index) };
1380 let data = unsafe { self.buffer.get_unchecked(*start..*end) };
1381 Row {
1382 data,
1383 config: &self.config,
1384 }
1385 }
1386
1387 pub fn row_len(&self, row: usize) -> usize {
1390 let end = self.checked_row_end(row);
1391
1392 self.offsets[end] - self.offsets[row]
1393 }
1394
1395 pub fn lengths(&self) -> RowLengthIter<'_> {
1397 self.offsets.windows(2).map(|w| w[1] - w[0])
1398 }
1399
1400 pub fn clear(&mut self) {
1402 self.offsets.truncate(1);
1403 self.buffer.clear();
1404 }
1405
1406 pub fn num_rows(&self) -> usize {
1408 self.offsets.len() - 1
1409 }
1410
1411 pub fn iter(&self) -> RowsIter<'_> {
1413 self.into_iter()
1414 }
1415
1416 pub fn size(&self) -> usize {
1420 std::mem::size_of::<Self>()
1422 + self.buffer.capacity()
1423 + self.offsets.capacity() * std::mem::size_of::<usize>()
1424 }
1425
1426 pub fn try_into_binary(self) -> Result<BinaryArray, ArrowError> {
1456 if self.buffer.len() > i32::MAX as usize {
1457 return Err(ArrowError::InvalidArgumentError(format!(
1458 "{}-byte rows buffer too long to convert into a i32-indexed BinaryArray",
1459 self.buffer.len()
1460 )));
1461 }
1462 let offsets_scalar = ScalarBuffer::from_iter(self.offsets.into_iter().map(i32::usize_as));
1464 let array = unsafe {
1466 BinaryArray::new_unchecked(
1467 OffsetBuffer::new_unchecked(offsets_scalar),
1468 Buffer::from_vec(self.buffer),
1469 None,
1470 )
1471 };
1472 Ok(array)
1473 }
1474}
1475
1476impl<'a> IntoIterator for &'a Rows {
1477 type Item = Row<'a>;
1478 type IntoIter = RowsIter<'a>;
1479
1480 fn into_iter(self) -> Self::IntoIter {
1481 RowsIter {
1482 rows: self,
1483 start: 0,
1484 end: self.num_rows(),
1485 }
1486 }
1487}
1488
1489#[derive(Debug)]
1491pub struct RowsIter<'a> {
1492 rows: &'a Rows,
1493 start: usize,
1494 end: usize,
1495}
1496
1497impl<'a> Iterator for RowsIter<'a> {
1498 type Item = Row<'a>;
1499
1500 fn next(&mut self) -> Option<Self::Item> {
1501 if self.end == self.start {
1502 return None;
1503 }
1504
1505 let row = unsafe { self.rows.row_unchecked(self.start) };
1507 self.start += 1;
1508 Some(row)
1509 }
1510
1511 fn size_hint(&self) -> (usize, Option<usize>) {
1512 let len = self.len();
1513 (len, Some(len))
1514 }
1515}
1516
1517impl ExactSizeIterator for RowsIter<'_> {
1518 fn len(&self) -> usize {
1519 self.end - self.start
1520 }
1521}
1522
1523impl DoubleEndedIterator for RowsIter<'_> {
1524 fn next_back(&mut self) -> Option<Self::Item> {
1525 if self.end == self.start {
1526 return None;
1527 }
1528
1529 self.end -= 1;
1530
1531 let row = unsafe { self.rows.row_unchecked(self.end) };
1534 Some(row)
1535 }
1536}
1537
1538#[derive(Debug, Copy, Clone)]
1547pub struct Row<'a> {
1548 data: &'a [u8],
1549 config: &'a RowConfig,
1550}
1551
1552impl<'a> Row<'a> {
1553 pub fn owned(&self) -> OwnedRow {
1555 OwnedRow {
1556 data: self.data.into(),
1557 config: self.config.clone(),
1558 }
1559 }
1560
1561 pub fn data(&self) -> &'a [u8] {
1563 self.data
1564 }
1565}
1566
1567impl PartialEq for Row<'_> {
1570 #[inline]
1571 fn eq(&self, other: &Self) -> bool {
1572 self.data.eq(other.data)
1573 }
1574}
1575
1576impl Eq for Row<'_> {}
1577
1578impl PartialOrd for Row<'_> {
1579 #[inline]
1580 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
1581 Some(self.cmp(other))
1582 }
1583}
1584
1585impl Ord for Row<'_> {
1586 #[inline]
1587 fn cmp(&self, other: &Self) -> Ordering {
1588 self.data.cmp(other.data)
1589 }
1590}
1591
1592impl Hash for Row<'_> {
1593 #[inline]
1594 fn hash<H: Hasher>(&self, state: &mut H) {
1595 self.data.hash(state)
1596 }
1597}
1598
1599impl AsRef<[u8]> for Row<'_> {
1600 #[inline]
1601 fn as_ref(&self) -> &[u8] {
1602 self.data
1603 }
1604}
1605
1606#[derive(Debug, Clone)]
1610pub struct OwnedRow {
1611 data: Box<[u8]>,
1612 config: RowConfig,
1613}
1614
1615impl OwnedRow {
1616 pub fn row(&self) -> Row<'_> {
1620 Row {
1621 data: &self.data,
1622 config: &self.config,
1623 }
1624 }
1625}
1626
1627impl PartialEq for OwnedRow {
1630 #[inline]
1631 fn eq(&self, other: &Self) -> bool {
1632 self.row().eq(&other.row())
1633 }
1634}
1635
1636impl Eq for OwnedRow {}
1637
1638impl PartialOrd for OwnedRow {
1639 #[inline]
1640 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
1641 Some(self.cmp(other))
1642 }
1643}
1644
1645impl Ord for OwnedRow {
1646 #[inline]
1647 fn cmp(&self, other: &Self) -> Ordering {
1648 self.row().cmp(&other.row())
1649 }
1650}
1651
1652impl Hash for OwnedRow {
1653 #[inline]
1654 fn hash<H: Hasher>(&self, state: &mut H) {
1655 self.row().hash(state)
1656 }
1657}
1658
1659impl AsRef<[u8]> for OwnedRow {
1660 #[inline]
1661 fn as_ref(&self) -> &[u8] {
1662 &self.data
1663 }
1664}
1665
1666#[inline]
1668fn null_sentinel(options: SortOptions) -> u8 {
1669 match options.nulls_first {
1670 true => 0,
1671 false => 0xFF,
1672 }
1673}
1674
1675enum LengthTracker {
1677 Fixed { length: usize, num_rows: usize },
1679 Variable {
1681 fixed_length: usize,
1682 lengths: Vec<usize>,
1683 },
1684}
1685
1686impl LengthTracker {
1687 fn new(num_rows: usize) -> Self {
1688 Self::Fixed {
1689 length: 0,
1690 num_rows,
1691 }
1692 }
1693
1694 fn push_fixed(&mut self, new_length: usize) {
1696 match self {
1697 LengthTracker::Fixed { length, .. } => *length += new_length,
1698 LengthTracker::Variable { fixed_length, .. } => *fixed_length += new_length,
1699 }
1700 }
1701
1702 fn push_variable(&mut self, new_lengths: impl ExactSizeIterator<Item = usize>) {
1704 match self {
1705 LengthTracker::Fixed { length, .. } => {
1706 *self = LengthTracker::Variable {
1707 fixed_length: *length,
1708 lengths: new_lengths.collect(),
1709 }
1710 }
1711 LengthTracker::Variable { lengths, .. } => {
1712 assert_eq!(lengths.len(), new_lengths.len());
1713 lengths
1714 .iter_mut()
1715 .zip(new_lengths)
1716 .for_each(|(length, new_length)| *length += new_length);
1717 }
1718 }
1719 }
1720
1721 fn materialized(&mut self) -> &mut [usize] {
1723 if let LengthTracker::Fixed { length, num_rows } = *self {
1724 *self = LengthTracker::Variable {
1725 fixed_length: length,
1726 lengths: vec![0; num_rows],
1727 };
1728 }
1729
1730 match self {
1731 LengthTracker::Variable { lengths, .. } => lengths,
1732 LengthTracker::Fixed { .. } => unreachable!(),
1733 }
1734 }
1735
1736 fn extend_offsets(&self, initial_offset: usize, offsets: &mut Vec<usize>) -> usize {
1754 match self {
1755 LengthTracker::Fixed { length, num_rows } => {
1756 offsets.extend((0..*num_rows).map(|i| initial_offset + i * length));
1757
1758 initial_offset + num_rows * length
1759 }
1760 LengthTracker::Variable {
1761 fixed_length,
1762 lengths,
1763 } => {
1764 let mut acc = initial_offset;
1765
1766 offsets.extend(lengths.iter().map(|length| {
1767 let current = acc;
1768 acc += length + fixed_length;
1769 current
1770 }));
1771
1772 acc
1773 }
1774 }
1775 }
1776}
1777
1778fn row_lengths(cols: &[ArrayRef], encoders: &[Encoder]) -> LengthTracker {
1780 use fixed::FixedLengthEncoding;
1781
1782 let num_rows = cols.first().map(|x| x.len()).unwrap_or(0);
1783 let mut tracker = LengthTracker::new(num_rows);
1784
1785 for (array, encoder) in cols.iter().zip(encoders) {
1786 match encoder {
1787 Encoder::Stateless => {
1788 downcast_primitive_array! {
1789 array => tracker.push_fixed(fixed::encoded_len(array)),
1790 DataType::Null => tracker.push_fixed(2)
1791 DataType::Boolean => tracker.push_fixed(bool::ENCODED_LEN),
1792 DataType::Binary => push_generic_byte_array_lengths(&mut tracker, as_generic_binary_array::<i32>(array)),
1793 DataType::LargeBinary => push_generic_byte_array_lengths(&mut tracker, as_generic_binary_array::<i64>(array)),
1794 DataType::BinaryView => push_byte_view_array_lengths(&mut tracker, array.as_binary_view()),
1795 DataType::Utf8 => push_generic_byte_array_lengths(&mut tracker, array.as_string::<i32>()),
1796 DataType::LargeUtf8 => push_generic_byte_array_lengths(&mut tracker, array.as_string::<i64>()),
1797 DataType::Utf8View => push_byte_view_array_lengths(&mut tracker, array.as_string_view()),
1798 DataType::FixedSizeBinary(len) => {
1799 let len = len.to_usize().unwrap();
1800 tracker.push_fixed(1 + len)
1801 }
1802 _ => unimplemented!("unsupported data type: {}", array.data_type()),
1803 }
1804 }
1805 Encoder::Dictionary(values, null) => {
1806 downcast_dictionary_array! {
1807 array => {
1808 tracker.push_variable(
1809 array.keys().iter().map(|v| match v {
1810 Some(k) => values.row_len(k.as_usize()),
1811 None => null.data.len(),
1812 })
1813 )
1814 }
1815 _ => unreachable!(),
1816 }
1817 }
1818 Encoder::Struct(rows, null) => {
1819 let array = as_struct_array(array);
1820 if rows.num_rows() > 0 {
1821 tracker.push_variable((0..array.len()).map(|idx| match array.is_valid(idx) {
1823 true => 1 + rows.row_len(idx),
1824 false => 1 + null.data.len(),
1825 }));
1826 } else {
1827 tracker.push_variable((0..array.len()).map(|idx| match array.is_valid(idx) {
1829 true => 1,
1830 false => 1 + null.data.len(),
1831 }));
1832 }
1833 }
1834 Encoder::List(rows) => match array.data_type() {
1835 DataType::List(_) => {
1836 list::compute_lengths(tracker.materialized(), rows, as_list_array(array))
1837 }
1838 DataType::LargeList(_) => {
1839 list::compute_lengths(tracker.materialized(), rows, as_large_list_array(array))
1840 }
1841 DataType::ListView(_) => {
1842 let list_view = array.as_list_view::<i32>();
1843 let (min_offset, _) = compute_list_view_bounds(list_view);
1844 list::compute_lengths_list_view(
1845 tracker.materialized(),
1846 rows,
1847 list_view,
1848 min_offset,
1849 )
1850 }
1851 DataType::LargeListView(_) => {
1852 let list_view = array.as_list_view::<i64>();
1853 let (min_offset, _) = compute_list_view_bounds(list_view);
1854 list::compute_lengths_list_view(
1855 tracker.materialized(),
1856 rows,
1857 list_view,
1858 min_offset,
1859 )
1860 }
1861 DataType::FixedSizeList(_, _) => compute_lengths_fixed_size_list(
1862 &mut tracker,
1863 rows,
1864 as_fixed_size_list_array(array),
1865 ),
1866 _ => unreachable!(),
1867 },
1868 Encoder::Map(rows) => {
1869 list::compute_lengths(tracker.materialized(), rows, as_map_array(array))
1870 }
1871 Encoder::RunEndEncoded(rows) => match array.data_type() {
1872 DataType::RunEndEncoded(r, _) => match r.data_type() {
1873 DataType::Int16 => run::compute_lengths(
1874 tracker.materialized(),
1875 rows,
1876 array.as_run::<Int16Type>(),
1877 ),
1878 DataType::Int32 => run::compute_lengths(
1879 tracker.materialized(),
1880 rows,
1881 array.as_run::<Int32Type>(),
1882 ),
1883 DataType::Int64 => run::compute_lengths(
1884 tracker.materialized(),
1885 rows,
1886 array.as_run::<Int64Type>(),
1887 ),
1888 _ => unreachable!("Unsupported run end index type: {r:?}"),
1889 },
1890 _ => unreachable!(),
1891 },
1892 Encoder::Union {
1893 child_rows,
1894 field_to_type_ids,
1895 type_ids,
1896 offsets,
1897 } => {
1898 let union_array = array
1899 .as_any()
1900 .downcast_ref::<UnionArray>()
1901 .expect("expected UnionArray");
1902
1903 let mut type_id_to_field_idx = [0usize; 128];
1904 for (field_idx, &type_id) in field_to_type_ids.iter().enumerate() {
1905 type_id_to_field_idx[type_id as usize] = field_idx;
1906 }
1907
1908 let lengths = (0..union_array.len()).map(|i| {
1909 let type_id = type_ids[i];
1910 let field_idx = type_id_to_field_idx[type_id as usize];
1911 let child_row_i = offsets.as_ref().map(|o| o[i] as usize).unwrap_or(i);
1912 let child_row_len = child_rows[field_idx].row_len(child_row_i);
1913
1914 1 + child_row_len
1916 });
1917
1918 tracker.push_variable(lengths);
1919 }
1920 }
1921 }
1922
1923 tracker
1924}
1925
1926fn push_generic_byte_array_lengths<T: ByteArrayType>(
1928 tracker: &mut LengthTracker,
1929 array: &GenericByteArray<T>,
1930) {
1931 if let Some(nulls) = array.nulls().filter(|n| n.null_count() > 0) {
1932 tracker.push_variable(
1933 array
1934 .offsets()
1935 .lengths()
1936 .zip(nulls.iter())
1937 .map(|(length, is_valid)| if is_valid { Some(length) } else { None })
1938 .map(variable::padded_length),
1939 )
1940 } else {
1941 tracker.push_variable(
1942 array
1943 .offsets()
1944 .lengths()
1945 .map(variable::non_null_padded_length),
1946 )
1947 }
1948}
1949
1950fn push_byte_view_array_lengths<T: ByteViewType>(
1952 tracker: &mut LengthTracker,
1953 array: &GenericByteViewArray<T>,
1954) {
1955 if let Some(nulls) = array.nulls().filter(|n| n.null_count() > 0) {
1956 tracker.push_variable(
1957 array
1958 .lengths()
1959 .zip(nulls.iter())
1960 .map(|(length, is_valid)| {
1961 if is_valid {
1962 Some(length as usize)
1963 } else {
1964 None
1965 }
1966 })
1967 .map(variable::padded_length),
1968 )
1969 } else {
1970 tracker.push_variable(
1971 array
1972 .lengths()
1973 .map(|len| variable::padded_length(Some(len as usize))),
1974 )
1975 }
1976}
1977
1978fn encode_column(
1980 data: &mut [u8],
1981 offsets: &mut [usize],
1982 column: &dyn Array,
1983 opts: SortOptions,
1984 encoder: &Encoder<'_>,
1985) {
1986 match encoder {
1987 Encoder::Stateless => {
1988 downcast_primitive_array! {
1989 column => {
1990 if let Some(nulls) = column.nulls().filter(|n| n.null_count() > 0){
1991 fixed::encode(data, offsets, column.values(), nulls, opts)
1992 } else {
1993 fixed::encode_not_null(data, offsets, column.values(), opts)
1994 }
1995 }
1996 DataType::Null => {
1997 for offset in offsets.iter_mut().skip(1) {
1998 variable::encode_null_value(&mut data[*offset..], opts);
1999 *offset += 2;
2000 }
2001 }
2002 DataType::Boolean => {
2003 if let Some(nulls) = column.nulls().filter(|n| n.null_count() > 0){
2004 fixed::encode_boolean(data, offsets, column.as_boolean().values(), nulls, opts)
2005 } else {
2006 fixed::encode_boolean_not_null(data, offsets, column.as_boolean().values(), opts)
2007 }
2008 }
2009 DataType::Binary => {
2010 variable::encode_generic_byte_array(data, offsets, as_generic_binary_array::<i32>(column), opts)
2011 }
2012 DataType::BinaryView => {
2013 variable::encode(data, offsets, column.as_binary_view().iter(), opts)
2014 }
2015 DataType::LargeBinary => {
2016 variable::encode_generic_byte_array(data, offsets, as_generic_binary_array::<i64>(column), opts)
2017 }
2018 DataType::Utf8 => variable::encode_generic_byte_array(
2019 data, offsets,
2020 column.as_string::<i32>(),
2021 opts,
2022 ),
2023 DataType::LargeUtf8 => variable::encode_generic_byte_array(
2024 data, offsets,
2025 column.as_string::<i64>(),
2026 opts,
2027 ),
2028 DataType::Utf8View => variable::encode(
2029 data, offsets,
2030 column.as_string_view().iter().map(|x| x.map(|x| x.as_bytes())),
2031 opts,
2032 ),
2033 DataType::FixedSizeBinary(_) => {
2034 let array = column.as_any().downcast_ref().unwrap();
2035 fixed::encode_fixed_size_binary(data, offsets, array, opts)
2036 }
2037 _ => unimplemented!("unsupported data type: {}", column.data_type()),
2038 }
2039 }
2040 Encoder::Dictionary(values, nulls) => {
2041 downcast_dictionary_array! {
2042 column => encode_dictionary_values(data, offsets, column, values, nulls),
2043 _ => unreachable!()
2044 }
2045 }
2046 Encoder::Struct(rows, null) => {
2047 fn struct_encode_helper<const NO_CHILD_FIELDS: bool>(
2048 array: &StructArray,
2049 offsets: &mut [usize],
2050 null_sentinel: u8,
2051 rows: &Rows,
2052 null: &Row<'_>,
2053 data: &mut [u8],
2054 ) {
2055 let empty_row = Row {
2056 data: &[],
2057 config: &rows.config,
2058 };
2059
2060 offsets
2061 .iter_mut()
2062 .skip(1)
2063 .enumerate()
2064 .for_each(|(idx, offset)| {
2065 let (row, sentinel) = match array.is_valid(idx) {
2066 true => (
2067 if NO_CHILD_FIELDS {
2068 empty_row
2069 } else {
2070 rows.row(idx)
2071 },
2072 0x01,
2073 ),
2074 false => (*null, null_sentinel),
2075 };
2076 let end_offset = *offset + 1 + row.as_ref().len();
2077 data[*offset] = sentinel;
2078 data[*offset + 1..end_offset].copy_from_slice(row.as_ref());
2079 *offset = end_offset;
2080 })
2081 }
2082
2083 let array = as_struct_array(column);
2084 let null_sentinel = null_sentinel(opts);
2085 if rows.num_rows() == 0 {
2086 struct_encode_helper::<true>(array, offsets, null_sentinel, rows, null, data);
2088 } else {
2089 struct_encode_helper::<false>(array, offsets, null_sentinel, rows, null, data);
2090 }
2091 }
2092 Encoder::List(rows) => match column.data_type() {
2093 DataType::List(_) => list::encode(data, offsets, rows, opts, as_list_array(column)),
2094 DataType::LargeList(_) => {
2095 list::encode(data, offsets, rows, opts, as_large_list_array(column))
2096 }
2097 DataType::ListView(_) => {
2098 let list_view = column.as_list_view::<i32>();
2099 let (min_offset, _) = compute_list_view_bounds(list_view);
2100 list::encode_list_view(data, offsets, rows, opts, list_view, min_offset)
2101 }
2102 DataType::LargeListView(_) => {
2103 let list_view = column.as_list_view::<i64>();
2104 let (min_offset, _) = compute_list_view_bounds(list_view);
2105 list::encode_list_view(data, offsets, rows, opts, list_view, min_offset)
2106 }
2107 DataType::FixedSizeList(_, _) => {
2108 encode_fixed_size_list(data, offsets, rows, opts, as_fixed_size_list_array(column))
2109 }
2110 _ => unreachable!(),
2111 },
2112 Encoder::Map(rows) => list::encode(data, offsets, rows, opts, as_map_array(column)),
2113 Encoder::RunEndEncoded(rows) => match column.data_type() {
2114 DataType::RunEndEncoded(r, _) => match r.data_type() {
2115 DataType::Int16 => {
2116 run::encode(data, offsets, rows, opts, column.as_run::<Int16Type>())
2117 }
2118 DataType::Int32 => {
2119 run::encode(data, offsets, rows, opts, column.as_run::<Int32Type>())
2120 }
2121 DataType::Int64 => {
2122 run::encode(data, offsets, rows, opts, column.as_run::<Int64Type>())
2123 }
2124 _ => unreachable!("Unsupported run end index type: {r:?}"),
2125 },
2126 _ => unreachable!(),
2127 },
2128 Encoder::Union {
2129 child_rows,
2130 field_to_type_ids,
2131 type_ids,
2132 offsets: offsets_buf,
2133 } => {
2134 let mut type_id_to_field_idx = [0usize; 128];
2135 for (field_idx, &type_id) in field_to_type_ids.iter().enumerate() {
2136 type_id_to_field_idx[type_id as usize] = field_idx;
2137 }
2138
2139 offsets
2140 .iter_mut()
2141 .skip(1)
2142 .enumerate()
2143 .for_each(|(i, offset)| {
2144 let type_id = type_ids[i];
2145 let field_idx = type_id_to_field_idx[type_id as usize];
2146
2147 let child_row_idx = offsets_buf.as_ref().map(|o| o[i] as usize).unwrap_or(i);
2148 let child_row = child_rows[field_idx].row(child_row_idx);
2149 let child_bytes = child_row.as_ref();
2150
2151 let type_id_byte = if opts.descending {
2152 !(type_id as u8)
2153 } else {
2154 type_id as u8
2155 };
2156 data[*offset] = type_id_byte;
2157
2158 let child_start = *offset + 1;
2159 let child_end = child_start + child_bytes.len();
2160 data[child_start..child_end].copy_from_slice(child_bytes);
2161
2162 *offset = child_end;
2163 });
2164 }
2165 }
2166}
2167
2168pub fn encode_dictionary_values<K: ArrowDictionaryKeyType>(
2170 data: &mut [u8],
2171 offsets: &mut [usize],
2172 column: &DictionaryArray<K>,
2173 values: &Rows,
2174 null: &Row<'_>,
2175) {
2176 for (offset, k) in offsets.iter_mut().skip(1).zip(column.keys()) {
2177 let row = match k {
2178 Some(k) => values.row(k.as_usize()).data,
2179 None => null.data,
2180 };
2181 let end_offset = *offset + row.len();
2182 data[*offset..end_offset].copy_from_slice(row);
2183 *offset = end_offset;
2184 }
2185}
2186
2187macro_rules! decode_primitive_helper {
2188 ($t:ty, $rows:ident, $data_type:ident, $options:ident) => {
2189 Arc::new(decode_primitive::<$t>($rows, $data_type, $options))
2190 };
2191}
2192
2193unsafe fn decode_column(
2199 field: &SortField,
2200 rows: &mut [&[u8]],
2201 codec: &Codec,
2202 validate_utf8: bool,
2203) -> Result<ArrayRef, ArrowError> {
2204 let options = field.options;
2205
2206 let array: ArrayRef = match codec {
2207 Codec::Stateless => {
2208 let data_type = field.data_type.clone();
2209 downcast_primitive! {
2210 data_type => (decode_primitive_helper, rows, data_type, options),
2211 DataType::Null => {
2212 variable::decode_null_value(rows, options);
2213 Arc::new(NullArray::new(rows.len()))
2214 }
2215 DataType::Boolean => Arc::new(decode_bool(rows, options)),
2216 DataType::Binary => Arc::new(decode_binary::<i32>(rows, options)),
2217 DataType::LargeBinary => Arc::new(decode_binary::<i64>(rows, options)),
2218 DataType::BinaryView => Arc::new(decode_binary_view(rows, options)),
2219 DataType::FixedSizeBinary(size) => Arc::new(decode_fixed_size_binary(rows, size, options)),
2220 DataType::Utf8 => Arc::new(unsafe{ decode_string::<i32>(rows, options, validate_utf8) }),
2221 DataType::LargeUtf8 => Arc::new(unsafe { decode_string::<i64>(rows, options, validate_utf8) }),
2222 DataType::Utf8View => Arc::new(unsafe { decode_string_view(rows, options, validate_utf8) }),
2223 _ => return Err(ArrowError::NotYetImplemented(format!("unsupported data type: {data_type}" )))
2224 }
2225 }
2226 Codec::Dictionary(converter, _) => {
2227 let cols = unsafe { converter.convert_raw(rows, validate_utf8) }?;
2228 cols.into_iter().next().unwrap()
2229 }
2230 Codec::Struct(converter, _) => {
2231 let nulls = fixed::decode_nulls(rows);
2232 rows.iter_mut().for_each(|row| *row = &row[1..]);
2233 let children = unsafe { converter.convert_raw(rows, validate_utf8) }?;
2234
2235 let corrected_fields: Vec<Field> = match &field.data_type {
2238 DataType::Struct(struct_fields) => struct_fields
2239 .iter()
2240 .zip(children.iter())
2241 .map(|(orig_field, child_array)| {
2242 orig_field
2243 .as_ref()
2244 .clone()
2245 .with_data_type(child_array.data_type().clone())
2246 })
2247 .collect(),
2248 _ => unreachable!("Only Struct types should be corrected here"),
2249 };
2250
2251 Arc::new(unsafe {
2252 StructArray::new_unchecked_with_length(
2253 corrected_fields.into(),
2254 children,
2255 nulls,
2256 rows.len(),
2257 )
2258 })
2259 }
2260 Codec::List(converter) => match &field.data_type {
2261 DataType::List(_) => Arc::new(unsafe {
2262 list::decode::<GenericListArray<i32>>(converter, rows, field, validate_utf8)
2263 }?),
2264 DataType::LargeList(_) => Arc::new(unsafe {
2265 list::decode::<GenericListArray<i64>>(converter, rows, field, validate_utf8)
2266 }?),
2267 DataType::ListView(_) => Arc::new(unsafe {
2268 list::decode_list_view::<i32>(converter, rows, field, validate_utf8)
2269 }?),
2270 DataType::LargeListView(_) => Arc::new(unsafe {
2271 list::decode_list_view::<i64>(converter, rows, field, validate_utf8)
2272 }?),
2273 DataType::FixedSizeList(_, value_length) => Arc::new(unsafe {
2274 list::decode_fixed_size_list(
2275 converter,
2276 rows,
2277 field,
2278 validate_utf8,
2279 value_length.as_usize(),
2280 )
2281 }?),
2282 _ => unreachable!(),
2283 },
2284 Codec::Map(converter) => {
2285 Arc::new(unsafe { list::decode::<MapArray>(converter, rows, field, validate_utf8) }?)
2286 }
2287 Codec::RunEndEncoded(converter) => match &field.data_type {
2288 DataType::RunEndEncoded(run_ends, _) => match run_ends.data_type() {
2289 DataType::Int16 => Arc::new(unsafe {
2290 run::decode::<Int16Type>(converter, rows, field, validate_utf8)
2291 }?),
2292 DataType::Int32 => Arc::new(unsafe {
2293 run::decode::<Int32Type>(converter, rows, field, validate_utf8)
2294 }?),
2295 DataType::Int64 => Arc::new(unsafe {
2296 run::decode::<Int64Type>(converter, rows, field, validate_utf8)
2297 }?),
2298 _ => unreachable!(),
2299 },
2300 _ => unreachable!(),
2301 },
2302 Codec::Union(converters, field_to_type_ids, null_rows) => {
2303 let len = rows.len();
2304
2305 let DataType::Union(union_fields, mode) = &field.data_type else {
2306 unreachable!()
2307 };
2308
2309 let mut type_id_to_field_idx = [0usize; 128];
2310 for (field_idx, &type_id) in field_to_type_ids.iter().enumerate() {
2311 type_id_to_field_idx[type_id as usize] = field_idx;
2312 }
2313
2314 let mut type_ids = Vec::with_capacity(len);
2315 let mut rows_by_field: Vec<Vec<(usize, &[u8])>> = vec![Vec::new(); converters.len()];
2316
2317 for (idx, row) in rows.iter_mut().enumerate() {
2318 let type_id_byte = {
2319 let id = row[0];
2320 if options.descending { !id } else { id }
2321 };
2322
2323 let type_id = type_id_byte as i8;
2324 type_ids.push(type_id);
2325
2326 let field_idx = type_id_to_field_idx[type_id as usize];
2327
2328 let child_row = &row[1..];
2329 rows_by_field[field_idx].push((idx, child_row));
2330 }
2331
2332 let mut child_arrays: Vec<ArrayRef> = Vec::with_capacity(converters.len());
2333 let mut offsets = (*mode == UnionMode::Dense).then(|| Vec::with_capacity(len));
2334
2335 for (field_idx, converter) in converters.iter().enumerate() {
2336 let field_rows = &rows_by_field[field_idx];
2337
2338 match &mode {
2339 UnionMode::Dense => {
2340 if field_rows.is_empty() {
2341 let (_, field) = union_fields.iter().nth(field_idx).unwrap();
2342 child_arrays.push(arrow_array::new_empty_array(field.data_type()));
2343 continue;
2344 }
2345
2346 let mut child_data = field_rows
2347 .iter()
2348 .map(|(_, bytes)| *bytes)
2349 .collect::<Vec<_>>();
2350
2351 let child_array =
2352 unsafe { converter.convert_raw(&mut child_data, validate_utf8) }?;
2353
2354 for ((row_idx, original_bytes), remaining_bytes) in
2356 field_rows.iter().zip(child_data)
2357 {
2358 let consumed_length = 1 + original_bytes.len() - remaining_bytes.len();
2359 rows[*row_idx] = &rows[*row_idx][consumed_length..];
2360 }
2361
2362 child_arrays.push(child_array.into_iter().next().unwrap());
2363 }
2364 UnionMode::Sparse => {
2365 let mut sparse_data: Vec<&[u8]> = Vec::with_capacity(len);
2366 let mut field_row_iter = field_rows.iter().peekable();
2367 let null_row_bytes: &[u8] = &null_rows[field_idx].data;
2368
2369 for idx in 0..len {
2370 if let Some((next_idx, bytes)) = field_row_iter.peek() {
2371 if *next_idx == idx {
2372 sparse_data.push(*bytes);
2373
2374 field_row_iter.next();
2375 continue;
2376 }
2377 }
2378 sparse_data.push(null_row_bytes);
2379 }
2380
2381 let child_array =
2382 unsafe { converter.convert_raw(&mut sparse_data, validate_utf8) }?;
2383
2384 for (row_idx, child_row) in field_rows.iter() {
2386 let remaining_len = sparse_data[*row_idx].len();
2387 let consumed_length = 1 + child_row.len() - remaining_len;
2388 rows[*row_idx] = &rows[*row_idx][consumed_length..];
2389 }
2390
2391 child_arrays.push(child_array.into_iter().next().unwrap());
2392 }
2393 }
2394 }
2395
2396 if let Some(ref mut offsets_vec) = offsets {
2398 let mut count = vec![0i32; converters.len()];
2399 for type_id in &type_ids {
2400 let field_idx = *type_id as usize;
2401 offsets_vec.push(count[field_idx]);
2402
2403 count[field_idx] += 1;
2404 }
2405 }
2406
2407 let type_ids_buffer = ScalarBuffer::from(type_ids);
2408 let offsets_buffer = offsets.map(ScalarBuffer::from);
2409
2410 let union_array = UnionArray::try_new(
2411 union_fields.clone(),
2412 type_ids_buffer,
2413 offsets_buffer,
2414 child_arrays,
2415 )?;
2416
2417 Arc::new(union_array)
2420 }
2421 };
2422 Ok(array)
2423}
2424
2425#[cfg(test)]
2426mod tests {
2427 use arrow_array::builder::*;
2428 use arrow_array::types::*;
2429 use arrow_array::*;
2430 use arrow_buffer::{Buffer, OffsetBuffer};
2431 use arrow_buffer::{NullBuffer, i256};
2432 use arrow_cast::display::{ArrayFormatter, FormatOptions};
2433 use arrow_ord::sort::{LexicographicalComparator, SortColumn};
2434 use rand::distr::uniform::SampleUniform;
2435 use rand::distr::{Distribution, StandardUniform};
2436 use rand::prelude::StdRng;
2437 use rand::{Rng, RngCore, SeedableRng};
2438
2439 use super::*;
2440
2441 fn all_sort_options() -> [SortOptions; 4] {
2442 [
2443 SortOptions {
2444 descending: false,
2445 nulls_first: false,
2446 },
2447 SortOptions {
2448 descending: false,
2449 nulls_first: true,
2450 },
2451 SortOptions {
2452 descending: true,
2453 nulls_first: false,
2454 },
2455 SortOptions {
2456 descending: true,
2457 nulls_first: true,
2458 },
2459 ]
2460 }
2461
2462 #[test]
2463 fn test_fixed_width() {
2464 let cols = [
2465 Arc::new(Int16Array::from_iter([
2466 Some(1),
2467 Some(2),
2468 None,
2469 Some(-5),
2470 Some(2),
2471 Some(2),
2472 Some(0),
2473 ])) as ArrayRef,
2474 Arc::new(Float32Array::from_iter([
2475 Some(1.3),
2476 Some(2.5),
2477 None,
2478 Some(4.),
2479 Some(0.1),
2480 Some(-4.),
2481 Some(-0.),
2482 ])) as ArrayRef,
2483 ];
2484
2485 let converter = RowConverter::new(vec![
2486 SortField::new(DataType::Int16),
2487 SortField::new(DataType::Float32),
2488 ])
2489 .unwrap();
2490 let rows = converter.convert_columns(&cols).unwrap();
2491
2492 assert_eq!(rows.offsets, &[0, 8, 16, 24, 32, 40, 48, 56]);
2493 assert_eq!(
2494 rows.buffer,
2495 &[
2496 1, 128, 1, 1, 191, 166, 102, 102, 1, 128, 2, 1, 192, 32, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 1, 127, 251, 1, 192, 128, 0, 0, 1, 128, 2, 1, 189, 204, 204, 205, 1, 128, 2, 1, 63, 127, 255, 255, 1, 128, 0, 1, 127, 255, 255, 255 ]
2511 );
2512
2513 assert!(rows.row(3) < rows.row(6));
2514 assert!(rows.row(0) < rows.row(1));
2515 assert!(rows.row(3) < rows.row(0));
2516 assert!(rows.row(4) < rows.row(1));
2517 assert!(rows.row(5) < rows.row(4));
2518
2519 let back = converter.convert_rows(&rows).unwrap();
2520 for (expected, actual) in cols.iter().zip(&back) {
2521 assert_eq!(expected, actual);
2522 }
2523 }
2524
2525 fn test_roundtrip(sort_option: SortOptions, col: ArrayRef) {
2526 let converter = RowConverter::new(vec![SortField::new_with_options(
2527 col.data_type().clone(),
2528 sort_option,
2529 )])
2530 .unwrap();
2531 let rows = converter.convert_columns(&[Arc::clone(&col)]).unwrap();
2532 let back = converter.convert_rows(&rows).unwrap();
2533 assert_eq!(back.len(), 1);
2534 assert_eq!(&back[0], &col);
2535 back[0].to_data().validate_full().unwrap();
2536 }
2537
2538 #[test]
2539 fn test_zero_width_fixed_size_binary_roundtrip() {
2540 for sort_option in all_sort_options() {
2541 for with_null in [true, false] {
2544 let nulls = if with_null {
2545 Some(NullBuffer::from(vec![true, false, true, false, true]))
2546 } else {
2547 None
2548 };
2549 let col: ArrayRef = Arc::new(
2550 FixedSizeBinaryArray::try_new_with_len(0, Buffer::default(), nulls, 5).unwrap(),
2551 );
2552
2553 test_roundtrip(sort_option, col);
2554 }
2555 }
2556 }
2557
2558 #[test]
2559 fn test_zero_width_fixed_size_list_roundtrip() {
2560 for sort_option in all_sort_options() {
2561 for with_null in [true, false] {
2564 let nulls = if with_null {
2565 Some(NullBuffer::from(vec![true, false, true, false, true]))
2566 } else {
2567 None
2568 };
2569 let col: ArrayRef = Arc::new(
2570 FixedSizeListArray::try_new_with_length(
2571 Arc::new(Field::new("item", DataType::Boolean, false)),
2572 0,
2573 new_empty_array(&DataType::Boolean),
2574 nulls,
2575 5,
2576 )
2577 .unwrap(),
2578 );
2579
2580 test_roundtrip(sort_option, col);
2581 }
2582 }
2583 }
2584
2585 #[test]
2586 fn test_decimal32() {
2587 let converter = RowConverter::new(vec![SortField::new(DataType::Decimal32(
2588 DECIMAL32_MAX_PRECISION,
2589 7,
2590 ))])
2591 .unwrap();
2592 let col = Arc::new(
2593 Decimal32Array::from_iter([
2594 None,
2595 Some(i32::MIN),
2596 Some(-13),
2597 Some(46_i32),
2598 Some(5456_i32),
2599 Some(i32::MAX),
2600 ])
2601 .with_precision_and_scale(9, 7)
2602 .unwrap(),
2603 ) as ArrayRef;
2604
2605 let rows = converter.convert_columns(&[Arc::clone(&col)]).unwrap();
2606 for i in 0..rows.num_rows() - 1 {
2607 assert!(rows.row(i) < rows.row(i + 1));
2608 }
2609
2610 let back = converter.convert_rows(&rows).unwrap();
2611 assert_eq!(back.len(), 1);
2612 assert_eq!(col.as_ref(), back[0].as_ref())
2613 }
2614
2615 #[test]
2616 fn test_decimal64() {
2617 let converter = RowConverter::new(vec![SortField::new(DataType::Decimal64(
2618 DECIMAL64_MAX_PRECISION,
2619 7,
2620 ))])
2621 .unwrap();
2622 let col = Arc::new(
2623 Decimal64Array::from_iter([
2624 None,
2625 Some(i64::MIN),
2626 Some(-13),
2627 Some(46_i64),
2628 Some(5456_i64),
2629 Some(i64::MAX),
2630 ])
2631 .with_precision_and_scale(18, 7)
2632 .unwrap(),
2633 ) as ArrayRef;
2634
2635 let rows = converter.convert_columns(&[Arc::clone(&col)]).unwrap();
2636 for i in 0..rows.num_rows() - 1 {
2637 assert!(rows.row(i) < rows.row(i + 1));
2638 }
2639
2640 let back = converter.convert_rows(&rows).unwrap();
2641 assert_eq!(back.len(), 1);
2642 assert_eq!(col.as_ref(), back[0].as_ref())
2643 }
2644
2645 #[test]
2646 fn test_decimal128() {
2647 let converter = RowConverter::new(vec![SortField::new(DataType::Decimal128(
2648 DECIMAL128_MAX_PRECISION,
2649 7,
2650 ))])
2651 .unwrap();
2652 let col = Arc::new(
2653 Decimal128Array::from_iter([
2654 None,
2655 Some(i128::MIN),
2656 Some(-13),
2657 Some(46_i128),
2658 Some(5456_i128),
2659 Some(i128::MAX),
2660 ])
2661 .with_precision_and_scale(38, 7)
2662 .unwrap(),
2663 ) as ArrayRef;
2664
2665 let rows = converter.convert_columns(&[Arc::clone(&col)]).unwrap();
2666 for i in 0..rows.num_rows() - 1 {
2667 assert!(rows.row(i) < rows.row(i + 1));
2668 }
2669
2670 let back = converter.convert_rows(&rows).unwrap();
2671 assert_eq!(back.len(), 1);
2672 assert_eq!(col.as_ref(), back[0].as_ref())
2673 }
2674
2675 #[test]
2676 fn test_decimal256() {
2677 let converter = RowConverter::new(vec![SortField::new(DataType::Decimal256(
2678 DECIMAL256_MAX_PRECISION,
2679 7,
2680 ))])
2681 .unwrap();
2682 let col = Arc::new(
2683 Decimal256Array::from_iter([
2684 None,
2685 Some(i256::MIN),
2686 Some(i256::from_parts(0, -1)),
2687 Some(i256::from_parts(u128::MAX, -1)),
2688 Some(i256::from_parts(u128::MAX, 0)),
2689 Some(i256::from_parts(0, 46_i128)),
2690 Some(i256::from_parts(5, 46_i128)),
2691 Some(i256::MAX),
2692 ])
2693 .with_precision_and_scale(DECIMAL256_MAX_PRECISION, 7)
2694 .unwrap(),
2695 ) as ArrayRef;
2696
2697 let rows = converter.convert_columns(&[Arc::clone(&col)]).unwrap();
2698 for i in 0..rows.num_rows() - 1 {
2699 assert!(rows.row(i) < rows.row(i + 1));
2700 }
2701
2702 let back = converter.convert_rows(&rows).unwrap();
2703 assert_eq!(back.len(), 1);
2704 assert_eq!(col.as_ref(), back[0].as_ref())
2705 }
2706
2707 #[test]
2708 fn test_bool() {
2709 let converter = RowConverter::new(vec![SortField::new(DataType::Boolean)]).unwrap();
2710
2711 let col = Arc::new(BooleanArray::from_iter([None, Some(false), Some(true)])) as ArrayRef;
2712
2713 let rows = converter.convert_columns(&[Arc::clone(&col)]).unwrap();
2714 assert!(rows.row(2) > rows.row(1));
2715 assert!(rows.row(2) > rows.row(0));
2716 assert!(rows.row(1) > rows.row(0));
2717
2718 let cols = converter.convert_rows(&rows).unwrap();
2719 assert_eq!(&cols[0], &col);
2720
2721 let converter = RowConverter::new(vec![SortField::new_with_options(
2722 DataType::Boolean,
2723 SortOptions::default().desc().with_nulls_first(false),
2724 )])
2725 .unwrap();
2726
2727 let rows = converter.convert_columns(&[Arc::clone(&col)]).unwrap();
2728 assert!(rows.row(2) < rows.row(1));
2729 assert!(rows.row(2) < rows.row(0));
2730 assert!(rows.row(1) < rows.row(0));
2731 let cols = converter.convert_rows(&rows).unwrap();
2732 assert_eq!(&cols[0], &col);
2733 }
2734
2735 #[test]
2736 fn test_timezone() {
2737 let a =
2738 TimestampNanosecondArray::from(vec![1, 2, 3, 4, 5]).with_timezone("+01:00".to_string());
2739 let d = a.data_type().clone();
2740
2741 let converter = RowConverter::new(vec![SortField::new(a.data_type().clone())]).unwrap();
2742 let rows = converter.convert_columns(&[Arc::new(a) as _]).unwrap();
2743 let back = converter.convert_rows(&rows).unwrap();
2744 assert_eq!(back.len(), 1);
2745 assert_eq!(back[0].data_type(), &d);
2746
2747 let mut a = PrimitiveDictionaryBuilder::<Int32Type, TimestampNanosecondType>::new();
2749 a.append(34).unwrap();
2750 a.append_null();
2751 a.append(345).unwrap();
2752
2753 let dict = a.finish();
2755 let values = TimestampNanosecondArray::from(dict.values().to_data());
2756 let dict_with_tz = dict.with_values(Arc::new(values.with_timezone("+02:00")));
2757 let v = DataType::Timestamp(TimeUnit::Nanosecond, Some("+02:00".into()));
2758 let d = DataType::Dictionary(Box::new(DataType::Int32), Box::new(v.clone()));
2759
2760 assert_eq!(dict_with_tz.data_type(), &d);
2761 let converter = RowConverter::new(vec![SortField::new(d.clone())]).unwrap();
2762 let rows = converter
2763 .convert_columns(&[Arc::new(dict_with_tz) as _])
2764 .unwrap();
2765 let back = converter.convert_rows(&rows).unwrap();
2766 assert_eq!(back.len(), 1);
2767 assert_eq!(back[0].data_type(), &v);
2768 }
2769
2770 #[test]
2771 fn test_null_encoding() {
2772 let col = Arc::new(NullArray::new(10));
2773 let converter = RowConverter::new(vec![SortField::new(DataType::Null)]).unwrap();
2774 let rows = converter.convert_columns(&[col]).unwrap();
2775 assert_eq!(rows.num_rows(), 10);
2776 assert_eq!(rows.row(1).data.len(), 2);
2778 }
2779
2780 #[test]
2781 fn test_variable_width() {
2782 let col = Arc::new(StringArray::from_iter([
2783 Some("hello"),
2784 Some("he"),
2785 None,
2786 Some("foo"),
2787 Some(""),
2788 ])) as ArrayRef;
2789
2790 let converter = RowConverter::new(vec![SortField::new(DataType::Utf8)]).unwrap();
2791 let rows = converter.convert_columns(&[Arc::clone(&col)]).unwrap();
2792
2793 assert!(rows.row(1) < rows.row(0));
2794 assert!(rows.row(2) < rows.row(4));
2795 assert!(rows.row(3) < rows.row(0));
2796 assert!(rows.row(3) < rows.row(1));
2797
2798 let cols = converter.convert_rows(&rows).unwrap();
2799 assert_eq!(&cols[0], &col);
2800
2801 let col = Arc::new(BinaryArray::from_iter([
2802 None,
2803 Some(vec![0_u8; 0]),
2804 Some(vec![0_u8; 6]),
2805 Some(vec![0_u8; variable::MINI_BLOCK_SIZE]),
2806 Some(vec![0_u8; variable::MINI_BLOCK_SIZE + 1]),
2807 Some(vec![0_u8; variable::BLOCK_SIZE]),
2808 Some(vec![0_u8; variable::BLOCK_SIZE + 1]),
2809 Some(vec![1_u8; 6]),
2810 Some(vec![1_u8; variable::MINI_BLOCK_SIZE]),
2811 Some(vec![1_u8; variable::MINI_BLOCK_SIZE + 1]),
2812 Some(vec![1_u8; variable::BLOCK_SIZE]),
2813 Some(vec![1_u8; variable::BLOCK_SIZE + 1]),
2814 Some(vec![0xFF_u8; 6]),
2815 Some(vec![0xFF_u8; variable::MINI_BLOCK_SIZE]),
2816 Some(vec![0xFF_u8; variable::MINI_BLOCK_SIZE + 1]),
2817 Some(vec![0xFF_u8; variable::BLOCK_SIZE]),
2818 Some(vec![0xFF_u8; variable::BLOCK_SIZE + 1]),
2819 ])) as ArrayRef;
2820
2821 let converter = RowConverter::new(vec![SortField::new(DataType::Binary)]).unwrap();
2822 let rows = converter.convert_columns(&[Arc::clone(&col)]).unwrap();
2823
2824 for i in 0..rows.num_rows() {
2825 for j in i + 1..rows.num_rows() {
2826 assert!(
2827 rows.row(i) < rows.row(j),
2828 "{} < {} - {:?} < {:?}",
2829 i,
2830 j,
2831 rows.row(i),
2832 rows.row(j)
2833 );
2834 }
2835 }
2836
2837 let cols = converter.convert_rows(&rows).unwrap();
2838 assert_eq!(&cols[0], &col);
2839
2840 let converter = RowConverter::new(vec![SortField::new_with_options(
2841 DataType::Binary,
2842 SortOptions::default().desc().with_nulls_first(false),
2843 )])
2844 .unwrap();
2845 let rows = converter.convert_columns(&[Arc::clone(&col)]).unwrap();
2846
2847 for i in 0..rows.num_rows() {
2848 for j in i + 1..rows.num_rows() {
2849 assert!(
2850 rows.row(i) > rows.row(j),
2851 "{} > {} - {:?} > {:?}",
2852 i,
2853 j,
2854 rows.row(i),
2855 rows.row(j)
2856 );
2857 }
2858 }
2859
2860 let cols = converter.convert_rows(&rows).unwrap();
2861 assert_eq!(&cols[0], &col);
2862 }
2863
2864 fn dictionary_eq(a: &dyn Array, b: &dyn Array) {
2866 match b.data_type() {
2867 DataType::Dictionary(_, v) => {
2868 assert_eq!(a.data_type(), v.as_ref());
2869 let b = arrow_cast::cast(b, v).unwrap();
2870 assert_eq!(a, b.as_ref())
2871 }
2872 _ => assert_eq!(a, b),
2873 }
2874 }
2875
2876 #[test]
2877 fn test_string_dictionary() {
2878 let a = Arc::new(DictionaryArray::<Int32Type>::from_iter([
2879 Some("foo"),
2880 Some("hello"),
2881 Some("he"),
2882 None,
2883 Some("hello"),
2884 Some(""),
2885 Some("hello"),
2886 Some("hello"),
2887 ])) as ArrayRef;
2888
2889 let field = SortField::new(a.data_type().clone());
2890 let converter = RowConverter::new(vec![field]).unwrap();
2891 let rows_a = converter.convert_columns(&[Arc::clone(&a)]).unwrap();
2892
2893 assert!(rows_a.row(3) < rows_a.row(5));
2894 assert!(rows_a.row(2) < rows_a.row(1));
2895 assert!(rows_a.row(0) < rows_a.row(1));
2896 assert!(rows_a.row(3) < rows_a.row(0));
2897
2898 assert_eq!(rows_a.row(1), rows_a.row(4));
2899 assert_eq!(rows_a.row(1), rows_a.row(6));
2900 assert_eq!(rows_a.row(1), rows_a.row(7));
2901
2902 let cols = converter.convert_rows(&rows_a).unwrap();
2903 dictionary_eq(&cols[0], &a);
2904
2905 let b = Arc::new(DictionaryArray::<Int32Type>::from_iter([
2906 Some("hello"),
2907 None,
2908 Some("cupcakes"),
2909 ])) as ArrayRef;
2910
2911 let rows_b = converter.convert_columns(&[Arc::clone(&b)]).unwrap();
2912 assert_eq!(rows_a.row(1), rows_b.row(0));
2913 assert_eq!(rows_a.row(3), rows_b.row(1));
2914 assert!(rows_b.row(2) < rows_a.row(0));
2915
2916 let cols = converter.convert_rows(&rows_b).unwrap();
2917 dictionary_eq(&cols[0], &b);
2918
2919 let converter = RowConverter::new(vec![SortField::new_with_options(
2920 a.data_type().clone(),
2921 SortOptions::default().desc().with_nulls_first(false),
2922 )])
2923 .unwrap();
2924
2925 let rows_c = converter.convert_columns(&[Arc::clone(&a)]).unwrap();
2926 assert!(rows_c.row(3) > rows_c.row(5));
2927 assert!(rows_c.row(2) > rows_c.row(1));
2928 assert!(rows_c.row(0) > rows_c.row(1));
2929 assert!(rows_c.row(3) > rows_c.row(0));
2930
2931 let cols = converter.convert_rows(&rows_c).unwrap();
2932 dictionary_eq(&cols[0], &a);
2933
2934 let converter = RowConverter::new(vec![SortField::new_with_options(
2935 a.data_type().clone(),
2936 SortOptions::default().desc().with_nulls_first(true),
2937 )])
2938 .unwrap();
2939
2940 let rows_c = converter.convert_columns(&[Arc::clone(&a)]).unwrap();
2941 assert!(rows_c.row(3) < rows_c.row(5));
2942 assert!(rows_c.row(2) > rows_c.row(1));
2943 assert!(rows_c.row(0) > rows_c.row(1));
2944 assert!(rows_c.row(3) < rows_c.row(0));
2945
2946 let cols = converter.convert_rows(&rows_c).unwrap();
2947 dictionary_eq(&cols[0], &a);
2948 }
2949
2950 #[test]
2951 fn test_struct() {
2952 let a = Arc::new(Int32Array::from(vec![1, 1, 2, 2])) as ArrayRef;
2954 let a_f = Arc::new(Field::new("int", DataType::Int32, false));
2955 let u = Arc::new(StringArray::from(vec!["a", "b", "c", "d"])) as ArrayRef;
2956 let u_f = Arc::new(Field::new("s", DataType::Utf8, false));
2957 let s1 = Arc::new(StructArray::from(vec![(a_f, a), (u_f, u)])) as ArrayRef;
2958
2959 let sort_fields = vec![SortField::new(s1.data_type().clone())];
2960 let converter = RowConverter::new(sort_fields).unwrap();
2961 let r1 = converter.convert_columns(&[Arc::clone(&s1)]).unwrap();
2962
2963 for (a, b) in r1.iter().zip(r1.iter().skip(1)) {
2964 assert!(a < b);
2965 }
2966
2967 let back = converter.convert_rows(&r1).unwrap();
2968 assert_eq!(back.len(), 1);
2969 assert_eq!(&back[0], &s1);
2970
2971 let data = s1
2973 .to_data()
2974 .into_builder()
2975 .null_bit_buffer(Some(Buffer::from_slice_ref([0b00001010])))
2976 .null_count(2)
2977 .build()
2978 .unwrap();
2979
2980 let s2 = Arc::new(StructArray::from(data)) as ArrayRef;
2981 let r2 = converter.convert_columns(&[Arc::clone(&s2)]).unwrap();
2982 assert_eq!(r2.row(0), r2.row(2)); assert!(r2.row(0) < r2.row(1)); assert_ne!(r1.row(0), r2.row(0)); assert_eq!(r1.row(1), r2.row(1)); let back = converter.convert_rows(&r2).unwrap();
2988 assert_eq!(back.len(), 1);
2989 assert_eq!(&back[0], &s2);
2990
2991 back[0].to_data().validate_full().unwrap();
2992 }
2993
2994 #[test]
2995 fn test_dictionary_in_struct() {
2996 let builder = StringDictionaryBuilder::<Int32Type>::new();
2997 let mut struct_builder = StructBuilder::new(
2998 vec![Field::new_dictionary(
2999 "foo",
3000 DataType::Int32,
3001 DataType::Utf8,
3002 true,
3003 )],
3004 vec![Box::new(builder)],
3005 );
3006
3007 let dict_builder = struct_builder
3008 .field_builder::<StringDictionaryBuilder<Int32Type>>(0)
3009 .unwrap();
3010
3011 dict_builder.append_value("a");
3013 dict_builder.append_null();
3014 dict_builder.append_value("a");
3015 dict_builder.append_value("b");
3016
3017 for _ in 0..4 {
3018 struct_builder.append(true);
3019 }
3020
3021 let s = Arc::new(struct_builder.finish()) as ArrayRef;
3022 let sort_fields = vec![SortField::new(s.data_type().clone())];
3023 let converter = RowConverter::new(sort_fields).unwrap();
3024 let r = converter.convert_columns(&[Arc::clone(&s)]).unwrap();
3025
3026 let back = converter.convert_rows(&r).unwrap();
3027 let [s2] = back.try_into().unwrap();
3028
3029 assert_ne!(&s.data_type(), &s2.data_type());
3032 s2.to_data().validate_full().unwrap();
3033
3034 let s1_struct = s.as_struct();
3038 let s1_0 = s1_struct.column(0);
3039 let s1_idx_0 = s1_0.as_dictionary::<Int32Type>();
3040 let keys = s1_idx_0.keys();
3041 let values = s1_idx_0.values().as_string::<i32>();
3042 let s2_struct = s2.as_struct();
3044 let s2_0 = s2_struct.column(0);
3045 let s2_idx_0 = s2_0.as_string::<i32>();
3046
3047 for i in 0..keys.len() {
3048 if keys.is_null(i) {
3049 assert!(s2_idx_0.is_null(i));
3050 } else {
3051 let dict_index = keys.value(i) as usize;
3052 assert_eq!(values.value(dict_index), s2_idx_0.value(i));
3053 }
3054 }
3055 }
3056
3057 #[test]
3058 fn test_dictionary_in_struct_empty() {
3059 let ty = DataType::Struct(
3060 vec![Field::new_dictionary(
3061 "foo",
3062 DataType::Int32,
3063 DataType::Int32,
3064 false,
3065 )]
3066 .into(),
3067 );
3068 let s = arrow_array::new_empty_array(&ty);
3069
3070 let sort_fields = vec![SortField::new(s.data_type().clone())];
3071 let converter = RowConverter::new(sort_fields).unwrap();
3072 let r = converter.convert_columns(&[Arc::clone(&s)]).unwrap();
3073
3074 let back = converter.convert_rows(&r).unwrap();
3075 let [s2] = back.try_into().unwrap();
3076
3077 assert_ne!(&s.data_type(), &s2.data_type());
3080 s2.to_data().validate_full().unwrap();
3081 assert_eq!(s.len(), 0);
3082 assert_eq!(s2.len(), 0);
3083 }
3084
3085 #[test]
3086 fn test_list_of_string_dictionary() {
3087 let mut builder = ListBuilder::<StringDictionaryBuilder<Int32Type>>::default();
3088 builder.values().append("a").unwrap();
3090 builder.values().append("b").unwrap();
3091 builder.values().append("zero").unwrap();
3092 builder.values().append_null();
3093 builder.values().append("c").unwrap();
3094 builder.values().append("b").unwrap();
3095 builder.values().append("d").unwrap();
3096 builder.append(true);
3097 builder.append(false);
3099 builder.values().append("e").unwrap();
3101 builder.values().append("zero").unwrap();
3102 builder.values().append("a").unwrap();
3103 builder.append(true);
3104
3105 let a = Arc::new(builder.finish()) as ArrayRef;
3106 let data_type = a.data_type().clone();
3107
3108 let field = SortField::new(data_type.clone());
3109 let converter = RowConverter::new(vec![field]).unwrap();
3110 let rows = converter.convert_columns(&[Arc::clone(&a)]).unwrap();
3111
3112 let back = converter.convert_rows(&rows).unwrap();
3113 assert_eq!(back.len(), 1);
3114 let [a2] = back.try_into().unwrap();
3115
3116 assert_ne!(&a.data_type(), &a2.data_type());
3119
3120 a2.to_data().validate_full().unwrap();
3121
3122 let a2_list = a2.as_list::<i32>();
3123 let a1_list = a.as_list::<i32>();
3124
3125 let a1_0 = a1_list.value(0);
3128 let a1_idx_0 = a1_0.as_dictionary::<Int32Type>();
3129 let keys = a1_idx_0.keys();
3130 let values = a1_idx_0.values().as_string::<i32>();
3131 let a2_0 = a2_list.value(0);
3132 let a2_idx_0 = a2_0.as_string::<i32>();
3133
3134 for i in 0..keys.len() {
3135 if keys.is_null(i) {
3136 assert!(a2_idx_0.is_null(i));
3137 } else {
3138 let dict_index = keys.value(i) as usize;
3139 assert_eq!(values.value(dict_index), a2_idx_0.value(i));
3140 }
3141 }
3142
3143 assert!(a1_list.is_null(1));
3145 assert!(a2_list.is_null(1));
3146
3147 let a1_2 = a1_list.value(2);
3149 let a1_idx_2 = a1_2.as_dictionary::<Int32Type>();
3150 let keys = a1_idx_2.keys();
3151 let values = a1_idx_2.values().as_string::<i32>();
3152 let a2_2 = a2_list.value(2);
3153 let a2_idx_2 = a2_2.as_string::<i32>();
3154
3155 for i in 0..keys.len() {
3156 if keys.is_null(i) {
3157 assert!(a2_idx_2.is_null(i));
3158 } else {
3159 let dict_index = keys.value(i) as usize;
3160 assert_eq!(values.value(dict_index), a2_idx_2.value(i));
3161 }
3162 }
3163 }
3164
3165 #[test]
3166 fn test_primitive_dictionary() {
3167 let mut builder = PrimitiveDictionaryBuilder::<Int32Type, Int32Type>::new();
3168 builder.append(2).unwrap();
3169 builder.append(3).unwrap();
3170 builder.append(0).unwrap();
3171 builder.append_null();
3172 builder.append(5).unwrap();
3173 builder.append(3).unwrap();
3174 builder.append(-1).unwrap();
3175
3176 let a = builder.finish();
3177 let data_type = a.data_type().clone();
3178 let columns = [Arc::new(a) as ArrayRef];
3179
3180 let field = SortField::new(data_type.clone());
3181 let converter = RowConverter::new(vec![field]).unwrap();
3182 let rows = converter.convert_columns(&columns).unwrap();
3183 assert!(rows.row(0) < rows.row(1));
3184 assert!(rows.row(2) < rows.row(0));
3185 assert!(rows.row(3) < rows.row(2));
3186 assert!(rows.row(6) < rows.row(2));
3187 assert!(rows.row(3) < rows.row(6));
3188
3189 let back = converter.convert_rows(&rows).unwrap();
3190 assert_eq!(back.len(), 1);
3191 back[0].to_data().validate_full().unwrap();
3192 }
3193
3194 #[test]
3195 fn test_dictionary_nulls() {
3196 let values = Int32Array::from_iter([Some(1), Some(-1), None, Some(4), None]).into_data();
3197 let keys =
3198 Int32Array::from_iter([Some(0), Some(0), Some(1), Some(2), Some(4), None]).into_data();
3199
3200 let data_type = DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Int32));
3201 let data = keys
3202 .into_builder()
3203 .data_type(data_type.clone())
3204 .child_data(vec![values])
3205 .build()
3206 .unwrap();
3207
3208 let columns = [Arc::new(DictionaryArray::<Int32Type>::from(data)) as ArrayRef];
3209 let field = SortField::new(data_type.clone());
3210 let converter = RowConverter::new(vec![field]).unwrap();
3211 let rows = converter.convert_columns(&columns).unwrap();
3212
3213 assert_eq!(rows.row(0), rows.row(1));
3214 assert_eq!(rows.row(3), rows.row(4));
3215 assert_eq!(rows.row(4), rows.row(5));
3216 assert!(rows.row(3) < rows.row(0));
3217 }
3218
3219 #[test]
3220 fn test_from_binary_shared_buffer() {
3221 let converter = RowConverter::new(vec![SortField::new(DataType::Binary)]).unwrap();
3222 let array = Arc::new(BinaryArray::from_iter_values([&[0xFF]])) as _;
3223 let rows = converter.convert_columns(&[array]).unwrap();
3224 let binary_rows = rows.try_into_binary().expect("known-small rows");
3225 let _binary_rows_shared_buffer = binary_rows.clone();
3226
3227 let parsed = converter.from_binary(binary_rows);
3228
3229 converter.convert_rows(parsed.iter()).unwrap();
3230 }
3231
3232 #[test]
3233 #[should_panic(expected = "Encountered non UTF-8 data")]
3234 fn test_invalid_utf8() {
3235 let converter = RowConverter::new(vec![SortField::new(DataType::Binary)]).unwrap();
3236 let array = Arc::new(BinaryArray::from_iter_values([&[0xFF]])) as _;
3237 let rows = converter.convert_columns(&[array]).unwrap();
3238 let binary_row = rows.row(0);
3239
3240 let converter = RowConverter::new(vec![SortField::new(DataType::Utf8)]).unwrap();
3241 let parser = converter.parser();
3242 let utf8_row = parser.parse(binary_row.as_ref());
3243
3244 converter.convert_rows(std::iter::once(utf8_row)).unwrap();
3245 }
3246
3247 #[test]
3248 #[should_panic(expected = "Encountered non UTF-8 data")]
3249 fn test_invalid_utf8_array() {
3250 let converter = RowConverter::new(vec![SortField::new(DataType::Binary)]).unwrap();
3251 let array = Arc::new(BinaryArray::from_iter_values([&[0xFF]])) as _;
3252 let rows = converter.convert_columns(&[array]).unwrap();
3253 let binary_rows = rows.try_into_binary().expect("known-small rows");
3254
3255 let converter = RowConverter::new(vec![SortField::new(DataType::Utf8)]).unwrap();
3256 let parsed = converter.from_binary(binary_rows);
3257
3258 converter.convert_rows(parsed.iter()).unwrap();
3259 }
3260
3261 #[test]
3262 #[should_panic(expected = "index out of bounds")]
3263 fn test_invalid_empty() {
3264 let binary_row: &[u8] = &[];
3265
3266 let converter = RowConverter::new(vec![SortField::new(DataType::Utf8)]).unwrap();
3267 let parser = converter.parser();
3268 let utf8_row = parser.parse(binary_row.as_ref());
3269
3270 converter.convert_rows(std::iter::once(utf8_row)).unwrap();
3271 }
3272
3273 #[test]
3274 #[should_panic(expected = "index out of bounds")]
3275 fn test_invalid_empty_array() {
3276 let row: &[u8] = &[];
3277 let binary_rows = BinaryArray::from(vec![row]);
3278
3279 let converter = RowConverter::new(vec![SortField::new(DataType::Utf8)]).unwrap();
3280 let parsed = converter.from_binary(binary_rows);
3281
3282 converter.convert_rows(parsed.iter()).unwrap();
3283 }
3284
3285 #[test]
3286 #[should_panic(expected = "index out of bounds")]
3287 fn test_invalid_truncated() {
3288 let binary_row: &[u8] = &[0x02];
3289
3290 let converter = RowConverter::new(vec![SortField::new(DataType::Utf8)]).unwrap();
3291 let parser = converter.parser();
3292 let utf8_row = parser.parse(binary_row.as_ref());
3293
3294 converter.convert_rows(std::iter::once(utf8_row)).unwrap();
3295 }
3296
3297 #[test]
3298 #[should_panic(expected = "index out of bounds")]
3299 fn test_invalid_truncated_array() {
3300 let row: &[u8] = &[0x02];
3301 let binary_rows = BinaryArray::from(vec![row]);
3302
3303 let converter = RowConverter::new(vec![SortField::new(DataType::Utf8)]).unwrap();
3304 let parsed = converter.from_binary(binary_rows);
3305
3306 converter.convert_rows(parsed.iter()).unwrap();
3307 }
3308
3309 #[test]
3310 #[should_panic(expected = "rows were not produced by this RowConverter")]
3311 fn test_different_converter() {
3312 let values = Arc::new(Int32Array::from_iter([Some(1), Some(-1)]));
3313 let converter = RowConverter::new(vec![SortField::new(DataType::Int32)]).unwrap();
3314 let rows = converter.convert_columns(&[values]).unwrap();
3315
3316 let converter = RowConverter::new(vec![SortField::new(DataType::Int32)]).unwrap();
3317 let _ = converter.convert_rows(&rows);
3318 }
3319
3320 fn test_single_list<O: OffsetSizeTrait>() {
3321 let mut builder = GenericListBuilder::<O, _>::new(Int32Builder::new());
3322 builder.values().append_value(32);
3323 builder.values().append_value(52);
3324 builder.values().append_value(32);
3325 builder.append(true);
3326 builder.values().append_value(32);
3327 builder.values().append_value(52);
3328 builder.values().append_value(12);
3329 builder.append(true);
3330 builder.values().append_value(32);
3331 builder.values().append_value(52);
3332 builder.append(true);
3333 builder.values().append_value(32); builder.values().append_value(52); builder.append(false);
3336 builder.values().append_value(32);
3337 builder.values().append_null();
3338 builder.append(true);
3339 builder.append(true);
3340 builder.values().append_value(17); builder.values().append_null(); builder.append(false);
3343
3344 let list = Arc::new(builder.finish()) as ArrayRef;
3345 let d = list.data_type().clone();
3346
3347 let converter = RowConverter::new(vec![SortField::new(d.clone())]).unwrap();
3348
3349 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
3350 assert!(rows.row(0) > rows.row(1)); assert!(rows.row(2) < rows.row(1)); assert!(rows.row(3) < rows.row(2)); assert!(rows.row(4) < rows.row(2)); assert!(rows.row(5) < rows.row(2)); assert!(rows.row(3) < rows.row(5)); assert_eq!(rows.row(3), rows.row(6)); let back = converter.convert_rows(&rows).unwrap();
3359 assert_eq!(back.len(), 1);
3360 back[0].to_data().validate_full().unwrap();
3361 assert_eq!(&back[0], &list);
3362
3363 let options = SortOptions::default().asc().with_nulls_first(false);
3364 let field = SortField::new_with_options(d.clone(), options);
3365 let converter = RowConverter::new(vec![field]).unwrap();
3366 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
3367
3368 assert!(rows.row(0) > rows.row(1)); assert!(rows.row(2) < rows.row(1)); assert!(rows.row(3) > rows.row(2)); assert!(rows.row(4) > rows.row(2)); assert!(rows.row(5) < rows.row(2)); assert!(rows.row(3) > rows.row(5)); assert_eq!(rows.row(3), rows.row(6)); let back = converter.convert_rows(&rows).unwrap();
3377 assert_eq!(back.len(), 1);
3378 back[0].to_data().validate_full().unwrap();
3379 assert_eq!(&back[0], &list);
3380
3381 let options = SortOptions::default().desc().with_nulls_first(false);
3382 let field = SortField::new_with_options(d.clone(), options);
3383 let converter = RowConverter::new(vec![field]).unwrap();
3384 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
3385
3386 assert!(rows.row(0) < rows.row(1)); assert!(rows.row(2) > rows.row(1)); assert!(rows.row(3) > rows.row(2)); assert!(rows.row(4) > rows.row(2)); assert!(rows.row(5) > rows.row(2)); assert!(rows.row(3) > rows.row(5)); assert_eq!(rows.row(3), rows.row(6)); let back = converter.convert_rows(&rows).unwrap();
3395 assert_eq!(back.len(), 1);
3396 back[0].to_data().validate_full().unwrap();
3397 assert_eq!(&back[0], &list);
3398
3399 let options = SortOptions::default().desc().with_nulls_first(true);
3400 let field = SortField::new_with_options(d, options);
3401 let converter = RowConverter::new(vec![field]).unwrap();
3402 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
3403
3404 assert!(rows.row(0) < rows.row(1)); assert!(rows.row(2) > rows.row(1)); assert!(rows.row(3) < rows.row(2)); assert!(rows.row(4) < rows.row(2)); assert!(rows.row(5) > rows.row(2)); assert!(rows.row(3) < rows.row(5)); assert_eq!(rows.row(3), rows.row(6)); let back = converter.convert_rows(&rows).unwrap();
3413 assert_eq!(back.len(), 1);
3414 back[0].to_data().validate_full().unwrap();
3415 assert_eq!(&back[0], &list);
3416
3417 let sliced_list = list.slice(1, 5);
3418 let rows_on_sliced_list = converter
3419 .convert_columns(&[Arc::clone(&sliced_list)])
3420 .unwrap();
3421
3422 assert!(rows_on_sliced_list.row(1) > rows_on_sliced_list.row(0)); assert!(rows_on_sliced_list.row(2) < rows_on_sliced_list.row(1)); assert!(rows_on_sliced_list.row(3) < rows_on_sliced_list.row(1)); assert!(rows_on_sliced_list.row(4) > rows_on_sliced_list.row(1)); assert!(rows_on_sliced_list.row(2) < rows_on_sliced_list.row(4)); let back = converter.convert_rows(&rows_on_sliced_list).unwrap();
3429 assert_eq!(back.len(), 1);
3430 back[0].to_data().validate_full().unwrap();
3431 assert_eq!(&back[0], &sliced_list);
3432 }
3433
3434 fn test_nested_list<O: OffsetSizeTrait>() {
3435 let mut builder =
3436 GenericListBuilder::<O, _>::new(GenericListBuilder::<O, _>::new(Int32Builder::new()));
3437
3438 builder.values().values().append_value(1);
3439 builder.values().values().append_value(2);
3440 builder.values().append(true);
3441 builder.values().values().append_value(1);
3442 builder.values().values().append_null();
3443 builder.values().append(true);
3444 builder.append(true);
3445
3446 builder.values().values().append_value(1);
3447 builder.values().values().append_null();
3448 builder.values().append(true);
3449 builder.values().values().append_value(1);
3450 builder.values().values().append_null();
3451 builder.values().append(true);
3452 builder.append(true);
3453
3454 builder.values().values().append_value(1);
3455 builder.values().values().append_null();
3456 builder.values().append(true);
3457 builder.values().append(false);
3458 builder.append(true);
3459 builder.append(false);
3460
3461 builder.values().values().append_value(1);
3462 builder.values().values().append_value(2);
3463 builder.values().append(true);
3464 builder.append(true);
3465
3466 let list = Arc::new(builder.finish()) as ArrayRef;
3467 let d = list.data_type().clone();
3468
3469 let options = SortOptions::default().asc().with_nulls_first(true);
3477 let field = SortField::new_with_options(d.clone(), options);
3478 let converter = RowConverter::new(vec![field]).unwrap();
3479 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
3480
3481 assert!(rows.row(0) > rows.row(1));
3482 assert!(rows.row(1) > rows.row(2));
3483 assert!(rows.row(2) > rows.row(3));
3484 assert!(rows.row(4) < rows.row(0));
3485 assert!(rows.row(4) > rows.row(1));
3486
3487 let back = converter.convert_rows(&rows).unwrap();
3488 assert_eq!(back.len(), 1);
3489 back[0].to_data().validate_full().unwrap();
3490 assert_eq!(&back[0], &list);
3491
3492 let options = SortOptions::default().desc().with_nulls_first(true);
3493 let field = SortField::new_with_options(d.clone(), options);
3494 let converter = RowConverter::new(vec![field]).unwrap();
3495 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
3496
3497 assert!(rows.row(0) > rows.row(1));
3498 assert!(rows.row(1) > rows.row(2));
3499 assert!(rows.row(2) > rows.row(3));
3500 assert!(rows.row(4) > rows.row(0));
3501 assert!(rows.row(4) > rows.row(1));
3502
3503 let back = converter.convert_rows(&rows).unwrap();
3504 assert_eq!(back.len(), 1);
3505 back[0].to_data().validate_full().unwrap();
3506 assert_eq!(&back[0], &list);
3507
3508 let options = SortOptions::default().desc().with_nulls_first(false);
3509 let field = SortField::new_with_options(d, options);
3510 let converter = RowConverter::new(vec![field]).unwrap();
3511 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
3512
3513 assert!(rows.row(0) < rows.row(1));
3514 assert!(rows.row(1) < rows.row(2));
3515 assert!(rows.row(2) < rows.row(3));
3516 assert!(rows.row(4) > rows.row(0));
3517 assert!(rows.row(4) < rows.row(1));
3518
3519 let back = converter.convert_rows(&rows).unwrap();
3520 assert_eq!(back.len(), 1);
3521 back[0].to_data().validate_full().unwrap();
3522 assert_eq!(&back[0], &list);
3523
3524 let sliced_list = list.slice(1, 3);
3525 let rows = converter
3526 .convert_columns(&[Arc::clone(&sliced_list)])
3527 .unwrap();
3528
3529 assert!(rows.row(0) < rows.row(1));
3530 assert!(rows.row(1) < rows.row(2));
3531
3532 let back = converter.convert_rows(&rows).unwrap();
3533 assert_eq!(back.len(), 1);
3534 back[0].to_data().validate_full().unwrap();
3535 assert_eq!(&back[0], &sliced_list);
3536 }
3537
3538 #[test]
3539 fn test_list() {
3540 test_single_list::<i32>();
3541 test_nested_list::<i32>();
3542 }
3543
3544 #[test]
3545 fn test_large_list() {
3546 test_single_list::<i64>();
3547 test_nested_list::<i64>();
3548 }
3549
3550 fn test_single_list_view<O: OffsetSizeTrait>() {
3551 let mut builder = GenericListViewBuilder::<O, _>::new(Int32Builder::new());
3552 builder.values().append_value(32);
3553 builder.values().append_value(52);
3554 builder.values().append_value(32);
3555 builder.append(true);
3556 builder.values().append_value(32);
3557 builder.values().append_value(52);
3558 builder.values().append_value(12);
3559 builder.append(true);
3560 builder.values().append_value(32);
3561 builder.values().append_value(52);
3562 builder.append(true);
3563 builder.values().append_value(32); builder.values().append_value(52); builder.append(false);
3566 builder.values().append_value(32);
3567 builder.values().append_null();
3568 builder.append(true);
3569 builder.append(true);
3570 builder.values().append_value(17); builder.values().append_null(); builder.append(false);
3573
3574 let list = Arc::new(builder.finish()) as ArrayRef;
3575 let d = list.data_type().clone();
3576
3577 let converter = RowConverter::new(vec![SortField::new(d.clone())]).unwrap();
3578
3579 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
3580 assert!(rows.row(0) > rows.row(1)); assert!(rows.row(2) < rows.row(1)); assert!(rows.row(3) < rows.row(2)); assert!(rows.row(4) < rows.row(2)); assert!(rows.row(5) < rows.row(2)); assert!(rows.row(3) < rows.row(5)); assert_eq!(rows.row(3), rows.row(6)); let back = converter.convert_rows(&rows).unwrap();
3589 assert_eq!(back.len(), 1);
3590 back[0].to_data().validate_full().unwrap();
3591
3592 let back_list_view = back[0]
3594 .as_any()
3595 .downcast_ref::<GenericListViewArray<O>>()
3596 .unwrap();
3597 let orig_list_view = list
3598 .as_any()
3599 .downcast_ref::<GenericListViewArray<O>>()
3600 .unwrap();
3601
3602 assert_eq!(back_list_view.len(), orig_list_view.len());
3603 for i in 0..back_list_view.len() {
3604 assert_eq!(back_list_view.is_valid(i), orig_list_view.is_valid(i));
3605 if back_list_view.is_valid(i) {
3606 assert_eq!(&back_list_view.value(i), &orig_list_view.value(i));
3607 }
3608 }
3609
3610 let options = SortOptions::default().asc().with_nulls_first(false);
3611 let field = SortField::new_with_options(d.clone(), options);
3612 let converter = RowConverter::new(vec![field]).unwrap();
3613 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
3614
3615 assert!(rows.row(0) > rows.row(1)); assert!(rows.row(2) < rows.row(1)); assert!(rows.row(3) > rows.row(2)); assert!(rows.row(4) > rows.row(2)); assert!(rows.row(5) < rows.row(2)); assert!(rows.row(3) > rows.row(5)); assert_eq!(rows.row(3), rows.row(6)); let back = converter.convert_rows(&rows).unwrap();
3624 assert_eq!(back.len(), 1);
3625 back[0].to_data().validate_full().unwrap();
3626
3627 let options = SortOptions::default().desc().with_nulls_first(false);
3628 let field = SortField::new_with_options(d.clone(), options);
3629 let converter = RowConverter::new(vec![field]).unwrap();
3630 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
3631
3632 assert!(rows.row(0) < rows.row(1)); assert!(rows.row(2) > rows.row(1)); assert!(rows.row(3) > rows.row(2)); assert!(rows.row(4) > rows.row(2)); assert!(rows.row(5) > rows.row(2)); assert!(rows.row(3) > rows.row(5)); assert_eq!(rows.row(3), rows.row(6)); let back = converter.convert_rows(&rows).unwrap();
3641 assert_eq!(back.len(), 1);
3642 back[0].to_data().validate_full().unwrap();
3643
3644 let options = SortOptions::default().desc().with_nulls_first(true);
3645 let field = SortField::new_with_options(d, options);
3646 let converter = RowConverter::new(vec![field]).unwrap();
3647 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
3648
3649 assert!(rows.row(0) < rows.row(1)); assert!(rows.row(2) > rows.row(1)); assert!(rows.row(3) < rows.row(2)); assert!(rows.row(4) < rows.row(2)); assert!(rows.row(5) > rows.row(2)); assert!(rows.row(3) < rows.row(5)); assert_eq!(rows.row(3), rows.row(6)); let back = converter.convert_rows(&rows).unwrap();
3658 assert_eq!(back.len(), 1);
3659 back[0].to_data().validate_full().unwrap();
3660
3661 let sliced_list = list.slice(1, 5);
3662 let rows_on_sliced_list = converter
3663 .convert_columns(&[Arc::clone(&sliced_list)])
3664 .unwrap();
3665
3666 assert!(rows_on_sliced_list.row(1) > rows_on_sliced_list.row(0)); assert!(rows_on_sliced_list.row(2) < rows_on_sliced_list.row(1)); assert!(rows_on_sliced_list.row(3) < rows_on_sliced_list.row(1)); assert!(rows_on_sliced_list.row(4) > rows_on_sliced_list.row(1)); assert!(rows_on_sliced_list.row(2) < rows_on_sliced_list.row(4)); let back = converter.convert_rows(&rows_on_sliced_list).unwrap();
3673 assert_eq!(back.len(), 1);
3674 back[0].to_data().validate_full().unwrap();
3675 }
3676
3677 fn test_nested_list_view<O: OffsetSizeTrait>() {
3678 let mut builder = GenericListViewBuilder::<O, _>::new(GenericListViewBuilder::<O, _>::new(
3679 Int32Builder::new(),
3680 ));
3681
3682 builder.values().values().append_value(1);
3684 builder.values().values().append_value(2);
3685 builder.values().append(true);
3686 builder.values().values().append_value(1);
3687 builder.values().values().append_null();
3688 builder.values().append(true);
3689 builder.append(true);
3690
3691 builder.values().values().append_value(1);
3693 builder.values().values().append_null();
3694 builder.values().append(true);
3695 builder.values().values().append_value(1);
3696 builder.values().values().append_null();
3697 builder.values().append(true);
3698 builder.append(true);
3699
3700 builder.values().values().append_value(1);
3702 builder.values().values().append_null();
3703 builder.values().append(true);
3704 builder.values().append(false);
3705 builder.append(true);
3706
3707 builder.append(false);
3709
3710 builder.values().values().append_value(1);
3712 builder.values().values().append_value(2);
3713 builder.values().append(true);
3714 builder.append(true);
3715
3716 let list = Arc::new(builder.finish()) as ArrayRef;
3717 let d = list.data_type().clone();
3718
3719 let options = SortOptions::default().asc().with_nulls_first(true);
3727 let field = SortField::new_with_options(d.clone(), options);
3728 let converter = RowConverter::new(vec![field]).unwrap();
3729 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
3730
3731 assert!(rows.row(0) > rows.row(1));
3732 assert!(rows.row(1) > rows.row(2));
3733 assert!(rows.row(2) > rows.row(3));
3734 assert!(rows.row(4) < rows.row(0));
3735 assert!(rows.row(4) > rows.row(1));
3736
3737 let back = converter.convert_rows(&rows).unwrap();
3738 assert_eq!(back.len(), 1);
3739 back[0].to_data().validate_full().unwrap();
3740
3741 let back_list_view = back[0]
3743 .as_any()
3744 .downcast_ref::<GenericListViewArray<O>>()
3745 .unwrap();
3746 let orig_list_view = list
3747 .as_any()
3748 .downcast_ref::<GenericListViewArray<O>>()
3749 .unwrap();
3750
3751 assert_eq!(back_list_view.len(), orig_list_view.len());
3752 for i in 0..back_list_view.len() {
3753 assert_eq!(back_list_view.is_valid(i), orig_list_view.is_valid(i));
3754 if back_list_view.is_valid(i) {
3755 assert_eq!(&back_list_view.value(i), &orig_list_view.value(i));
3756 }
3757 }
3758
3759 let options = SortOptions::default().desc().with_nulls_first(true);
3760 let field = SortField::new_with_options(d.clone(), options);
3761 let converter = RowConverter::new(vec![field]).unwrap();
3762 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
3763
3764 assert!(rows.row(0) > rows.row(1));
3765 assert!(rows.row(1) > rows.row(2));
3766 assert!(rows.row(2) > rows.row(3));
3767 assert!(rows.row(4) > rows.row(0));
3768 assert!(rows.row(4) > rows.row(1));
3769
3770 let back = converter.convert_rows(&rows).unwrap();
3771 assert_eq!(back.len(), 1);
3772 back[0].to_data().validate_full().unwrap();
3773
3774 let back_list_view = back[0]
3776 .as_any()
3777 .downcast_ref::<GenericListViewArray<O>>()
3778 .unwrap();
3779
3780 assert_eq!(back_list_view.len(), orig_list_view.len());
3781 for i in 0..back_list_view.len() {
3782 assert_eq!(back_list_view.is_valid(i), orig_list_view.is_valid(i));
3783 if back_list_view.is_valid(i) {
3784 assert_eq!(&back_list_view.value(i), &orig_list_view.value(i));
3785 }
3786 }
3787
3788 let options = SortOptions::default().desc().with_nulls_first(false);
3789 let field = SortField::new_with_options(d.clone(), options);
3790 let converter = RowConverter::new(vec![field]).unwrap();
3791 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
3792
3793 assert!(rows.row(0) < rows.row(1));
3794 assert!(rows.row(1) < rows.row(2));
3795 assert!(rows.row(2) < rows.row(3));
3796 assert!(rows.row(4) > rows.row(0));
3797 assert!(rows.row(4) < rows.row(1));
3798
3799 let back = converter.convert_rows(&rows).unwrap();
3800 assert_eq!(back.len(), 1);
3801 back[0].to_data().validate_full().unwrap();
3802
3803 let back_list_view = back[0]
3805 .as_any()
3806 .downcast_ref::<GenericListViewArray<O>>()
3807 .unwrap();
3808
3809 assert_eq!(back_list_view.len(), orig_list_view.len());
3810 for i in 0..back_list_view.len() {
3811 assert_eq!(back_list_view.is_valid(i), orig_list_view.is_valid(i));
3812 if back_list_view.is_valid(i) {
3813 assert_eq!(&back_list_view.value(i), &orig_list_view.value(i));
3814 }
3815 }
3816
3817 let sliced_list = list.slice(1, 3);
3818 let rows = converter
3819 .convert_columns(&[Arc::clone(&sliced_list)])
3820 .unwrap();
3821
3822 assert!(rows.row(0) < rows.row(1));
3823 assert!(rows.row(1) < rows.row(2));
3824
3825 let back = converter.convert_rows(&rows).unwrap();
3826 assert_eq!(back.len(), 1);
3827 back[0].to_data().validate_full().unwrap();
3828 }
3829
3830 #[test]
3831 fn test_list_view() {
3832 test_single_list_view::<i32>();
3833 test_nested_list_view::<i32>();
3834 }
3835
3836 #[test]
3837 fn test_large_list_view() {
3838 test_single_list_view::<i64>();
3839 test_nested_list_view::<i64>();
3840 }
3841
3842 fn test_list_view_with_shared_values<O: OffsetSizeTrait>() {
3843 let values = Int32Array::from(vec![1, 2, 3, 4, 5, 6, 7, 8]);
3845 let field = Arc::new(Field::new_list_field(DataType::Int32, true));
3846
3847 let offsets = ScalarBuffer::<O>::from(vec![
3855 O::from_usize(0).unwrap(),
3856 O::from_usize(0).unwrap(),
3857 O::from_usize(5).unwrap(),
3858 O::from_usize(2).unwrap(),
3859 O::from_usize(1).unwrap(),
3860 O::from_usize(2).unwrap(),
3861 ]);
3862 let sizes = ScalarBuffer::<O>::from(vec![
3863 O::from_usize(3).unwrap(),
3864 O::from_usize(3).unwrap(),
3865 O::from_usize(2).unwrap(),
3866 O::from_usize(2).unwrap(),
3867 O::from_usize(4).unwrap(),
3868 O::from_usize(1).unwrap(),
3869 ]);
3870
3871 let list_view: GenericListViewArray<O> =
3872 GenericListViewArray::try_new(field, offsets, sizes, Arc::new(values), None).unwrap();
3873
3874 let d = list_view.data_type().clone();
3875 let list = Arc::new(list_view) as ArrayRef;
3876
3877 let converter = RowConverter::new(vec![SortField::new(d.clone())]).unwrap();
3878 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
3879
3880 assert_eq!(rows.row(0), rows.row(1));
3882
3883 assert!(rows.row(0) < rows.row(2));
3885
3886 assert!(rows.row(3) > rows.row(0));
3888
3889 assert!(rows.row(4) > rows.row(0));
3891
3892 assert!(rows.row(5) < rows.row(3));
3894
3895 assert!(rows.row(5) > rows.row(4));
3897
3898 let back = converter.convert_rows(&rows).unwrap();
3900 assert_eq!(back.len(), 1);
3901 back[0].to_data().validate_full().unwrap();
3902
3903 let back_list_view = back[0]
3905 .as_any()
3906 .downcast_ref::<GenericListViewArray<O>>()
3907 .unwrap();
3908 let orig_list_view = list
3909 .as_any()
3910 .downcast_ref::<GenericListViewArray<O>>()
3911 .unwrap();
3912
3913 assert_eq!(back_list_view.len(), orig_list_view.len());
3914 for i in 0..back_list_view.len() {
3915 assert_eq!(back_list_view.is_valid(i), orig_list_view.is_valid(i));
3916 if back_list_view.is_valid(i) {
3917 assert_eq!(&back_list_view.value(i), &orig_list_view.value(i));
3918 }
3919 }
3920
3921 let options = SortOptions::default().desc();
3923 let field = SortField::new_with_options(d, options);
3924 let converter = RowConverter::new(vec![field]).unwrap();
3925 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
3926
3927 assert_eq!(rows.row(0), rows.row(1)); assert!(rows.row(0) > rows.row(2)); assert!(rows.row(3) < rows.row(0)); let back = converter.convert_rows(&rows).unwrap();
3933 assert_eq!(back.len(), 1);
3934 back[0].to_data().validate_full().unwrap();
3935 }
3936
3937 #[test]
3938 fn test_list_view_shared_values() {
3939 test_list_view_with_shared_values::<i32>();
3940 }
3941
3942 #[test]
3943 fn test_large_list_view_shared_values() {
3944 test_list_view_with_shared_values::<i64>();
3945 }
3946
3947 #[test]
3948 fn test_fixed_size_list() {
3949 let mut builder = FixedSizeListBuilder::new(Int32Builder::new(), 3);
3950 builder.values().append_value(32);
3951 builder.values().append_value(52);
3952 builder.values().append_value(32);
3953 builder.append(true);
3954 builder.values().append_value(32);
3955 builder.values().append_value(52);
3956 builder.values().append_value(12);
3957 builder.append(true);
3958 builder.values().append_value(32);
3959 builder.values().append_value(52);
3960 builder.values().append_null();
3961 builder.append(true);
3962 builder.values().append_value(32); builder.values().append_value(52); builder.values().append_value(13); builder.append(false);
3966 builder.values().append_value(32);
3967 builder.values().append_null();
3968 builder.values().append_null();
3969 builder.append(true);
3970 builder.values().append_null();
3971 builder.values().append_null();
3972 builder.values().append_null();
3973 builder.append(true);
3974 builder.values().append_value(17); builder.values().append_null(); builder.values().append_value(77); builder.append(false);
3978
3979 let list = Arc::new(builder.finish()) as ArrayRef;
3980 let d = list.data_type().clone();
3981
3982 let converter = RowConverter::new(vec![SortField::new(d.clone())]).unwrap();
3984
3985 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
3986 assert!(rows.row(0) > rows.row(1)); assert!(rows.row(2) < rows.row(1)); assert!(rows.row(3) < rows.row(2)); assert!(rows.row(4) < rows.row(2)); assert!(rows.row(5) < rows.row(2)); assert!(rows.row(3) < rows.row(5)); assert_eq!(rows.row(3), rows.row(6)); let back = converter.convert_rows(&rows).unwrap();
3995 assert_eq!(back.len(), 1);
3996 back[0].to_data().validate_full().unwrap();
3997 assert_eq!(&back[0], &list);
3998
3999 let options = SortOptions::default().asc().with_nulls_first(false);
4001 let field = SortField::new_with_options(d.clone(), options);
4002 let converter = RowConverter::new(vec![field]).unwrap();
4003 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
4004 assert!(rows.row(0) > rows.row(1)); assert!(rows.row(2) > rows.row(1)); assert!(rows.row(3) > rows.row(2)); assert!(rows.row(4) > rows.row(2)); assert!(rows.row(5) > rows.row(2)); assert!(rows.row(3) > rows.row(5)); assert_eq!(rows.row(3), rows.row(6)); let back = converter.convert_rows(&rows).unwrap();
4013 assert_eq!(back.len(), 1);
4014 back[0].to_data().validate_full().unwrap();
4015 assert_eq!(&back[0], &list);
4016
4017 let options = SortOptions::default().desc().with_nulls_first(false);
4019 let field = SortField::new_with_options(d.clone(), options);
4020 let converter = RowConverter::new(vec![field]).unwrap();
4021 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
4022 assert!(rows.row(0) < rows.row(1)); assert!(rows.row(2) > rows.row(1)); assert!(rows.row(3) > rows.row(2)); assert!(rows.row(4) > rows.row(2)); assert!(rows.row(5) > rows.row(2)); assert!(rows.row(3) > rows.row(5)); assert_eq!(rows.row(3), rows.row(6)); let back = converter.convert_rows(&rows).unwrap();
4031 assert_eq!(back.len(), 1);
4032 back[0].to_data().validate_full().unwrap();
4033 assert_eq!(&back[0], &list);
4034
4035 let options = SortOptions::default().desc().with_nulls_first(true);
4037 let field = SortField::new_with_options(d, options);
4038 let converter = RowConverter::new(vec![field]).unwrap();
4039 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
4040
4041 assert!(rows.row(0) < rows.row(1)); assert!(rows.row(2) < rows.row(1)); assert!(rows.row(3) < rows.row(2)); assert!(rows.row(4) < rows.row(2)); assert!(rows.row(5) < rows.row(2)); assert!(rows.row(3) < rows.row(5)); assert_eq!(rows.row(3), rows.row(6)); let back = converter.convert_rows(&rows).unwrap();
4050 assert_eq!(back.len(), 1);
4051 back[0].to_data().validate_full().unwrap();
4052 assert_eq!(&back[0], &list);
4053
4054 let sliced_list = list.slice(1, 5);
4055 let rows_on_sliced_list = converter
4056 .convert_columns(&[Arc::clone(&sliced_list)])
4057 .unwrap();
4058
4059 assert!(rows_on_sliced_list.row(2) < rows_on_sliced_list.row(1)); assert!(rows_on_sliced_list.row(3) < rows_on_sliced_list.row(1)); assert!(rows_on_sliced_list.row(4) < rows_on_sliced_list.row(1)); assert!(rows_on_sliced_list.row(2) < rows_on_sliced_list.row(4)); let back = converter.convert_rows(&rows_on_sliced_list).unwrap();
4065 assert_eq!(back.len(), 1);
4066 back[0].to_data().validate_full().unwrap();
4067 assert_eq!(&back[0], &sliced_list);
4068 }
4069
4070 #[test]
4071 fn test_two_fixed_size_lists() {
4072 let mut first = FixedSizeListBuilder::new(UInt8Builder::new(), 1);
4073 first.values().append_value(100);
4075 first.append(true);
4076 first.values().append_value(101);
4078 first.append(true);
4079 first.values().append_value(102);
4081 first.append(true);
4082 first.values().append_null();
4084 first.append(true);
4085 first.values().append_null(); first.append(false);
4088 let first = Arc::new(first.finish()) as ArrayRef;
4089 let first_type = first.data_type().clone();
4090
4091 let mut second = FixedSizeListBuilder::new(UInt8Builder::new(), 1);
4092 second.values().append_value(200);
4094 second.append(true);
4095 second.values().append_value(201);
4097 second.append(true);
4098 second.values().append_value(202);
4100 second.append(true);
4101 second.values().append_null();
4103 second.append(true);
4104 second.values().append_null(); second.append(false);
4107 let second = Arc::new(second.finish()) as ArrayRef;
4108 let second_type = second.data_type().clone();
4109
4110 let converter = RowConverter::new(vec![
4111 SortField::new(first_type.clone()),
4112 SortField::new(second_type.clone()),
4113 ])
4114 .unwrap();
4115
4116 let rows = converter
4117 .convert_columns(&[Arc::clone(&first), Arc::clone(&second)])
4118 .unwrap();
4119
4120 let back = converter.convert_rows(&rows).unwrap();
4121 assert_eq!(back.len(), 2);
4122 back[0].to_data().validate_full().unwrap();
4123 assert_eq!(&back[0], &first);
4124 back[1].to_data().validate_full().unwrap();
4125 assert_eq!(&back[1], &second);
4126 }
4127
4128 #[test]
4129 fn test_fixed_size_list_with_variable_width_content() {
4130 let mut first = FixedSizeListBuilder::new(
4131 StructBuilder::from_fields(
4132 vec![
4133 Field::new(
4134 "timestamp",
4135 DataType::Timestamp(TimeUnit::Microsecond, Some(Arc::from("UTC"))),
4136 false,
4137 ),
4138 Field::new("offset_minutes", DataType::Int16, false),
4139 Field::new("time_zone", DataType::Utf8, false),
4140 ],
4141 1,
4142 ),
4143 1,
4144 );
4145 first
4147 .values()
4148 .field_builder::<TimestampMicrosecondBuilder>(0)
4149 .unwrap()
4150 .append_null();
4151 first
4152 .values()
4153 .field_builder::<Int16Builder>(1)
4154 .unwrap()
4155 .append_null();
4156 first
4157 .values()
4158 .field_builder::<StringBuilder>(2)
4159 .unwrap()
4160 .append_null();
4161 first.values().append(false);
4162 first.append(false);
4163 first
4165 .values()
4166 .field_builder::<TimestampMicrosecondBuilder>(0)
4167 .unwrap()
4168 .append_null();
4169 first
4170 .values()
4171 .field_builder::<Int16Builder>(1)
4172 .unwrap()
4173 .append_null();
4174 first
4175 .values()
4176 .field_builder::<StringBuilder>(2)
4177 .unwrap()
4178 .append_null();
4179 first.values().append(false);
4180 first.append(true);
4181 first
4183 .values()
4184 .field_builder::<TimestampMicrosecondBuilder>(0)
4185 .unwrap()
4186 .append_value(0);
4187 first
4188 .values()
4189 .field_builder::<Int16Builder>(1)
4190 .unwrap()
4191 .append_value(0);
4192 first
4193 .values()
4194 .field_builder::<StringBuilder>(2)
4195 .unwrap()
4196 .append_value("UTC");
4197 first.values().append(true);
4198 first.append(true);
4199 first
4201 .values()
4202 .field_builder::<TimestampMicrosecondBuilder>(0)
4203 .unwrap()
4204 .append_value(1126351800123456);
4205 first
4206 .values()
4207 .field_builder::<Int16Builder>(1)
4208 .unwrap()
4209 .append_value(120);
4210 first
4211 .values()
4212 .field_builder::<StringBuilder>(2)
4213 .unwrap()
4214 .append_value("Europe/Warsaw");
4215 first.values().append(true);
4216 first.append(true);
4217 let first = Arc::new(first.finish()) as ArrayRef;
4218 let first_type = first.data_type().clone();
4219
4220 let mut second = StringBuilder::new();
4221 second.append_value("somewhere near");
4222 second.append_null();
4223 second.append_value("Greenwich");
4224 second.append_value("Warsaw");
4225 let second = Arc::new(second.finish()) as ArrayRef;
4226 let second_type = second.data_type().clone();
4227
4228 let converter = RowConverter::new(vec![
4229 SortField::new(first_type.clone()),
4230 SortField::new(second_type.clone()),
4231 ])
4232 .unwrap();
4233
4234 let rows = converter
4235 .convert_columns(&[Arc::clone(&first), Arc::clone(&second)])
4236 .unwrap();
4237
4238 let back = converter.convert_rows(&rows).unwrap();
4239 assert_eq!(back.len(), 2);
4240 back[0].to_data().validate_full().unwrap();
4241 assert_eq!(&back[0], &first);
4242 back[1].to_data().validate_full().unwrap();
4243 assert_eq!(&back[1], &second);
4244 }
4245
4246 #[test]
4247 fn test_single_map() {
4248 let mut builder = MapBuilder::new(None, StringBuilder::new(), Int32Builder::new());
4249 builder.keys().append_value("hello");
4251 builder.values().append_value(1);
4252 builder.keys().append_value("world");
4253 builder.values().append_value(2);
4254 builder.append(true).unwrap();
4255
4256 builder.keys().append_value("foo");
4258 builder.values().append_value(3);
4259 builder.append(true).unwrap();
4260
4261 builder.append(true).unwrap();
4263
4264 builder.keys().append_value("masked_key");
4266 builder.values().append_value(999);
4267 builder.append(false).unwrap();
4268
4269 builder.append(false).unwrap();
4271
4272 builder.keys().append_value("bar");
4274 builder.values().append_null();
4275 builder.append(true).unwrap();
4276
4277 builder.keys().append_value("other_masked");
4279 builder.values().append_value(0);
4280 builder.append(false).unwrap();
4281
4282 builder.keys().append_value("a");
4284 builder.values().append_value(10);
4285 builder.keys().append_value("b");
4286 builder.values().append_value(20);
4287 builder.keys().append_value("c");
4288 builder.values().append_value(30);
4289 builder.append(true).unwrap();
4290
4291 let map = Arc::new(builder.finish()) as ArrayRef;
4292 let d = map.data_type().clone();
4293
4294 let converter = RowConverter::new(vec![SortField::new(d.clone())]).unwrap();
4295
4296 let rows = converter.convert_columns(&[Arc::clone(&map)]).unwrap();
4297
4298 assert_eq!(rows.row(3), rows.row(4));
4300 assert_eq!(rows.row(4), rows.row(6));
4301
4302 let back = converter.convert_rows(&rows).unwrap();
4303 assert_eq!(back.len(), 1);
4304 back[0].to_data().validate_full().unwrap();
4305 assert_eq!(&back[0], &map);
4306
4307 let sliced_map = map.slice(1, map.len() - 2);
4308 let rows_on_sliced = converter
4309 .convert_columns(&[Arc::clone(&sliced_map)])
4310 .unwrap();
4311
4312 let back = converter.convert_rows(&rows_on_sliced).unwrap();
4313 assert_eq!(back.len(), 1);
4314 back[0].to_data().validate_full().unwrap();
4315 assert_eq!(&back[0], &sliced_map);
4316 }
4317
4318 #[test]
4319 fn two_maps_with_different_keys_order_should_sort_by_entry_order() {
4320 let map_1: ArrayRef =
4321 Arc::new(MapArray::from_vec_of_maps::<StringArray, Int32Array, _, _>(
4322 vec![Some(vec![("hello", Some(1)), ("world", Some(2))])],
4323 false,
4324 ));
4325 let map_2: ArrayRef =
4327 Arc::new(MapArray::from_vec_of_maps::<StringArray, Int32Array, _, _>(
4328 vec![Some(vec![("world", Some(2)), ("hello", Some(1))])],
4329 false,
4330 ));
4331
4332 let converter = RowConverter::new(vec![SortField::new(map_1.data_type().clone())]).unwrap();
4333
4334 let map_1_rows = converter.convert_columns(&[Arc::clone(&map_1)]).unwrap();
4335 let map_2_rows = converter.convert_columns(&[Arc::clone(&map_2)]).unwrap();
4336
4337 assert_ne!(map_1_rows.row(0), map_2_rows.row(0));
4338 assert!(map_1_rows.row(0) < map_2_rows.row(0));
4339
4340 let back_1 = converter.convert_rows(&map_1_rows).unwrap();
4341 let back_2 = converter.convert_rows(&map_2_rows).unwrap();
4342 assert_eq!(&back_1[0], &map_1);
4343 assert_eq!(&back_2[0], &map_2);
4344 }
4345
4346 #[test]
4347 fn test_nested_map() {
4348 let mut builder = MapBuilder::new(
4350 None,
4351 StringBuilder::new(),
4352 MapBuilder::new(None, StringBuilder::new(), Int32Builder::new()),
4353 );
4354
4355 builder.keys().append_value("outer1");
4357 builder.values().keys().append_value("inner_a");
4358 builder.values().values().append_value(1);
4359 builder.values().keys().append_value("inner_b");
4360 builder.values().values().append_value(2);
4361 builder.values().append(true).unwrap();
4362 builder.keys().append_value("outer2");
4363 builder.values().keys().append_value("inner_c");
4364 builder.values().values().append_value(3);
4365 builder.values().append(true).unwrap();
4366 builder.append(true).unwrap();
4367
4368 builder.keys().append_value("x");
4370 builder.values().append(true).unwrap();
4371 builder.append(true).unwrap();
4372
4373 builder.keys().append_value("y");
4375 builder.values().keys().append_value("masked"); builder.values().values().append_value(0); builder.values().append(false).unwrap();
4378 builder.append(true).unwrap();
4379
4380 builder.keys().append_value("y");
4382 builder.values().append(false).unwrap(); builder.append(true).unwrap();
4384
4385 builder.keys().append_value("masked_outer"); builder.values().keys().append_value("masked_inner"); builder.values().values().append_value(0); builder.values().append(true).unwrap(); builder.append(false).unwrap();
4391
4392 builder.keys().append_value("masked_outer"); builder.values().append(false).unwrap(); builder.append(false).unwrap();
4396
4397 builder.append(false).unwrap(); builder.append(true).unwrap();
4402
4403 let map = Arc::new(builder.finish()) as ArrayRef;
4404 let d = map.data_type().clone();
4405
4406 let converter = RowConverter::new(vec![SortField::new(d.clone())]).unwrap();
4407
4408 let rows = converter.convert_columns(&[Arc::clone(&map)]).unwrap();
4409
4410 let back = converter.convert_rows(&rows).unwrap();
4411 assert_eq!(back.len(), 1);
4412 back[0].to_data().validate_full().unwrap();
4413 assert_eq!(&back[0], &map);
4414
4415 let sliced_map = map.slice(1, 3);
4416 let rows_on_sliced = converter
4417 .convert_columns(&[Arc::clone(&sliced_map)])
4418 .unwrap();
4419
4420 let back = converter.convert_rows(&rows_on_sliced).unwrap();
4421 assert_eq!(back.len(), 1);
4422 back[0].to_data().validate_full().unwrap();
4423 assert_eq!(&back[0], &sliced_map);
4424 }
4425
4426 #[test]
4427 fn test_single_map_with_non_nullable_values() {
4428 let value_field = Arc::new(Field::new("values", DataType::Int32, false));
4430 let mut builder = MapBuilder::new(None, StringBuilder::new(), Int32Builder::new())
4431 .with_values_field(value_field);
4432 builder.keys().append_value("a");
4434 builder.values().append_value(1);
4435 builder.keys().append_value("b");
4436 builder.values().append_value(2);
4437 builder.append(true).unwrap();
4438 builder.append(false).unwrap();
4440 builder.keys().append_value("c");
4442 builder.values().append_value(3);
4443 builder.append(true).unwrap();
4444 builder.append(true).unwrap();
4446 builder.keys().append_value("masked"); builder.values().append_value(0); builder.append(false).unwrap();
4450
4451 let map = Arc::new(builder.finish()) as ArrayRef;
4452 let d = map.data_type().clone();
4453
4454 let converter = RowConverter::new(vec![SortField::new(d.clone())]).unwrap();
4455
4456 let rows = converter.convert_columns(&[Arc::clone(&map)]).unwrap();
4457
4458 let back = converter.convert_rows(&rows).unwrap();
4459 assert_eq!(back.len(), 1);
4460 back[0].to_data().validate_full().unwrap();
4461 assert_eq!(&back[0], &map);
4462 }
4463
4464 #[test]
4465 fn test_single_map_with_non_nullable_map_but_with_nullable_values() {
4466 let value_field = Arc::new(Field::new("values", DataType::Int32, true));
4468 let mut builder = MapBuilder::new(None, StringBuilder::new(), Int32Builder::new())
4469 .with_values_field(value_field);
4470
4471 builder.keys().append_value("a");
4473 builder.values().append_value(1);
4474 builder.keys().append_value("b");
4475 builder.values().append_null();
4476 builder.append(true).unwrap();
4477 builder.keys().append_value("c");
4479 builder.values().append_null();
4480 builder.keys().append_value("d");
4481 builder.values().append_null();
4482 builder.append(true).unwrap();
4483 builder.append(true).unwrap();
4485 builder.keys().append_value("e");
4487 builder.values().append_value(5);
4488 builder.append(true).unwrap();
4489
4490 let map = Arc::new(builder.finish()) as ArrayRef;
4491 let d = map.data_type().clone();
4492
4493 let converter = RowConverter::new(vec![SortField::new(d.clone())]).unwrap();
4494
4495 let rows = converter.convert_columns(&[Arc::clone(&map)]).unwrap();
4496
4497 let back = converter.convert_rows(&rows).unwrap();
4498 assert_eq!(back.len(), 1);
4499 back[0].to_data().validate_full().unwrap();
4500 assert_eq!(&back[0], &map);
4501 }
4502
4503 #[test]
4504 fn test_map_all_nulls() {
4505 let mut builder = MapBuilder::new(None, StringBuilder::new(), Int32Builder::new());
4506 builder.keys().append_value("m1"); builder.values().append_value(1); builder.append(false).unwrap();
4510 builder.keys().append_value("m2"); builder.values().append_value(2); builder.append(false).unwrap();
4513
4514 builder.append(false).unwrap(); builder.keys().append_value("m3"); builder.values().append_value(3); builder.append(false).unwrap();
4519
4520 let map = Arc::new(builder.finish()) as ArrayRef;
4521 let d = map.data_type().clone();
4522
4523 let converter = RowConverter::new(vec![SortField::new(d.clone())]).unwrap();
4524
4525 let rows = converter.convert_columns(&[Arc::clone(&map)]).unwrap();
4526
4527 rows.iter().for_each(|row| assert_eq!(row, rows.row(0)));
4529
4530 let back = converter.convert_rows(&rows).unwrap();
4531 assert_eq!(back.len(), 1);
4532 back[0].to_data().validate_full().unwrap();
4533 assert_eq!(&back[0], &map);
4534 }
4535
4536 #[test]
4537 fn test_map_all_empty() {
4538 let mut builder = MapBuilder::new(None, StringBuilder::new(), Int32Builder::new());
4539 builder.append(true).unwrap();
4541 builder.append(true).unwrap();
4542 builder.append(true).unwrap();
4543
4544 let map = Arc::new(builder.finish()) as ArrayRef;
4545 let d = map.data_type().clone();
4546
4547 let converter = RowConverter::new(vec![SortField::new(d.clone())]).unwrap();
4548
4549 let rows = converter.convert_columns(&[Arc::clone(&map)]).unwrap();
4550
4551 rows.iter().for_each(|row| assert_eq!(row, rows.row(0)));
4553
4554 let back = converter.convert_rows(&rows).unwrap();
4555 assert_eq!(back.len(), 1);
4556 back[0].to_data().validate_full().unwrap();
4557 assert_eq!(&back[0], &map);
4558 }
4559
4560 #[test]
4561 fn test_map_empty_array() {
4562 let builder = MapBuilder::new(None, StringBuilder::new(), Int32Builder::new());
4564 let map = Arc::new(builder.finish_cloned()) as ArrayRef;
4565 let d = map.data_type().clone();
4566
4567 let converter = RowConverter::new(vec![SortField::new(d.clone())]).unwrap();
4568
4569 let rows = converter.convert_columns(&[Arc::clone(&map)]).unwrap();
4570
4571 let back = converter.convert_rows(&rows).unwrap();
4572 assert_eq!(back.len(), 1);
4573 back[0].to_data().validate_full().unwrap();
4574 assert_eq!(&back[0], &map);
4575 }
4576
4577 fn generate_primitive_array<K>(
4578 rng: &mut impl RngCore,
4579 len: usize,
4580 valid_percent: f64,
4581 ) -> PrimitiveArray<K>
4582 where
4583 K: ArrowPrimitiveType,
4584 StandardUniform: Distribution<K::Native>,
4585 {
4586 (0..len)
4587 .map(|_| rng.random_bool(valid_percent).then(|| rng.random()))
4588 .collect()
4589 }
4590
4591 fn generate_all_unique_primitive_array<K>(
4592 rng: &mut impl RngCore,
4593 len: usize,
4594 ) -> PrimitiveArray<K>
4595 where
4596 K: ArrowPrimitiveType,
4597 K::Native: Hash + Eq,
4598 StandardUniform: Distribution<K::Native>,
4599 {
4600 let possible_number_of_values = 2i32.saturating_pow(size_of::<K::Native>() as u32 * 8);
4601 assert!(
4602 len <= possible_number_of_values as usize,
4603 "len {len} is larger than the number of possible values {possible_number_of_values}"
4604 );
4605
4606 let mut seen = std::collections::HashSet::new();
4607 (0..len)
4608 .map(|_| {
4609 let mut value;
4610 loop {
4611 value = rng.random();
4612
4613 if seen.insert(value) {
4614 break;
4615 }
4616 }
4617
4618 Some(value)
4619 })
4620 .collect()
4621 }
4622
4623 fn generate_boolean_array(
4624 rng: &mut impl RngCore,
4625 len: usize,
4626 valid_percent: f64,
4627 ) -> BooleanArray {
4628 (0..len)
4629 .map(|_| rng.random_bool(valid_percent).then(|| rng.random_bool(0.5)))
4630 .collect()
4631 }
4632
4633 fn generate_strings<O: OffsetSizeTrait>(
4634 rng: &mut impl RngCore,
4635 len: usize,
4636 valid_percent: f64,
4637 ) -> GenericStringArray<O> {
4638 (0..len)
4639 .map(|_| {
4640 rng.random_bool(valid_percent).then(|| {
4641 let len = rng.random_range(0..100);
4642 let bytes = (0..len).map(|_| rng.random_range(0..128)).collect();
4643 String::from_utf8(bytes).unwrap()
4644 })
4645 })
4646 .collect()
4647 }
4648
4649 fn generate_string_view(
4650 rng: &mut impl RngCore,
4651 len: usize,
4652 valid_percent: f64,
4653 ) -> StringViewArray {
4654 (0..len)
4655 .map(|_| {
4656 rng.random_bool(valid_percent).then(|| {
4657 let len = rng.random_range(0..100);
4658 let bytes = (0..len).map(|_| rng.random_range(0..128)).collect();
4659 String::from_utf8(bytes).unwrap()
4660 })
4661 })
4662 .collect()
4663 }
4664
4665 fn generate_byte_view(
4666 rng: &mut impl RngCore,
4667 len: usize,
4668 valid_percent: f64,
4669 ) -> BinaryViewArray {
4670 (0..len)
4671 .map(|_| {
4672 rng.random_bool(valid_percent).then(|| {
4673 let len = rng.random_range(0..100);
4674 let bytes: Vec<_> = (0..len).map(|_| rng.random_range(0..128)).collect();
4675 bytes
4676 })
4677 })
4678 .collect()
4679 }
4680
4681 fn generate_fixed_stringview_column(len: usize) -> StringViewArray {
4682 let edge_cases = vec![
4683 Some("bar".to_string()),
4684 Some("bar\0".to_string()),
4685 Some("LongerThan12Bytes".to_string()),
4686 Some("LongerThan12Bytez".to_string()),
4687 Some("LongerThan12Bytes\0".to_string()),
4688 Some("LongerThan12Byt".to_string()),
4689 Some("backend one".to_string()),
4690 Some("backend two".to_string()),
4691 Some("a".repeat(257)),
4692 Some("a".repeat(300)),
4693 ];
4694
4695 let mut values = Vec::with_capacity(len);
4697 for i in 0..len {
4698 values.push(
4699 edge_cases
4700 .get(i % edge_cases.len())
4701 .cloned()
4702 .unwrap_or(None),
4703 );
4704 }
4705
4706 StringViewArray::from(values)
4707 }
4708
4709 fn generate_dictionary<K>(
4710 rng: &mut impl RngCore,
4711 values: ArrayRef,
4712 len: usize,
4713 valid_percent: f64,
4714 ) -> DictionaryArray<K>
4715 where
4716 K: ArrowDictionaryKeyType,
4717 K::Native: SampleUniform,
4718 {
4719 let min_key = K::Native::from_usize(0).unwrap();
4720 let max_key = K::Native::from_usize(values.len()).unwrap();
4721 let keys: PrimitiveArray<K> = (0..len)
4722 .map(|_| {
4723 rng.random_bool(valid_percent)
4724 .then(|| rng.random_range(min_key..max_key))
4725 })
4726 .collect();
4727
4728 let data_type =
4729 DataType::Dictionary(Box::new(K::DATA_TYPE), Box::new(values.data_type().clone()));
4730
4731 let data = keys
4732 .into_data()
4733 .into_builder()
4734 .data_type(data_type)
4735 .add_child_data(values.to_data())
4736 .build()
4737 .unwrap();
4738
4739 DictionaryArray::from(data)
4740 }
4741
4742 fn generate_fixed_size_binary(
4743 rng: &mut impl RngCore,
4744 len: usize,
4745 valid_percent: f64,
4746 ) -> FixedSizeBinaryArray {
4747 let width = rng.random_range(0..20);
4748 let mut builder = FixedSizeBinaryBuilder::new(width);
4749
4750 let mut b = vec![0; width as usize];
4751 for _ in 0..len {
4752 match rng.random_bool(valid_percent) {
4753 true => {
4754 b.iter_mut().for_each(|x| *x = rng.random());
4755 builder.append_value(&b).unwrap();
4756 }
4757 false => builder.append_null(),
4758 }
4759 }
4760
4761 builder.finish()
4762 }
4763
4764 fn generate_struct(rng: &mut impl RngCore, len: usize, valid_percent: f64) -> StructArray {
4765 let nulls = NullBuffer::from_iter((0..len).map(|_| rng.random_bool(valid_percent)));
4766 let a = generate_primitive_array::<Int32Type>(rng, len, valid_percent);
4767 let b = generate_strings::<i32>(rng, len, valid_percent);
4768 let fields = Fields::from(vec![
4769 Field::new("a", DataType::Int32, true),
4770 Field::new("b", DataType::Utf8, true),
4771 ]);
4772 let values = vec![Arc::new(a) as _, Arc::new(b) as _];
4773 StructArray::new(fields, values, Some(nulls))
4774 }
4775
4776 fn generate_list<R: RngCore, F>(
4777 rng: &mut R,
4778 len: usize,
4779 valid_percent: f64,
4780 values: F,
4781 ) -> ListArray
4782 where
4783 F: FnOnce(&mut R, usize) -> ArrayRef,
4784 {
4785 let offsets = OffsetBuffer::<i32>::from_lengths((0..len).map(|_| rng.random_range(0..10)));
4786 let values_len = offsets.last().unwrap().to_usize().unwrap();
4787 let values = values(rng, values_len);
4788 let nulls = NullBuffer::from_iter((0..len).map(|_| rng.random_bool(valid_percent)));
4789 let field = Arc::new(Field::new_list_field(values.data_type().clone(), true));
4790 ListArray::new(field, offsets, values, Some(nulls))
4791 }
4792
4793 fn generate_list_view<F>(
4794 rng: &mut impl RngCore,
4795 len: usize,
4796 valid_percent: f64,
4797 values: F,
4798 ) -> ListViewArray
4799 where
4800 F: FnOnce(usize) -> ArrayRef,
4801 {
4802 let sizes: Vec<i32> = (0..len).map(|_| rng.random_range(0..10)).collect();
4804 let values_len: usize = sizes.iter().map(|s| *s as usize).sum::<usize>().max(1);
4805 let values = values(values_len);
4806
4807 let offsets: Vec<i32> = sizes
4809 .iter()
4810 .map(|&size| {
4811 if size == 0 {
4812 0
4813 } else {
4814 rng.random_range(0..=(values_len as i32 - size))
4815 }
4816 })
4817 .collect();
4818
4819 let nulls = NullBuffer::from_iter((0..len).map(|_| rng.random_bool(valid_percent)));
4820 let field = Arc::new(Field::new_list_field(values.data_type().clone(), true));
4821 ListViewArray::new(
4822 field,
4823 ScalarBuffer::from(offsets),
4824 ScalarBuffer::from(sizes),
4825 values,
4826 Some(nulls),
4827 )
4828 }
4829
4830 fn generate_map<R: RngCore, KeysFn, ValuesFn>(
4831 rng: &mut R,
4832 len: usize,
4833 valid_percent: f64,
4834 gen_keys: KeysFn,
4835 gen_values: ValuesFn,
4836 ) -> MapArray
4837 where
4838 KeysFn: FnOnce(&mut R, usize) -> ArrayRef,
4839 ValuesFn: FnOnce(&mut R, usize) -> ArrayRef,
4840 {
4841 let offsets = OffsetBuffer::<i32>::from_lengths((0..len).map(|_| rng.random_range(0..10)));
4842 let entries_len = offsets.last().unwrap().to_usize().unwrap();
4843 let keys = gen_keys(rng, entries_len);
4844 let values = gen_values(rng, entries_len);
4845 let nulls = NullBuffer::from_iter((0..len).map(|_| rng.random_bool(valid_percent)));
4846 let field = Arc::new(Field::new_map(
4847 "",
4848 "entries",
4849 Field::new("keys", keys.data_type().clone(), false),
4850 Field::new("values", values.data_type().clone(), true),
4851 false,
4852 true,
4853 ));
4854 let DataType::Map(struct_field, _) = field.data_type() else {
4855 unreachable!();
4856 };
4857
4858 let DataType::Struct(fields) = struct_field.data_type() else {
4859 unreachable!();
4860 };
4861
4862 let entries = StructArray::new(fields.clone(), vec![keys, values], None);
4863
4864 let map_array = MapArray::new(struct_field.clone(), offsets, entries, Some(nulls), false);
4865
4866 assert_valid_map(&map_array);
4867
4868 map_array
4869 }
4870
4871 fn assert_valid_map(array: &MapArray) {
4881 let keys_arrow_row_converter =
4882 RowConverter::new(vec![SortField::new(array.key_type().clone())]).unwrap();
4883
4884 array.iter().enumerate().flat_map(|(index, entry)| entry.map(|entry| (index, Arc::clone(entry.column(0))))).for_each(|(entry_index, keys)| {
4885 let keys_as_rows = keys_arrow_row_converter.convert_columns(&[Arc::clone(&keys)]).expect("should be able to convert keys");
4886
4887 for i in 0..keys_as_rows.num_rows() {
4888 for j in (i + 1)..keys_as_rows.num_rows() {
4889 if keys_as_rows.row(i) == keys_as_rows.row(j) {
4890 let key_i = keys.slice(i, 1);
4891 let key_j = keys.slice(j, 1);
4892
4893 assert_ne!(keys_as_rows.row(i), keys_as_rows.row(j), "map keys should be unique, but key {i} and key {j} are equal in entry {entry_index}. key {i} value is {key_i:?} and key {j} value is {key_j:?}");
4894 }
4895 }
4896 }
4897 })
4898 }
4899
4900 fn generate_nulls(rng: &mut impl RngCore, len: usize) -> Option<NullBuffer> {
4901 Some(NullBuffer::from_iter(
4902 (0..len).map(|_| rng.random_bool(0.8)),
4903 ))
4904 }
4905
4906 fn change_underlying_null_values_for_primitive<T: ArrowPrimitiveType>(
4907 array: &PrimitiveArray<T>,
4908 ) -> PrimitiveArray<T> {
4909 let (dt, values, nulls) = array.clone().into_parts();
4910
4911 let new_values = ScalarBuffer::<T::Native>::from_iter(
4912 values
4913 .iter()
4914 .zip(nulls.as_ref().unwrap().iter())
4915 .map(|(val, is_valid)| {
4916 if is_valid {
4917 *val
4918 } else {
4919 val.add_wrapping(T::Native::usize_as(1))
4920 }
4921 }),
4922 );
4923
4924 PrimitiveArray::new(new_values, nulls).with_data_type(dt)
4925 }
4926
4927 fn change_underline_null_values_for_byte_array<T: ByteArrayType>(
4928 array: &GenericByteArray<T>,
4929 ) -> GenericByteArray<T> {
4930 let (offsets, values, nulls) = array.clone().into_parts();
4931
4932 let new_offsets = OffsetBuffer::<T::Offset>::from_lengths(
4933 offsets
4934 .lengths()
4935 .zip(nulls.as_ref().unwrap().iter())
4936 .map(|(len, is_valid)| if is_valid { len } else { len + 1 }),
4937 );
4938
4939 let mut new_bytes = Vec::<u8>::with_capacity(new_offsets[new_offsets.len() - 1].as_usize());
4940
4941 offsets
4942 .windows(2)
4943 .zip(nulls.as_ref().unwrap().iter())
4944 .for_each(|(start_and_end, is_valid)| {
4945 let start = start_and_end[0].as_usize();
4946 let end = start_and_end[1].as_usize();
4947 new_bytes.extend_from_slice(&values.as_slice()[start..end]);
4948
4949 if !is_valid {
4951 new_bytes.push(b'c');
4952 }
4953 });
4954
4955 GenericByteArray::<T>::new(new_offsets, Buffer::from_vec(new_bytes), nulls)
4956 }
4957
4958 fn change_underline_null_values_for_list_array<O: OffsetSizeTrait>(
4959 array: &GenericListArray<O>,
4960 ) -> GenericListArray<O> {
4961 let (field, offsets, values, nulls) = array.clone().into_parts();
4962
4963 let (new_values, new_offsets) = {
4964 let concat_values = offsets
4965 .windows(2)
4966 .zip(nulls.as_ref().unwrap().iter())
4967 .map(|(start_and_end, is_valid)| {
4968 let start = start_and_end[0].as_usize();
4969 let end = start_and_end[1].as_usize();
4970 if is_valid {
4971 return (start, end - start);
4972 }
4973
4974 if end == values.len() {
4976 (start, (end - start).saturating_sub(1))
4977 } else {
4978 (start, end - start + 1)
4979 }
4980 })
4981 .map(|(start, length)| values.slice(start, length))
4982 .collect::<Vec<_>>();
4983
4984 let new_offsets =
4985 OffsetBuffer::<O>::from_lengths(concat_values.iter().map(|s| s.len()));
4986
4987 let new_values = {
4988 let values = concat_values.iter().map(|a| a.as_ref()).collect::<Vec<_>>();
4989 arrow_select::concat::concat(&values).expect("should be able to concat")
4990 };
4991
4992 (new_values, new_offsets)
4993 };
4994
4995 GenericListArray::<O>::new(field, new_offsets, new_values, nulls)
4996 }
4997
4998 fn change_underline_null_values_for_map_array(array: &MapArray) -> MapArray {
4999 let (field, offsets, entries, nulls, ordered) = array.clone().into_parts();
5000 assert!(
5001 !ordered,
5002 "can't replace underlying null values for ordered map array as this can violate the ordering"
5003 );
5004
5005 let (new_entries, new_offsets) = {
5006 let concat_values = offsets
5007 .windows(2)
5008 .zip(nulls.as_ref().unwrap().iter())
5009 .map(|(start_and_end, is_valid)| {
5010 let start = start_and_end[0].as_usize();
5011 let end = start_and_end[1].as_usize();
5012 if is_valid {
5013 return (start, end - start);
5014 }
5015
5016 if end == entries.len() {
5018 (start, (end - start).saturating_sub(1))
5019 } else {
5020 (start, end - start + 1)
5022 }
5023 })
5024 .map(|(start, length)| entries.slice(start, length))
5025 .collect::<Vec<_>>();
5026
5027 let new_offsets = OffsetBuffer::from_lengths(concat_values.iter().map(|s| s.len()));
5028
5029 let new_values = {
5030 let values = concat_values
5031 .iter()
5032 .map(|a| a as &dyn Array)
5033 .collect::<Vec<_>>();
5034 arrow_select::concat::concat(&values).expect("should be able to concat")
5035 };
5036
5037 (new_values.as_struct().clone(), new_offsets)
5038 };
5039
5040 let new_map = MapArray::new(field, new_offsets, new_entries, nulls, ordered);
5041
5042 assert_valid_map(&new_map);
5043
5044 new_map
5045 }
5046
5047 fn change_underline_null_values(array: &ArrayRef) -> ArrayRef {
5048 if array.null_count() == 0 {
5049 return Arc::clone(array);
5050 }
5051
5052 downcast_primitive_array!(
5053 array => {
5054 let output = change_underlying_null_values_for_primitive(array);
5055
5056 Arc::new(output)
5057 }
5058
5059 DataType::Utf8 => {
5060 Arc::new(change_underline_null_values_for_byte_array(array.as_string::<i32>()))
5061 }
5062 DataType::LargeUtf8 => {
5063 Arc::new(change_underline_null_values_for_byte_array(array.as_string::<i64>()))
5064 }
5065 DataType::Binary => {
5066 Arc::new(change_underline_null_values_for_byte_array(array.as_binary::<i32>()))
5067 }
5068 DataType::LargeBinary => {
5069 Arc::new(change_underline_null_values_for_byte_array(array.as_binary::<i64>()))
5070 }
5071 DataType::List(_) => {
5072 Arc::new(change_underline_null_values_for_list_array(array.as_list::<i32>()))
5073 }
5074 DataType::LargeList(_) => {
5075 Arc::new(change_underline_null_values_for_list_array(array.as_list::<i64>()))
5076 }
5077 DataType::Map(_, _) => {
5078 Arc::new(change_underline_null_values_for_map_array(array.as_map()))
5079 }
5080 _ => {
5081 Arc::clone(array)
5082 }
5083 )
5084 }
5085
5086 fn generate_column(rng: &mut (impl RngCore + Clone), len: usize) -> ArrayRef {
5087 match rng.random_range(0..24) {
5088 0 => Arc::new(generate_primitive_array::<Int32Type>(rng, len, 0.8)),
5089 1 => Arc::new(generate_primitive_array::<UInt32Type>(rng, len, 0.8)),
5090 2 => Arc::new(generate_primitive_array::<Int64Type>(rng, len, 0.8)),
5091 3 => Arc::new(generate_primitive_array::<UInt64Type>(rng, len, 0.8)),
5092 4 => Arc::new(generate_primitive_array::<Float32Type>(rng, len, 0.8)),
5093 5 => Arc::new(generate_primitive_array::<Float64Type>(rng, len, 0.8)),
5094 6 => Arc::new(generate_strings::<i32>(rng, len, 0.8)),
5095 7 => {
5096 let dict_values_len = rng.random_range(1..len);
5097 let strings = Arc::new(generate_strings::<i32>(rng, dict_values_len, 1.0));
5099 Arc::new(generate_dictionary::<Int64Type>(rng, strings, len, 0.8))
5100 }
5101 8 => {
5102 let dict_values_len = rng.random_range(1..len);
5103 let values = Arc::new(generate_primitive_array::<Int64Type>(
5105 rng,
5106 dict_values_len,
5107 1.0,
5108 ));
5109 Arc::new(generate_dictionary::<Int64Type>(rng, values, len, 0.8))
5110 }
5111 9 => Arc::new(generate_fixed_size_binary(rng, len, 0.8)),
5112 10 => Arc::new(generate_struct(rng, len, 0.8)),
5113 11 => Arc::new(generate_list(rng, len, 0.8, |rng, values_len| {
5114 Arc::new(generate_primitive_array::<Int64Type>(rng, values_len, 0.8))
5115 })),
5116 12 => Arc::new(generate_list(rng, len, 0.8, |rng, values_len| {
5117 Arc::new(generate_strings::<i32>(rng, values_len, 0.8))
5118 })),
5119 13 => Arc::new(generate_list(rng, len, 0.8, |rng, values_len| {
5120 Arc::new(generate_struct(rng, values_len, 0.8))
5121 })),
5122 14 => Arc::new(generate_string_view(rng, len, 0.8)),
5123 15 => Arc::new(generate_byte_view(rng, len, 0.8)),
5124 16 => Arc::new(generate_fixed_stringview_column(len)),
5125 17 => Arc::new(
5126 generate_list(&mut rng.clone(), len + 1000, 0.8, |rng, values_len| {
5127 Arc::new(generate_primitive_array::<Int64Type>(rng, values_len, 0.8))
5128 })
5129 .slice(500, len),
5130 ),
5131 18 => Arc::new(generate_boolean_array(rng, len, 0.8)),
5132 19 => Arc::new(generate_list_view(
5133 &mut rng.clone(),
5134 len,
5135 0.8,
5136 |values_len| Arc::new(generate_primitive_array::<Int64Type>(rng, values_len, 0.8)),
5137 )),
5138 20 => Arc::new(generate_list_view(
5139 &mut rng.clone(),
5140 len,
5141 0.8,
5142 |values_len| Arc::new(generate_strings::<i32>(rng, values_len, 0.8)),
5143 )),
5144 21 => Arc::new(generate_list_view(
5145 &mut rng.clone(),
5146 len,
5147 0.8,
5148 |values_len| Arc::new(generate_struct(rng, values_len, 0.8)),
5149 )),
5150 22 => Arc::new(
5151 generate_list_view(&mut rng.clone(), len + 1000, 0.8, |values_len| {
5152 Arc::new(generate_primitive_array::<Int64Type>(rng, values_len, 0.8))
5153 })
5154 .slice(500, len),
5155 ),
5156 23 => Arc::new(generate_map(
5157 rng,
5158 len,
5159 0.9,
5160 |rng, keys_len| {
5162 Arc::new(generate_all_unique_primitive_array::<Int64Type>(
5163 rng, keys_len,
5164 ))
5165 },
5166 |rng, values_len| Arc::new(generate_strings::<i32>(rng, values_len, 0.7)),
5167 )),
5168 _ => unreachable!(),
5169 }
5170 }
5171
5172 fn print_row(cols: &[SortColumn], row: usize) -> String {
5173 let t: Vec<_> = cols
5174 .iter()
5175 .map(|x| match x.values.is_valid(row) {
5176 true => {
5177 let opts = FormatOptions::default().with_null("NULL");
5178 let formatter = ArrayFormatter::try_new(x.values.as_ref(), &opts).unwrap();
5179 formatter.value(row).to_string()
5180 }
5181 false => "NULL".to_string(),
5182 })
5183 .collect();
5184 t.join(",")
5185 }
5186
5187 fn print_col_types(cols: &[SortColumn]) -> String {
5188 let t: Vec<_> = cols
5189 .iter()
5190 .map(|x| x.values.data_type().to_string())
5191 .collect();
5192 t.join(",")
5193 }
5194
5195 #[derive(Debug, PartialEq)]
5196 enum Nulls {
5197 AsIs,
5199
5200 Different,
5202
5203 None,
5205 }
5206
5207 #[test]
5208 #[cfg_attr(miri, ignore)]
5209 fn fuzz_test() {
5210 let mut rng = StdRng::seed_from_u64(42);
5211 for _ in 0..100 {
5212 for null_behavior in [Nulls::AsIs, Nulls::Different, Nulls::None] {
5213 let num_columns = rng.random_range(1..5);
5214 let len = rng.random_range(5..100);
5215 let mut arrays: Vec<_> = (0..num_columns)
5216 .map(|_| generate_column(&mut rng, len))
5217 .collect();
5218
5219 match null_behavior {
5220 Nulls::AsIs => {
5221 }
5223 Nulls::Different => {
5224 arrays = arrays
5226 .into_iter()
5227 .map(|a| replace_array_nulls(a, generate_nulls(&mut rng, len)))
5228 .collect()
5229 }
5230 Nulls::None => {
5231 arrays = arrays
5233 .into_iter()
5234 .map(|a| replace_array_nulls(a, None))
5235 .collect()
5236 }
5237 }
5238
5239 let options: Vec<_> = (0..num_columns)
5240 .map(|_| SortOptions {
5241 descending: rng.random_bool(0.5),
5242 nulls_first: rng.random_bool(0.5),
5243 })
5244 .collect();
5245
5246 let sort_columns: Vec<_> = options
5247 .iter()
5248 .zip(&arrays)
5249 .map(|(o, c)| SortColumn {
5250 values: Arc::clone(c),
5251 options: Some(*o),
5252 })
5253 .collect();
5254
5255 let comparator = LexicographicalComparator::try_new(&sort_columns).unwrap();
5256
5257 let columns: Vec<SortField> = options
5258 .into_iter()
5259 .zip(&arrays)
5260 .map(|(o, a)| SortField::new_with_options(a.data_type().clone(), o))
5261 .collect();
5262
5263 let converter = RowConverter::new(columns).unwrap();
5264 let rows = converter.convert_columns(&arrays).unwrap();
5265
5266 if !matches!(null_behavior, Nulls::None) {
5269 assert_same_rows_when_changing_input_underlying_null_values(
5270 &arrays, &converter, &rows,
5271 );
5272 }
5273
5274 for i in 0..len {
5275 for j in 0..len {
5276 let row_i = rows.row(i);
5277 let row_j = rows.row(j);
5278 let row_cmp = row_i.cmp(&row_j);
5279 let lex_cmp = comparator.compare(i, j);
5280 assert_eq!(
5281 row_cmp,
5282 lex_cmp,
5283 "({:?} vs {:?}) vs ({:?} vs {:?}) for types {}",
5284 print_row(&sort_columns, i),
5285 print_row(&sort_columns, j),
5286 row_i,
5287 row_j,
5288 print_col_types(&sort_columns)
5289 );
5290 }
5291 }
5292
5293 {
5295 let mut rows_iter = rows.iter();
5296 let mut rows_lengths_iter = rows.lengths();
5297 for (index, row) in rows_iter.by_ref().enumerate() {
5298 let len = rows_lengths_iter
5299 .next()
5300 .expect("Reached end of length iterator while still have rows");
5301 assert_eq!(
5302 row.data.len(),
5303 len,
5304 "Row length mismatch: {} vs {}",
5305 row.data.len(),
5306 len
5307 );
5308 assert_eq!(
5309 len,
5310 rows.row_len(index),
5311 "Row length mismatch at index {}: {} vs {}",
5312 index,
5313 len,
5314 rows.row_len(index)
5315 );
5316 }
5317
5318 assert_eq!(
5319 rows_lengths_iter.next(),
5320 None,
5321 "Length iterator did not reach end"
5322 );
5323 }
5324
5325 let back = converter.convert_rows(&rows).unwrap();
5328 for (actual, expected) in back.iter().zip(&arrays) {
5329 actual.to_data().validate_full().unwrap();
5330 dictionary_eq(actual, expected)
5331 }
5332
5333 let rows = rows.try_into_binary().expect("reasonable size");
5336 let parser = converter.parser();
5337 let back = converter
5338 .convert_rows(rows.iter().map(|b| parser.parse(b.expect("valid bytes"))))
5339 .unwrap();
5340 for (actual, expected) in back.iter().zip(&arrays) {
5341 actual.to_data().validate_full().unwrap();
5342 dictionary_eq(actual, expected)
5343 }
5344
5345 let rows = converter.from_binary(rows);
5346 let back = converter.convert_rows(&rows).unwrap();
5347 for (actual, expected) in back.iter().zip(&arrays) {
5348 actual.to_data().validate_full().unwrap();
5349 dictionary_eq(actual, expected)
5350 }
5351 }
5352 }
5353 }
5354
5355 fn replace_array_nulls(array: ArrayRef, new_nulls: Option<NullBuffer>) -> ArrayRef {
5356 make_array(
5357 array
5358 .into_data()
5359 .into_builder()
5360 .nulls(new_nulls)
5362 .build()
5363 .unwrap(),
5364 )
5365 }
5366
5367 fn assert_same_rows_when_changing_input_underlying_null_values(
5368 arrays: &[ArrayRef],
5369 converter: &RowConverter,
5370 rows: &Rows,
5371 ) {
5372 let arrays_with_different_data_behind_nulls = arrays
5373 .iter()
5374 .map(|arr| change_underline_null_values(arr))
5375 .collect::<Vec<_>>();
5376
5377 if arrays
5379 .iter()
5380 .zip(arrays_with_different_data_behind_nulls.iter())
5381 .all(|(a, b)| Arc::ptr_eq(a, b))
5382 {
5383 return;
5384 }
5385
5386 let rows_with_different_nulls = converter
5387 .convert_columns(&arrays_with_different_data_behind_nulls)
5388 .unwrap();
5389
5390 assert_eq!(
5391 rows.iter().collect::<Vec<_>>(),
5392 rows_with_different_nulls.iter().collect::<Vec<_>>(),
5393 "Different underlying nulls should not output different rows"
5394 )
5395 }
5396
5397 #[test]
5398 fn test_clear() {
5399 let converter = RowConverter::new(vec![SortField::new(DataType::Int32)]).unwrap();
5400 let mut rows = converter.empty_rows(3, 128);
5401
5402 let first = Int32Array::from(vec![None, Some(2), Some(4)]);
5403 let second = Int32Array::from(vec![Some(2), None, Some(4)]);
5404 let arrays = [Arc::new(first) as ArrayRef, Arc::new(second) as ArrayRef];
5405
5406 for array in arrays.iter() {
5407 rows.clear();
5408 converter
5409 .append(&mut rows, std::slice::from_ref(array))
5410 .unwrap();
5411 let back = converter.convert_rows(&rows).unwrap();
5412 assert_eq!(&back[0], array);
5413 }
5414
5415 let mut rows_expected = converter.empty_rows(3, 128);
5416 converter.append(&mut rows_expected, &arrays[1..]).unwrap();
5417
5418 for (i, (actual, expected)) in rows.iter().zip(rows_expected.iter()).enumerate() {
5419 assert_eq!(
5420 actual, expected,
5421 "For row {i}: expected {expected:?}, actual: {actual:?}",
5422 );
5423 }
5424 }
5425
5426 #[test]
5427 fn test_append_codec_dictionary_binary() {
5428 use DataType::*;
5429 let converter = RowConverter::new(vec![SortField::new(Dictionary(
5431 Box::new(Int32),
5432 Box::new(Binary),
5433 ))])
5434 .unwrap();
5435 let mut rows = converter.empty_rows(4, 128);
5436
5437 let keys = Int32Array::from_iter_values([0, 1, 2, 3]);
5438 let values = BinaryArray::from(vec![
5439 Some("a".as_bytes()),
5440 Some(b"b"),
5441 Some(b"c"),
5442 Some(b"d"),
5443 ]);
5444 let dict_array = DictionaryArray::new(keys, Arc::new(values));
5445
5446 rows.clear();
5447 let array = Arc::new(dict_array) as ArrayRef;
5448 converter
5449 .append(&mut rows, std::slice::from_ref(&array))
5450 .unwrap();
5451 let back = converter.convert_rows(&rows).unwrap();
5452
5453 dictionary_eq(&back[0], &array);
5454 }
5455
5456 #[test]
5457 fn test_list_prefix() {
5458 let mut a = ListBuilder::new(Int8Builder::new());
5459 a.append_value([None]);
5460 a.append_value([None, None]);
5461 let a = a.finish();
5462
5463 let converter = RowConverter::new(vec![SortField::new(a.data_type().clone())]).unwrap();
5464 let rows = converter.convert_columns(&[Arc::new(a) as _]).unwrap();
5465 assert_eq!(rows.row(0).cmp(&rows.row(1)), Ordering::Less);
5466 }
5467
5468 #[test]
5469 fn test_utf8_validation_doesnt_affect_values_buffer_size() {
5470 fn assert_values_buffer_lens(col: ArrayRef) -> usize {
5471 let converter = RowConverter::new(vec![SortField::new(DataType::Utf8View)]).unwrap();
5473
5474 let rows = converter.convert_columns(&[col]).unwrap();
5476 let converted = converter.convert_rows(&rows).unwrap();
5477 let unchecked_values_len = converted[0].as_string_view().data_buffers()[0].len();
5478
5479 let rows = rows.try_into_binary().expect("reasonable size");
5481 let parser = converter.parser();
5482 let converted = converter
5483 .convert_rows(rows.iter().map(|b| parser.parse(b.expect("valid bytes"))))
5484 .unwrap();
5485 let checked_values_len = converted[0].as_string_view().data_buffers()[0].len();
5486 assert_eq!(unchecked_values_len, checked_values_len);
5488 checked_values_len
5489 }
5490
5491 let col = Arc::new(StringViewArray::from_iter([
5493 Some("hello"), None, Some("short"), Some("tiny"), ])) as ArrayRef;
5498
5499 let values_len = assert_values_buffer_lens(col);
5500 assert_eq!(values_len, 0);
5502
5503 let col = Arc::new(StringViewArray::from_iter([
5505 Some("1234567890123"), Some("12345678901234"), ])) as ArrayRef;
5508
5509 let values_len = assert_values_buffer_lens(col);
5510 assert_eq!(values_len, 13 + 14);
5511
5512 let col = Arc::new(StringViewArray::from_iter([
5514 Some("tiny"), Some("thisisexact13"), None,
5517 Some("short"), ])) as ArrayRef;
5519
5520 let values_len = assert_values_buffer_lens(col);
5521 assert_eq!(values_len, 13);
5523 }
5524
5525 #[test]
5526 fn test_sparse_union() {
5527 let int_array = Int32Array::from(vec![Some(1), None, Some(3), None, Some(5)]);
5529 let str_array = StringArray::from(vec![None, Some("b"), None, Some("d"), None]);
5530
5531 let type_ids = vec![0, 1, 0, 1, 0].into();
5533
5534 let union_fields = [
5535 (0, Arc::new(Field::new("int", DataType::Int32, false))),
5536 (1, Arc::new(Field::new("str", DataType::Utf8, false))),
5537 ]
5538 .into_iter()
5539 .collect();
5540
5541 let union_array = UnionArray::try_new(
5542 union_fields,
5543 type_ids,
5544 None,
5545 vec![Arc::new(int_array) as ArrayRef, Arc::new(str_array)],
5546 )
5547 .unwrap();
5548
5549 let union_type = union_array.data_type().clone();
5550 let converter = RowConverter::new(vec![SortField::new(union_type)]).unwrap();
5551
5552 let rows = converter
5553 .convert_columns(&[Arc::new(union_array.clone())])
5554 .unwrap();
5555
5556 let back = converter.convert_rows(&rows).unwrap();
5558 let back_union = back[0].as_any().downcast_ref::<UnionArray>().unwrap();
5559
5560 assert_eq!(union_array.len(), back_union.len());
5561 for i in 0..union_array.len() {
5562 assert_eq!(union_array.type_id(i), back_union.type_id(i));
5563 }
5564 }
5565
5566 #[test]
5567 fn test_sparse_union_with_nulls() {
5568 let int_array = Int32Array::from(vec![Some(1), None, Some(3), None, Some(5)]);
5570 let str_array = StringArray::from(vec![None::<&str>; 5]);
5571
5572 let type_ids = vec![0, 1, 0, 1, 0].into();
5574
5575 let union_fields = [
5576 (0, Arc::new(Field::new("int", DataType::Int32, true))),
5577 (1, Arc::new(Field::new("str", DataType::Utf8, true))),
5578 ]
5579 .into_iter()
5580 .collect();
5581
5582 let union_array = UnionArray::try_new(
5583 union_fields,
5584 type_ids,
5585 None,
5586 vec![Arc::new(int_array) as ArrayRef, Arc::new(str_array)],
5587 )
5588 .unwrap();
5589
5590 let union_type = union_array.data_type().clone();
5591 let converter = RowConverter::new(vec![SortField::new(union_type)]).unwrap();
5592
5593 let rows = converter
5594 .convert_columns(&[Arc::new(union_array.clone())])
5595 .unwrap();
5596
5597 let back = converter.convert_rows(&rows).unwrap();
5599 let back_union = back[0].as_any().downcast_ref::<UnionArray>().unwrap();
5600
5601 assert_eq!(union_array.len(), back_union.len());
5602 for i in 0..union_array.len() {
5603 let expected_null = union_array.is_null(i);
5604 let actual_null = back_union.is_null(i);
5605 assert_eq!(expected_null, actual_null, "Null mismatch at index {i}");
5606 if !expected_null {
5607 assert_eq!(union_array.type_id(i), back_union.type_id(i));
5608 }
5609 }
5610 }
5611
5612 #[test]
5613 fn test_dense_union() {
5614 let int_array = Int32Array::from(vec![1, 3, 5]);
5616 let str_array = StringArray::from(vec!["a", "b"]);
5617
5618 let type_ids = vec![0, 1, 0, 1, 0].into();
5619
5620 let offsets = vec![0, 0, 1, 1, 2].into();
5622
5623 let union_fields = [
5624 (0, Arc::new(Field::new("int", DataType::Int32, false))),
5625 (1, Arc::new(Field::new("str", DataType::Utf8, false))),
5626 ]
5627 .into_iter()
5628 .collect();
5629
5630 let union_array = UnionArray::try_new(
5631 union_fields,
5632 type_ids,
5633 Some(offsets), vec![Arc::new(int_array) as ArrayRef, Arc::new(str_array)],
5635 )
5636 .unwrap();
5637
5638 let union_type = union_array.data_type().clone();
5639 let converter = RowConverter::new(vec![SortField::new(union_type)]).unwrap();
5640
5641 let rows = converter
5642 .convert_columns(&[Arc::new(union_array.clone())])
5643 .unwrap();
5644
5645 let back = converter.convert_rows(&rows).unwrap();
5647 let back_union = back[0].as_any().downcast_ref::<UnionArray>().unwrap();
5648
5649 assert_eq!(union_array.len(), back_union.len());
5650 for i in 0..union_array.len() {
5651 assert_eq!(union_array.type_id(i), back_union.type_id(i));
5652 }
5653 }
5654
5655 #[test]
5656 fn test_dense_union_with_nulls() {
5657 let int_array = Int32Array::from(vec![Some(1), None, Some(5)]);
5659 let str_array = StringArray::from(vec![Some("a"), None]);
5660
5661 let type_ids = vec![0, 1, 0, 1, 0].into();
5663 let offsets = vec![0, 0, 1, 1, 2].into();
5664
5665 let union_fields = [
5666 (0, Arc::new(Field::new("int", DataType::Int32, true))),
5667 (1, Arc::new(Field::new("str", DataType::Utf8, true))),
5668 ]
5669 .into_iter()
5670 .collect();
5671
5672 let union_array = UnionArray::try_new(
5673 union_fields,
5674 type_ids,
5675 Some(offsets),
5676 vec![Arc::new(int_array) as ArrayRef, Arc::new(str_array)],
5677 )
5678 .unwrap();
5679
5680 let union_type = union_array.data_type().clone();
5681 let converter = RowConverter::new(vec![SortField::new(union_type)]).unwrap();
5682
5683 let rows = converter
5684 .convert_columns(&[Arc::new(union_array.clone())])
5685 .unwrap();
5686
5687 let back = converter.convert_rows(&rows).unwrap();
5689 let back_union = back[0].as_any().downcast_ref::<UnionArray>().unwrap();
5690
5691 assert_eq!(union_array.len(), back_union.len());
5692 for i in 0..union_array.len() {
5693 let expected_null = union_array.is_null(i);
5694 let actual_null = back_union.is_null(i);
5695 assert_eq!(expected_null, actual_null, "Null mismatch at index {i}");
5696 if !expected_null {
5697 assert_eq!(union_array.type_id(i), back_union.type_id(i));
5698 }
5699 }
5700 }
5701
5702 #[test]
5703 fn test_union_ordering() {
5704 let int_array = Int32Array::from(vec![100, 5, 20]);
5705 let str_array = StringArray::from(vec!["z", "a"]);
5706
5707 let type_ids = vec![0, 1, 0, 1, 0].into();
5709 let offsets = vec![0, 0, 1, 1, 2].into();
5710
5711 let union_fields = [
5712 (0, Arc::new(Field::new("int", DataType::Int32, false))),
5713 (1, Arc::new(Field::new("str", DataType::Utf8, false))),
5714 ]
5715 .into_iter()
5716 .collect();
5717
5718 let union_array = UnionArray::try_new(
5719 union_fields,
5720 type_ids,
5721 Some(offsets),
5722 vec![Arc::new(int_array) as ArrayRef, Arc::new(str_array)],
5723 )
5724 .unwrap();
5725
5726 let union_type = union_array.data_type().clone();
5727 let converter = RowConverter::new(vec![SortField::new(union_type)]).unwrap();
5728
5729 let rows = converter.convert_columns(&[Arc::new(union_array)]).unwrap();
5730
5731 assert!(rows.row(2) < rows.row(1));
5743
5744 assert!(rows.row(0) < rows.row(3));
5746
5747 assert!(rows.row(2) < rows.row(4));
5750 assert!(rows.row(4) < rows.row(0));
5752
5753 assert!(rows.row(3) < rows.row(1));
5756 }
5757
5758 #[test]
5759 fn test_row_converter_roundtrip_with_many_union_columns() {
5760 let fields1 = UnionFields::try_new(
5762 vec![0, 1],
5763 vec![
5764 Field::new("int", DataType::Int32, true),
5765 Field::new("string", DataType::Utf8, true),
5766 ],
5767 )
5768 .unwrap();
5769
5770 let int_array1 = Int32Array::from(vec![Some(67), None]);
5771 let string_array1 = StringArray::from(vec![None::<&str>, Some("hello")]);
5772 let type_ids1 = vec![0i8, 1].into();
5773
5774 let union_array1 = UnionArray::try_new(
5775 fields1.clone(),
5776 type_ids1,
5777 None,
5778 vec![
5779 Arc::new(int_array1) as ArrayRef,
5780 Arc::new(string_array1) as ArrayRef,
5781 ],
5782 )
5783 .unwrap();
5784
5785 let fields2 = UnionFields::try_new(
5787 vec![0, 1],
5788 vec![
5789 Field::new("int", DataType::Int32, true),
5790 Field::new("string", DataType::Utf8, true),
5791 ],
5792 )
5793 .unwrap();
5794
5795 let int_array2 = Int32Array::from(vec![Some(100), None]);
5796 let string_array2 = StringArray::from(vec![None::<&str>, Some("world")]);
5797 let type_ids2 = vec![0i8, 1].into();
5798
5799 let union_array2 = UnionArray::try_new(
5800 fields2.clone(),
5801 type_ids2,
5802 None,
5803 vec![
5804 Arc::new(int_array2) as ArrayRef,
5805 Arc::new(string_array2) as ArrayRef,
5806 ],
5807 )
5808 .unwrap();
5809
5810 let field1 = Field::new("col1", DataType::Union(fields1, UnionMode::Sparse), true);
5812 let field2 = Field::new("col2", DataType::Union(fields2, UnionMode::Sparse), true);
5813
5814 let sort_field1 = SortField::new(field1.data_type().clone());
5815 let sort_field2 = SortField::new(field2.data_type().clone());
5816
5817 let converter = RowConverter::new(vec![sort_field1, sort_field2]).unwrap();
5818
5819 let rows = converter
5820 .convert_columns(&[
5821 Arc::new(union_array1.clone()) as ArrayRef,
5822 Arc::new(union_array2.clone()) as ArrayRef,
5823 ])
5824 .unwrap();
5825
5826 let out = converter.convert_rows(&rows).unwrap();
5828
5829 let [col1, col2] = out.as_slice() else {
5830 panic!("expected 2 columns")
5831 };
5832
5833 let col1 = col1.as_any().downcast_ref::<UnionArray>().unwrap();
5834 let col2 = col2.as_any().downcast_ref::<UnionArray>().unwrap();
5835
5836 for (expected, got) in [union_array1, union_array2].iter().zip([col1, col2]) {
5837 assert_eq!(expected.len(), got.len());
5838 assert_eq!(expected.type_ids(), got.type_ids());
5839
5840 for i in 0..expected.len() {
5841 assert_eq!(expected.value(i).as_ref(), got.value(i).as_ref());
5842 }
5843 }
5844 }
5845
5846 #[test]
5847 fn test_row_converter_roundtrip_with_one_union_column() {
5848 let fields = UnionFields::try_new(
5849 vec![0, 1],
5850 vec![
5851 Field::new("int", DataType::Int32, true),
5852 Field::new("string", DataType::Utf8, true),
5853 ],
5854 )
5855 .unwrap();
5856
5857 let int_array = Int32Array::from(vec![Some(67), None]);
5858 let string_array = StringArray::from(vec![None::<&str>, Some("hello")]);
5859 let type_ids = vec![0i8, 1].into();
5860
5861 let union_array = UnionArray::try_new(
5862 fields.clone(),
5863 type_ids,
5864 None,
5865 vec![
5866 Arc::new(int_array) as ArrayRef,
5867 Arc::new(string_array) as ArrayRef,
5868 ],
5869 )
5870 .unwrap();
5871
5872 let field = Field::new("col", DataType::Union(fields, UnionMode::Sparse), true);
5873 let sort_field = SortField::new(field.data_type().clone());
5874 let converter = RowConverter::new(vec![sort_field]).unwrap();
5875
5876 let rows = converter
5877 .convert_columns(&[Arc::new(union_array.clone()) as ArrayRef])
5878 .unwrap();
5879
5880 let out = converter.convert_rows(&rows).unwrap();
5882
5883 let [col1] = out.as_slice() else {
5884 panic!("expected 1 column")
5885 };
5886
5887 let col = col1.as_any().downcast_ref::<UnionArray>().unwrap();
5888 assert_eq!(col.len(), union_array.len());
5889 assert_eq!(col.type_ids(), union_array.type_ids());
5890
5891 for i in 0..col.len() {
5892 assert_eq!(col.value(i).as_ref(), union_array.value(i).as_ref());
5893 }
5894 }
5895
5896 #[test]
5897 fn test_row_converter_roundtrip_with_non_default_union_type_ids() {
5898 let fields = UnionFields::try_new(
5900 vec![70, 85],
5901 vec![
5902 Field::new("int", DataType::Int32, true),
5903 Field::new("string", DataType::Utf8, true),
5904 ],
5905 )
5906 .unwrap();
5907
5908 let int_array = Int32Array::from(vec![Some(67), None]);
5909 let string_array = StringArray::from(vec![None::<&str>, Some("hello")]);
5910 let type_ids = vec![70i8, 85].into();
5911
5912 let union_array = UnionArray::try_new(
5913 fields.clone(),
5914 type_ids,
5915 None,
5916 vec![
5917 Arc::new(int_array) as ArrayRef,
5918 Arc::new(string_array) as ArrayRef,
5919 ],
5920 )
5921 .unwrap();
5922
5923 let field = Field::new("col", DataType::Union(fields, UnionMode::Sparse), true);
5924 let sort_field = SortField::new(field.data_type().clone());
5925 let converter = RowConverter::new(vec![sort_field]).unwrap();
5926
5927 let rows = converter
5928 .convert_columns(&[Arc::new(union_array.clone()) as ArrayRef])
5929 .unwrap();
5930
5931 let out = converter.convert_rows(&rows).unwrap();
5933
5934 let [col1] = out.as_slice() else {
5935 panic!("expected 1 column")
5936 };
5937
5938 let col = col1.as_any().downcast_ref::<UnionArray>().unwrap();
5939 assert_eq!(col.len(), union_array.len());
5940 assert_eq!(col.type_ids(), union_array.type_ids());
5941
5942 for i in 0..col.len() {
5943 assert_eq!(col.value(i).as_ref(), union_array.value(i).as_ref());
5944 }
5945 }
5946
5947 #[test]
5948 fn rows_size_should_count_for_capacity() {
5949 let row_converter = RowConverter::new(vec![SortField::new(DataType::UInt8)]).unwrap();
5950
5951 let empty_rows_size_with_preallocate_rows_and_data = {
5952 let rows = row_converter.empty_rows(1000, 1000);
5953
5954 rows.size()
5955 };
5956 let empty_rows_size_with_preallocate_rows = {
5957 let rows = row_converter.empty_rows(1000, 0);
5958
5959 rows.size()
5960 };
5961 let empty_rows_size_with_preallocate_data = {
5962 let rows = row_converter.empty_rows(0, 1000);
5963
5964 rows.size()
5965 };
5966 let empty_rows_size_without_preallocate = {
5967 let rows = row_converter.empty_rows(0, 0);
5968
5969 rows.size()
5970 };
5971
5972 assert!(
5973 empty_rows_size_with_preallocate_rows_and_data > empty_rows_size_with_preallocate_rows,
5974 "{empty_rows_size_with_preallocate_rows_and_data} should be larger than {empty_rows_size_with_preallocate_rows}"
5975 );
5976 assert!(
5977 empty_rows_size_with_preallocate_rows_and_data > empty_rows_size_with_preallocate_data,
5978 "{empty_rows_size_with_preallocate_rows_and_data} should be larger than {empty_rows_size_with_preallocate_data}"
5979 );
5980 assert!(
5981 empty_rows_size_with_preallocate_rows > empty_rows_size_without_preallocate,
5982 "{empty_rows_size_with_preallocate_rows} should be larger than {empty_rows_size_without_preallocate}"
5983 );
5984 assert!(
5985 empty_rows_size_with_preallocate_data > empty_rows_size_without_preallocate,
5986 "{empty_rows_size_with_preallocate_data} should be larger than {empty_rows_size_without_preallocate}"
5987 );
5988 }
5989
5990 #[test]
5991 fn test_struct_no_child_fields() {
5992 fn run_test(array: ArrayRef) {
5993 let sort_fields = vec![SortField::new(array.data_type().clone())];
5994 let converter = RowConverter::new(sort_fields).unwrap();
5995 let r = converter.convert_columns(&[Arc::clone(&array)]).unwrap();
5996
5997 let back = converter.convert_rows(&r).unwrap();
5998 assert_eq!(back.len(), 1);
5999 assert_eq!(&back[0], &array);
6000 }
6001
6002 let s = Arc::new(StructArray::new_empty_fields(5, None)) as ArrayRef;
6003 run_test(s);
6004
6005 let s = Arc::new(StructArray::new_empty_fields(
6006 5,
6007 Some(vec![true, false, true, false, false].into()),
6008 )) as ArrayRef;
6009 run_test(s);
6010 }
6011
6012 #[test]
6013 fn reserve_should_increase_capacity_to_the_requested_size() {
6014 let row_converter = RowConverter::new(vec![SortField::new(DataType::UInt8)]).unwrap();
6015 let mut empty_rows = row_converter.empty_rows(0, 0);
6016 empty_rows.reserve(50, 50);
6017 let before_size = empty_rows.size();
6018 empty_rows.reserve(50, 50);
6019 assert_eq!(
6020 empty_rows.size(),
6021 before_size,
6022 "Size should not change when reserving already reserved space"
6023 );
6024 empty_rows.reserve(10, 20);
6025 assert_eq!(
6026 empty_rows.size(),
6027 before_size,
6028 "Size should not change when already have space for the expected reserved data"
6029 );
6030
6031 empty_rows.reserve(100, 20);
6032 assert!(
6033 empty_rows.size() > before_size,
6034 "Size should increase when reserving more space than previously reserved"
6035 );
6036
6037 let before_size = empty_rows.size();
6038
6039 empty_rows.reserve(20, 100);
6040 assert!(
6041 empty_rows.size() > before_size,
6042 "Size should increase when reserving more space than previously reserved"
6043 );
6044 }
6045
6046 #[test]
6047 fn empty_rows_should_return_empty_lengths_iterator() {
6048 let rows = RowConverter::new(vec![SortField::new(DataType::UInt8)])
6049 .unwrap()
6050 .empty_rows(0, 0);
6051 let mut lengths_iter = rows.lengths();
6052 assert_eq!(lengths_iter.next(), None);
6053 }
6054
6055 #[test]
6056 #[should_panic(expected = "row index out of bounds")]
6057 fn row_should_panic_on_overflowing_index() {
6058 let rows = RowConverter::new(vec![SortField::new(DataType::Int32)])
6059 .unwrap()
6060 .empty_rows(0, 0);
6061 rows.row(usize::MAX);
6062 }
6063
6064 #[test]
6065 #[should_panic(expected = "row index out of bounds")]
6066 fn row_len_should_panic_on_overflowing_index() {
6067 let rows = RowConverter::new(vec![SortField::new(DataType::Int32)])
6068 .unwrap()
6069 .empty_rows(0, 0);
6070 rows.row_len(usize::MAX);
6071 }
6072
6073 #[test]
6074 fn test_nested_null_list() {
6075 let null_array = Arc::new(NullArray::new(3));
6076 let list: ArrayRef = Arc::new(ListArray::new(
6078 Field::new_list_field(DataType::Null, true).into(),
6079 OffsetBuffer::from_lengths(vec![1, 0, 2]),
6080 null_array,
6081 None,
6082 ));
6083
6084 let converter = RowConverter::new(vec![SortField::new(list.data_type().clone())]).unwrap();
6085 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
6086 let back = converter.convert_rows(&rows).unwrap();
6087
6088 assert_eq!(&list, &back[0]);
6089 }
6090
6091 #[test]
6093 fn test_double_nested_null_list() {
6094 let null_array = Arc::new(NullArray::new(1));
6095 let nested_field = Arc::new(Field::new_list_field(DataType::Null, true));
6097 let nested_list = Arc::new(ListArray::new(
6098 nested_field.clone(),
6099 OffsetBuffer::from_lengths(vec![1]),
6100 null_array,
6101 None,
6102 ));
6103 let list = Arc::new(ListArray::new(
6105 Field::new_list_field(DataType::List(nested_field), true).into(),
6106 OffsetBuffer::from_lengths(vec![1]),
6107 nested_list,
6108 None,
6109 )) as ArrayRef;
6110
6111 let converter = RowConverter::new(vec![SortField::new(list.data_type().clone())]).unwrap();
6112 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
6113 let back = converter.convert_rows(&rows).unwrap();
6114
6115 assert_eq!(&list, &back[0]);
6116 }
6117
6118 #[test]
6120 fn test_large_list_null() {
6121 let null_array = Arc::new(NullArray::new(3));
6122 let list: ArrayRef = Arc::new(LargeListArray::new(
6124 Field::new_list_field(DataType::Null, true).into(),
6125 OffsetBuffer::from_lengths(vec![1, 0, 2]),
6126 null_array,
6127 None,
6128 ));
6129
6130 let converter = RowConverter::new(vec![SortField::new(list.data_type().clone())]).unwrap();
6131 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
6132 let back = converter.convert_rows(&rows).unwrap();
6133
6134 assert_eq!(&list, &back[0]);
6135 }
6136
6137 #[test]
6139 fn test_fixed_size_list_null() {
6140 let null_array = Arc::new(NullArray::new(6));
6141 let list: ArrayRef = Arc::new(FixedSizeListArray::new(
6143 Arc::new(Field::new_list_field(DataType::Null, true)),
6144 2,
6145 null_array,
6146 None,
6147 ));
6148
6149 let converter = RowConverter::new(vec![SortField::new(list.data_type().clone())]).unwrap();
6150 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
6151 let back = converter.convert_rows(&rows).unwrap();
6152
6153 assert_eq!(&list, &back[0]);
6154 }
6155
6156 #[test]
6158 fn test_fixed_size_list_of_dictionaries_round_trips() {
6159 let dict_dt = DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8));
6162 let element_field = Arc::new(Field::new("item", dict_dt.clone(), true));
6163 let fsl_dt = DataType::FixedSizeList(Arc::clone(&element_field), 2);
6164
6165 let values = Arc::new(StringArray::from(vec!["a", "b"]));
6166 let keys = Int32Array::from(vec![0, 1]);
6167 let dict = DictionaryArray::<Int32Type>::try_new(keys, values).unwrap();
6168 let fsl: ArrayRef = Arc::new(FixedSizeListArray::new(
6169 Arc::clone(&element_field),
6170 2,
6171 Arc::new(dict),
6172 None,
6173 ));
6174
6175 assert!(RowConverter::supports_fields(&[SortField::new(
6176 fsl_dt.clone()
6177 )]));
6178
6179 let converter = RowConverter::new(vec![SortField::new(fsl_dt.clone())]).unwrap();
6180 let rows = converter.convert_columns(&[Arc::clone(&fsl)]).unwrap();
6181
6182 let back = converter.convert_rows(&rows).unwrap();
6185 assert_eq!(back.len(), 1);
6186
6187 let out = back[0]
6191 .as_any()
6192 .downcast_ref::<FixedSizeListArray>()
6193 .expect("decoded array must be a FixedSizeListArray");
6194 assert_eq!(out.len(), 1);
6195 assert_eq!(out.value_length(), 2);
6196 assert_eq!(out.values().data_type(), &DataType::Utf8);
6200
6201 let values = out
6203 .values()
6204 .as_any()
6205 .downcast_ref::<StringArray>()
6206 .expect("child must be a StringArray after flattening");
6207 assert_eq!(values.value(0), "a");
6208 assert_eq!(values.value(1), "b");
6209 }
6210
6211 #[test]
6213 fn test_list_null_variations() {
6214 let null_array = Arc::new(NullArray::new(3));
6216 let list: ArrayRef = Arc::new(ListArray::new(
6217 Field::new_list_field(DataType::Null, true).into(),
6218 OffsetBuffer::from_lengths(vec![1, 0, 2]),
6219 null_array,
6220 None,
6221 ));
6222
6223 let converter = RowConverter::new(vec![SortField::new(list.data_type().clone())]).unwrap();
6224 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
6225 let back = converter.convert_rows(&rows).unwrap();
6226 assert_eq!(&list, &back[0]);
6227
6228 let null_array = Arc::new(NullArray::new(3));
6230 let list: ArrayRef = Arc::new(ListArray::new(
6231 Field::new_list_field(DataType::Null, true).into(),
6232 OffsetBuffer::from_lengths(vec![1, 0, 2]),
6233 null_array,
6234 Some(vec![true, false, true].into()),
6235 ));
6236
6237 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
6238 let back = converter.convert_rows(&rows).unwrap();
6239 assert_eq!(&list, &back[0]);
6240
6241 let null_array = Arc::new(NullArray::new(0));
6243 let list: ArrayRef = Arc::new(ListArray::new(
6244 Field::new_list_field(DataType::Null, true).into(),
6245 OffsetBuffer::from_lengths(vec![]),
6246 null_array,
6247 None,
6248 ));
6249
6250 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
6251 let back = converter.convert_rows(&rows).unwrap();
6252 assert_eq!(&list, &back[0]);
6253
6254 let null_array = Arc::new(NullArray::new(0));
6256 let list: ArrayRef = Arc::new(ListArray::new(
6257 Field::new_list_field(DataType::Null, true).into(),
6258 OffsetBuffer::from_lengths(vec![0, 0, 0]),
6259 null_array,
6260 None,
6261 ));
6262
6263 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
6264 let back = converter.convert_rows(&rows).unwrap();
6265 assert_eq!(&list, &back[0]);
6266 }
6267
6268 #[test]
6270 fn test_list_null_descending() {
6271 let null_array = Arc::new(NullArray::new(3));
6272 let list: ArrayRef = Arc::new(ListArray::new(
6274 Field::new_list_field(DataType::Null, true).into(),
6275 OffsetBuffer::from_lengths(vec![1, 0, 2]),
6276 null_array,
6277 None,
6278 ));
6279
6280 let options = SortOptions::default().with_descending(true);
6281 let field = SortField::new_with_options(list.data_type().clone(), options);
6282 let converter = RowConverter::new(vec![field]).unwrap();
6283 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
6284 let back = converter.convert_rows(&rows).unwrap();
6285
6286 assert_eq!(&list, &back[0]);
6287 }
6288
6289 #[test]
6291 fn test_struct_with_null_field() {
6292 let null_array = Arc::new(NullArray::new(3));
6294 let int_array = Arc::new(Int32Array::from(vec![1, 2, 3]));
6295
6296 let struct_array: ArrayRef = Arc::new(StructArray::new(
6297 vec![
6298 Arc::new(Field::new("a", DataType::Null, true)),
6299 Arc::new(Field::new("b", DataType::Int32, true)),
6300 ]
6301 .into(),
6302 vec![null_array, int_array],
6303 Some(vec![true, true, false].into()), ));
6305
6306 let converter =
6307 RowConverter::new(vec![SortField::new(struct_array.data_type().clone())]).unwrap();
6308 let rows = converter
6309 .convert_columns(&[Arc::clone(&struct_array)])
6310 .unwrap();
6311 let back = converter.convert_rows(&rows).unwrap();
6312
6313 assert_eq!(&struct_array, &back[0]);
6314 }
6315
6316 #[test]
6318 fn test_nested_struct_with_null() {
6319 let inner_null = Arc::new(NullArray::new(2));
6321 let inner_struct = Arc::new(StructArray::new(
6322 vec![Arc::new(Field::new("x", DataType::Null, true))].into(),
6323 vec![inner_null],
6324 None,
6325 ));
6326
6327 let y_array = Arc::new(Int32Array::from(vec![10, 20]));
6329 let outer_struct: ArrayRef = Arc::new(StructArray::new(
6330 vec![
6331 Arc::new(Field::new("inner", inner_struct.data_type().clone(), true)),
6332 Arc::new(Field::new("y", DataType::Int32, true)),
6333 ]
6334 .into(),
6335 vec![inner_struct, y_array],
6336 None,
6337 ));
6338
6339 let converter =
6340 RowConverter::new(vec![SortField::new(outer_struct.data_type().clone())]).unwrap();
6341 let rows = converter
6342 .convert_columns(&[Arc::clone(&outer_struct)])
6343 .unwrap();
6344 let back = converter.convert_rows(&rows).unwrap();
6345
6346 assert_eq!(&outer_struct, &back[0]);
6347 }
6348
6349 #[test]
6351 fn test_map_null_variations() {
6352 let keys = Arc::new(StringArray::from(vec!["a", "b", "c"])) as ArrayRef;
6354 let null_values = Arc::new(NullArray::new(3)) as ArrayRef;
6355
6356 let offsets = OffsetBuffer::new(vec![0, 1, 1, 3].into());
6357 let entries_fields = vec![
6358 Arc::new(Field::new("keys", DataType::Utf8, false)),
6359 Arc::new(Field::new("values", DataType::Null, true)),
6360 ];
6361 let struct_field = Arc::new(Field::new(
6362 "entries",
6363 DataType::Struct(entries_fields.clone().into()),
6364 false,
6365 ));
6366 let entries = StructArray::new(entries_fields.into(), vec![keys, null_values], None);
6367
6368 let map: ArrayRef = Arc::new(MapArray::new(
6369 struct_field.clone(),
6370 offsets,
6371 entries,
6372 None,
6373 false,
6374 ));
6375
6376 let converter = RowConverter::new(vec![SortField::new(map.data_type().clone())]).unwrap();
6377 let rows = converter.convert_columns(&[Arc::clone(&map)]).unwrap();
6378 let back = converter.convert_rows(&rows).unwrap();
6379 assert_eq!(back.len(), 1);
6380 back[0].to_data().validate_full().unwrap();
6381 assert_eq!(&map, &back[0]);
6382
6383 let keys = Arc::new(StringArray::from(vec!["a", "b", "c"])) as ArrayRef;
6385 let null_values = Arc::new(NullArray::new(3)) as ArrayRef;
6386
6387 let offsets = OffsetBuffer::new(vec![0, 1, 1, 3].into());
6388 let entries_fields = vec![
6389 Arc::new(Field::new("keys", DataType::Utf8, false)),
6390 Arc::new(Field::new("values", DataType::Null, true)),
6391 ];
6392 let struct_field = Arc::new(Field::new(
6393 "entries",
6394 DataType::Struct(entries_fields.clone().into()),
6395 false,
6396 ));
6397 let entries = StructArray::new(entries_fields.into(), vec![keys, null_values], None);
6398
6399 let map: ArrayRef = Arc::new(MapArray::new(
6400 struct_field.clone(),
6401 offsets,
6402 entries,
6403 Some(vec![true, false, true].into()),
6404 false,
6405 ));
6406
6407 let rows = converter.convert_columns(&[Arc::clone(&map)]).unwrap();
6408 let back = converter.convert_rows(&rows).unwrap();
6409 assert_eq!(back.len(), 1);
6410 back[0].to_data().validate_full().unwrap();
6411 assert_eq!(&map, &back[0]);
6412
6413 let keys = Arc::new(StringArray::from(Vec::<&str>::new())) as ArrayRef;
6415 let null_values = Arc::new(NullArray::new(0)) as ArrayRef;
6416
6417 let offsets = OffsetBuffer::new(vec![0i32].into());
6418 let entries_fields = vec![
6419 Arc::new(Field::new("keys", DataType::Utf8, false)),
6420 Arc::new(Field::new("values", DataType::Null, true)),
6421 ];
6422 let struct_field = Arc::new(Field::new(
6423 "entries",
6424 DataType::Struct(entries_fields.clone().into()),
6425 false,
6426 ));
6427 let entries = StructArray::new(entries_fields.into(), vec![keys, null_values], None);
6428
6429 let map: ArrayRef = Arc::new(MapArray::new(struct_field, offsets, entries, None, false));
6430
6431 let rows = converter.convert_columns(&[Arc::clone(&map)]).unwrap();
6432 let back = converter.convert_rows(&rows).unwrap();
6433 assert_eq!(back.len(), 1);
6434 back[0].to_data().validate_full().unwrap();
6435 assert_eq!(&map, &back[0]);
6436 }
6437
6438 #[test]
6440 fn test_map_null_descending() {
6441 let keys = Arc::new(StringArray::from(vec!["a", "b", "c"])) as ArrayRef;
6443 let null_values = Arc::new(NullArray::new(3)) as ArrayRef;
6444
6445 let offsets = OffsetBuffer::new(vec![0, 1, 1, 3].into());
6446 let entries_fields = vec![
6447 Arc::new(Field::new("keys", DataType::Utf8, false)),
6448 Arc::new(Field::new("values", DataType::Null, true)),
6449 ];
6450 let struct_field = Arc::new(Field::new(
6451 "entries",
6452 DataType::Struct(entries_fields.clone().into()),
6453 false,
6454 ));
6455 let entries = StructArray::new(entries_fields.into(), vec![keys, null_values], None);
6456
6457 let map: ArrayRef = Arc::new(MapArray::new(struct_field, offsets, entries, None, false));
6458
6459 let options = SortOptions::default().with_descending(true);
6460 let field = SortField::new_with_options(map.data_type().clone(), options);
6461 let converter = RowConverter::new(vec![field]).unwrap();
6462 let rows = converter.convert_columns(&[Arc::clone(&map)]).unwrap();
6463 let back = converter.convert_rows(&rows).unwrap();
6464 assert_eq!(back.len(), 1);
6465 back[0].to_data().validate_full().unwrap();
6466 assert_eq!(&map, &back[0]);
6467 }
6468
6469 #[test]
6471 fn test_map_null_all_empty() {
6472 let keys = Arc::new(StringArray::from(Vec::<&str>::new())) as ArrayRef;
6473 let null_values = Arc::new(NullArray::new(0)) as ArrayRef;
6474
6475 let offsets = OffsetBuffer::new(vec![0, 0, 0, 0].into());
6476 let entries_fields = vec![
6477 Arc::new(Field::new("keys", DataType::Utf8, false)),
6478 Arc::new(Field::new("values", DataType::Null, true)),
6479 ];
6480 let struct_field = Arc::new(Field::new(
6481 "entries",
6482 DataType::Struct(entries_fields.clone().into()),
6483 false,
6484 ));
6485 let entries = StructArray::new(entries_fields.into(), vec![keys, null_values], None);
6486
6487 let map: ArrayRef = Arc::new(MapArray::new(struct_field, offsets, entries, None, false));
6488
6489 let converter = RowConverter::new(vec![SortField::new(map.data_type().clone())]).unwrap();
6490 let rows = converter.convert_columns(&[Arc::clone(&map)]).unwrap();
6491
6492 assert_eq!(rows.row(0), rows.row(1));
6494 assert_eq!(rows.row(1), rows.row(2));
6495
6496 let back = converter.convert_rows(&rows).unwrap();
6497 assert_eq!(back.len(), 1);
6498 back[0].to_data().validate_full().unwrap();
6499 assert_eq!(&map, &back[0]);
6500 }
6501
6502 #[test]
6504 fn test_nested_map_null() {
6505 let inner_keys = Arc::new(StringArray::from(vec!["a", "b", "c"])) as ArrayRef;
6507 let inner_null_values = Arc::new(NullArray::new(3)) as ArrayRef;
6508
6509 let inner_entries_fields = vec![
6510 Arc::new(Field::new("keys", DataType::Utf8, false)),
6511 Arc::new(Field::new("values", DataType::Null, true)),
6512 ];
6513 let inner_struct_field = Arc::new(Field::new(
6514 "entries",
6515 DataType::Struct(inner_entries_fields.clone().into()),
6516 false,
6517 ));
6518 let inner_entries = StructArray::new(
6519 inner_entries_fields.clone().into(),
6520 vec![inner_keys, inner_null_values],
6521 None,
6522 );
6523
6524 let inner_map = Arc::new(MapArray::new(
6526 inner_struct_field.clone(),
6527 OffsetBuffer::new(vec![0, 1, 3].into()),
6528 inner_entries,
6529 None,
6530 false,
6531 )) as ArrayRef;
6532
6533 let outer_keys = Arc::new(StringArray::from(vec!["x", "y"])) as ArrayRef;
6535
6536 let inner_map_type = DataType::Map(inner_struct_field.clone(), false);
6537 let outer_entries_fields = vec![
6538 Arc::new(Field::new("keys", DataType::Utf8, false)),
6539 Arc::new(Field::new("values", inner_map_type, true)),
6540 ];
6541 let outer_struct_field = Arc::new(Field::new(
6542 "entries",
6543 DataType::Struct(outer_entries_fields.clone().into()),
6544 false,
6545 ));
6546 let outer_entries = StructArray::new(
6547 outer_entries_fields.into(),
6548 vec![outer_keys, inner_map],
6549 None,
6550 );
6551
6552 let map: ArrayRef = Arc::new(MapArray::new(
6554 outer_struct_field,
6555 OffsetBuffer::new(vec![0, 1, 2].into()),
6556 outer_entries,
6557 None,
6558 false,
6559 ));
6560
6561 let converter = RowConverter::new(vec![SortField::new(map.data_type().clone())]).unwrap();
6562 let rows = converter.convert_columns(&[Arc::clone(&map)]).unwrap();
6563 let back = converter.convert_rows(&rows).unwrap();
6564 assert_eq!(back.len(), 1);
6565 back[0].to_data().validate_full().unwrap();
6566 assert_eq!(&map, &back[0]);
6567 }
6568
6569 #[test]
6571 fn test_list_of_map_null() {
6572 let keys = Arc::new(StringArray::from(vec!["a", "b", "c"])) as ArrayRef;
6574 let null_values = Arc::new(NullArray::new(3)) as ArrayRef;
6575
6576 let entries_fields = vec![
6577 Arc::new(Field::new("keys", DataType::Utf8, false)),
6578 Arc::new(Field::new("values", DataType::Null, true)),
6579 ];
6580 let struct_field = Arc::new(Field::new(
6581 "entries",
6582 DataType::Struct(entries_fields.clone().into()),
6583 false,
6584 ));
6585 let entries = StructArray::new(entries_fields.into(), vec![keys, null_values], None);
6586
6587 let map_array = Arc::new(MapArray::new(
6589 struct_field.clone(),
6590 OffsetBuffer::new(vec![0, 1, 1, 3].into()),
6591 entries,
6592 None,
6593 false,
6594 )) as ArrayRef;
6595
6596 let map_type = DataType::Map(struct_field, false);
6597 let list: ArrayRef = Arc::new(ListArray::new(
6599 Arc::new(Field::new_list_field(map_type, true)),
6600 OffsetBuffer::new(vec![0, 1, 3].into()),
6601 map_array,
6602 None,
6603 ));
6604
6605 let converter = RowConverter::new(vec![SortField::new(list.data_type().clone())]).unwrap();
6606 let rows = converter.convert_columns(&[Arc::clone(&list)]).unwrap();
6607 let back = converter.convert_rows(&rows).unwrap();
6608 assert_eq!(&list, &back[0]);
6609 }
6610
6611 #[test]
6613 fn test_map_of_list_null() {
6614 let null_array = Arc::new(NullArray::new(3)) as ArrayRef;
6616 let list_array = Arc::new(ListArray::new(
6618 Arc::new(Field::new_list_field(DataType::Null, true)),
6619 OffsetBuffer::from_lengths(vec![1, 0, 2]),
6620 null_array,
6621 None,
6622 )) as ArrayRef;
6623
6624 let keys = Arc::new(StringArray::from(vec!["a", "b", "c"])) as ArrayRef;
6625
6626 let list_type = list_array.data_type().clone();
6627 let entries_fields = vec![
6628 Arc::new(Field::new("keys", DataType::Utf8, false)),
6629 Arc::new(Field::new("values", list_type, true)),
6630 ];
6631 let struct_field = Arc::new(Field::new(
6632 "entries",
6633 DataType::Struct(entries_fields.clone().into()),
6634 false,
6635 ));
6636 let entries = StructArray::new(entries_fields.into(), vec![keys, list_array], None);
6637
6638 let map: ArrayRef = Arc::new(MapArray::new(
6640 struct_field,
6641 OffsetBuffer::new(vec![0, 3].into()),
6642 entries,
6643 None,
6644 false,
6645 ));
6646
6647 let converter = RowConverter::new(vec![SortField::new(map.data_type().clone())]).unwrap();
6648 let rows = converter.convert_columns(&[Arc::clone(&map)]).unwrap();
6649 let back = converter.convert_rows(&rows).unwrap();
6650 assert_eq!(back.len(), 1);
6651 back[0].to_data().validate_full().unwrap();
6652 assert_eq!(&map, &back[0]);
6653 }
6654
6655 #[test]
6656 fn empty_row_iter_next_back() {
6657 let rows = RowConverter::new(vec![SortField::new(DataType::UInt8)])
6658 .unwrap()
6659 .empty_rows(0, 0);
6660 let mut rows_iter = rows.iter();
6661 assert_eq!(rows_iter.next_back(), None);
6662 assert_eq!(rows_iter.next_back(), None);
6663 assert_eq!(rows_iter.next_back(), None);
6664 }
6665
6666 #[test]
6668 fn test_row_parser_skip_utf8_validation_roundtrip() {
6669 let converter = RowConverter::new(vec![SortField::new(DataType::Utf8)]).unwrap();
6670 let array = StringArray::from(vec!["arrow", "rust"]);
6671 let rows = converter.convert_columns(&[Arc::new(array) as _]).unwrap();
6672 let binary = rows.try_into_binary().expect("fits in i32 offsets");
6673
6674 let parser = unsafe { RowParser::with_skip_utf8_validate(Arc::clone(&converter.fields)) };
6676
6677 let decoded = converter
6678 .convert_rows(binary.iter().map(|b| parser.parse(b.unwrap())))
6679 .unwrap();
6680 let got: Vec<_> = decoded[0].as_string::<i32>().iter().flatten().collect();
6681 assert_eq!(got, vec!["arrow", "rust"]);
6682 }
6683
6684 #[test]
6685 fn row_iter_next_back() {
6686 let row_converter = RowConverter::new(vec![SortField::new(DataType::UInt8)]).unwrap();
6687 let mut rng = StdRng::seed_from_u64(42);
6688 let array = generate_primitive_array::<UInt8Type>(&mut rng, 100, 0.8);
6689 let rows = row_converter.convert_columns(&[Arc::new(array)]).unwrap();
6690
6691 let mut rows_iter = rows.iter();
6692 let mut bytes: Vec<u8> = vec![];
6693
6694 while let Some(row) = rows_iter.next_back() {
6695 bytes.extend(row.data.iter().rev());
6696 }
6697
6698 bytes.reverse();
6699
6700 assert_eq!(
6701 bytes,
6702 &rows.buffer.as_slice()[..*rows.offsets.last().unwrap()]
6703 );
6704
6705 assert_eq!(rows_iter.next_back(), None);
6706 assert_eq!(rows_iter.next(), None);
6707 }
6708}