1use akar_catalog::CatalogColumn;
7use akar_common::types::{Date, Interval, LogicalTypeID, Timestamp, Value};
8use arrow::array::*;
9use arrow::datatypes::{DataType as ArrowDataType, TimeUnit};
10use arrow::record_batch::RecordBatch;
11
12#[derive(Debug)]
14pub enum ParquetReaderError {
15 ParquetError(String),
17 ColumnNotFound {
19 column_name: String,
20 available: Vec<String>,
21 },
22 TypeMismatch {
24 column_name: String,
25 arrow_type: String,
26 expected_type: String,
27 },
28 ConversionError {
30 column_name: String,
31 row: usize,
32 message: String,
33 },
34}
35
36impl std::fmt::Display for ParquetReaderError {
37 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
38 match self {
39 ParquetReaderError::ParquetError(e) => write!(f, "Parquet error: {e}"),
40 ParquetReaderError::ColumnNotFound { column_name, available } => write!(
41 f,
42 "Column '{}' not found in Parquet file. Available columns: [{}]",
43 column_name,
44 available.join(", ")
45 ),
46 ParquetReaderError::TypeMismatch {
47 column_name,
48 arrow_type,
49 expected_type,
50 } => write!(
51 f,
52 "Type mismatch for column '{}': Parquet has {arrow_type}, expected {expected_type}",
53 column_name
54 ),
55 ParquetReaderError::ConversionError {
56 column_name,
57 row,
58 message,
59 } => write!(
60 f,
61 "Conversion error for column '{}' at row {row}: {message}",
62 column_name
63 ),
64 }
65 }
66}
67
68impl std::error::Error for ParquetReaderError {}
69
70impl From<parquet::errors::ParquetError> for ParquetReaderError {
71 fn from(e: parquet::errors::ParquetError) -> Self {
72 ParquetReaderError::ParquetError(e.to_string())
73 }
74}
75
76impl From<arrow::error::ArrowError> for ParquetReaderError {
77 fn from(e: arrow::error::ArrowError) -> Self {
78 ParquetReaderError::ParquetError(e.to_string())
79 }
80}
81
82pub type ParquetResult<T> = Result<T, ParquetReaderError>;
84
85pub fn read_parquet(
102 path: &str,
103 vfs: &akar_common::file_system::VirtualFileSystemRegistry,
104 columns: &[CatalogColumn],
105) -> ParquetResult<Vec<Vec<Value>>> {
106 let mut file = vfs
107 .open_read(path)
108 .map_err(|e| ParquetReaderError::ParquetError(format!("Cannot open file: {e}")))?;
109
110 let mut buffer = Vec::new();
113 file.read_to_end(&mut buffer)
114 .map_err(|e| ParquetReaderError::ParquetError(format!("Read error: {e}")))?;
115 let bytes = bytes::Bytes::from(buffer);
116
117 let builder = parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder::try_new(bytes)?;
118 let schema = builder.schema().clone();
119 let reader = builder.build()?;
120
121 let arrow_fields: Vec<(String, &ArrowDataType)> = schema
123 .fields()
124 .iter()
125 .map(|f| (f.name().clone(), f.data_type()))
126 .collect();
127
128 let available_names: Vec<String> = arrow_fields.iter().map(|(n, _)| n.clone()).collect();
129
130 let mut col_indices: Vec<usize> = Vec::with_capacity(columns.len());
132 for col in columns {
133 let pos = arrow_fields
134 .iter()
135 .position(|(name, _)| name.eq_ignore_ascii_case(&col.name));
136 match pos {
137 Some(idx) => {
138 let (_, arrow_type) = &arrow_fields[idx];
139 validate_type_compatibility(arrow_type, col.logical_type).map_err(|_| {
140 ParquetReaderError::TypeMismatch {
141 column_name: col.name.clone(),
142 arrow_type: format!("{arrow_type:?}"),
143 expected_type: format!("{:?}", col.logical_type),
144 }
145 })?;
146 col_indices.push(idx);
147 }
148 None => {
149 return Err(ParquetReaderError::ColumnNotFound {
150 column_name: col.name.clone(),
151 available: available_names.clone(),
152 });
153 }
154 }
155 }
156
157 let mut results: Vec<Vec<Value>> = Vec::new();
159
160 for batch_result in reader {
161 let batch: RecordBatch = batch_result?;
162 let num_rows = batch.num_rows();
163
164 if results.is_empty() {
166 results.reserve(num_rows * 4); }
168
169 for row_idx in 0..num_rows {
170 let mut row = Vec::with_capacity(columns.len());
171 for (catalog_idx, &arrow_col_idx) in col_indices.iter().enumerate() {
172 let col = &columns[catalog_idx];
173 let array = batch.column(arrow_col_idx);
174 let value = arrow_array_to_value(array, row_idx, &col.name, col.logical_type, results.len())?;
175 row.push(value);
176 }
177 results.push(row);
178 }
179 }
180
181 Ok(results)
182}
183
184pub struct ParquetStreamReader {
190 reader: parquet::arrow::arrow_reader::ParquetRecordBatchReader,
191 col_indices: Vec<usize>,
192 columns: Vec<CatalogColumn>,
193}
194
195impl Iterator for ParquetStreamReader {
196 type Item = ParquetResult<Vec<Vec<Value>>>;
197
198 fn next(&mut self) -> Option<Self::Item> {
199 match self.reader.next() {
200 Some(Ok(batch)) => {
201 let num_rows = batch.num_rows();
202 let mut rows = Vec::with_capacity(num_rows);
203 for row_idx in 0..num_rows {
204 let mut row = Vec::with_capacity(self.columns.len());
205 for (catalog_idx, &arrow_col_idx) in self.col_indices.iter().enumerate() {
206 let col = &self.columns[catalog_idx];
207 let array = batch.column(arrow_col_idx);
208 let value = match arrow_array_to_value(array, row_idx, &col.name, col.logical_type, rows.len())
209 {
210 Ok(v) => v,
211 Err(e) => return Some(Err(e)),
212 };
213 row.push(value);
214 }
215 rows.push(row);
216 }
217 Some(Ok(rows))
218 }
219 Some(Err(e)) => Some(Err(ParquetReaderError::ParquetError(e.to_string()))),
220 None => None,
221 }
222 }
223}
224
225pub fn stream_parquet(
231 path: &str,
232 vfs: &akar_common::file_system::VirtualFileSystemRegistry,
233 columns: &[CatalogColumn],
234) -> ParquetResult<ParquetStreamReader> {
235 let mut file = vfs
236 .open_read(path)
237 .map_err(|e| ParquetReaderError::ParquetError(format!("Cannot open file: {e}")))?;
238
239 let mut buffer = Vec::new();
240 file.read_to_end(&mut buffer)
241 .map_err(|e| ParquetReaderError::ParquetError(format!("Read error: {e}")))?;
242 let bytes = bytes::Bytes::from(buffer);
243
244 let builder = parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder::try_new(bytes)?;
245 let schema = builder.schema().clone();
246 let reader = builder.build()?;
247
248 let arrow_fields: Vec<(String, &arrow::datatypes::DataType)> = schema
249 .fields()
250 .iter()
251 .map(|f| (f.name().clone(), f.data_type()))
252 .collect();
253
254 let available_names: Vec<String> = arrow_fields.iter().map(|(n, _)| n.clone()).collect();
255
256 let mut col_indices: Vec<usize> = Vec::with_capacity(columns.len());
257 for col in columns {
258 let pos = arrow_fields
259 .iter()
260 .position(|(name, _)| name.eq_ignore_ascii_case(&col.name));
261 match pos {
262 Some(idx) => {
263 let (_, arrow_type) = &arrow_fields[idx];
264 validate_type_compatibility(arrow_type, col.logical_type).map_err(|_| {
265 ParquetReaderError::TypeMismatch {
266 column_name: col.name.clone(),
267 arrow_type: format!("{arrow_type:?}"),
268 expected_type: format!("{:?}", col.logical_type),
269 }
270 })?;
271 col_indices.push(idx);
272 }
273 None => {
274 return Err(ParquetReaderError::ColumnNotFound {
275 column_name: col.name.clone(),
276 available: available_names.clone(),
277 });
278 }
279 }
280 }
281
282 Ok(ParquetStreamReader {
283 reader,
284 col_indices,
285 columns: columns.to_vec(),
286 })
287}
288
289fn validate_type_compatibility(arrow_type: &ArrowDataType, expected: LogicalTypeID) -> Result<(), ()> {
293 match (arrow_type, expected) {
294 (ArrowDataType::Boolean, LogicalTypeID::Bool) => Ok(()),
295 (ArrowDataType::Int8, LogicalTypeID::Int8) => Ok(()),
296 (ArrowDataType::Int16, LogicalTypeID::Int16) => Ok(()),
297 (ArrowDataType::Int32, LogicalTypeID::Int32) => Ok(()),
298 (ArrowDataType::Int64, LogicalTypeID::Int64 | LogicalTypeID::Serial) => Ok(()),
299 (ArrowDataType::UInt8, LogicalTypeID::UInt8) => Ok(()),
300 (ArrowDataType::UInt16, LogicalTypeID::UInt16) => Ok(()),
301 (ArrowDataType::UInt32, LogicalTypeID::UInt32) => Ok(()),
302 (ArrowDataType::UInt64, LogicalTypeID::UInt64) => Ok(()),
303 (ArrowDataType::Float32, LogicalTypeID::Float) => Ok(()),
304 (ArrowDataType::Float64, LogicalTypeID::Double) => Ok(()),
305 (ArrowDataType::Utf8 | ArrowDataType::LargeUtf8, LogicalTypeID::String) => Ok(()),
306 (ArrowDataType::Binary | ArrowDataType::LargeBinary, LogicalTypeID::Blob) => Ok(()),
307 (ArrowDataType::Date32 | ArrowDataType::Date64, LogicalTypeID::Date) => Ok(()),
308 (ArrowDataType::Timestamp(TimeUnit::Second, _), LogicalTypeID::TimestampSec) => Ok(()),
309 (ArrowDataType::Timestamp(TimeUnit::Millisecond, _), LogicalTypeID::TimestampMs) => Ok(()),
310 (ArrowDataType::Timestamp(TimeUnit::Microsecond, _), LogicalTypeID::Timestamp) => Ok(()),
311 (ArrowDataType::Timestamp(TimeUnit::Nanosecond, _), LogicalTypeID::TimestampNs) => Ok(()),
312 (ArrowDataType::Duration(_), LogicalTypeID::Interval) => Ok(()),
313 (ArrowDataType::List(_), LogicalTypeID::List) => Ok(()),
314 (ArrowDataType::Struct(_), LogicalTypeID::Struct) => Ok(()),
315 (ArrowDataType::Map(_, _), LogicalTypeID::Map) => Ok(()),
316 (ArrowDataType::Int8, LogicalTypeID::Int64 | LogicalTypeID::Int32 | LogicalTypeID::Int16) => Ok(()),
318 (ArrowDataType::Int16, LogicalTypeID::Int64 | LogicalTypeID::Int32) => Ok(()),
319 (ArrowDataType::Int32, LogicalTypeID::Int64) => Ok(()),
320 (ArrowDataType::UInt8, LogicalTypeID::UInt64 | LogicalTypeID::UInt32 | LogicalTypeID::UInt16) => Ok(()),
321 (ArrowDataType::UInt16, LogicalTypeID::UInt64 | LogicalTypeID::UInt32) => Ok(()),
322 (ArrowDataType::UInt32, LogicalTypeID::UInt64) => Ok(()),
323 (ArrowDataType::Float32, LogicalTypeID::Double) => Ok(()),
324 (_, LogicalTypeID::String) => Ok(()),
326 _ => Err(()),
327 }
328}
329
330fn arrow_array_to_value(
334 array: &dyn Array,
335 row: usize,
336 column_name: &str,
337 target_type: LogicalTypeID,
338 _global_row: usize,
339) -> ParquetResult<Value> {
340 if array.is_null(row) {
342 return Ok(Value::Null);
343 }
344
345 match target_type {
346 LogicalTypeID::Bool => {
347 let arr = downcast::<BooleanArray>(array, column_name)?;
348 Ok(Value::Bool(arr.value(row)))
349 }
350 LogicalTypeID::Int64 | LogicalTypeID::Serial => {
351 let val = cast_int_to_i64(array, row, column_name)?;
352 Ok(Value::Int64(val))
353 }
354 LogicalTypeID::Int32 => {
355 let val = cast_int_to_i64(array, row, column_name)?;
356 Ok(Value::Int32(val as i32))
357 }
358 LogicalTypeID::Int16 => {
359 let val = cast_int_to_i64(array, row, column_name)?;
360 Ok(Value::Int16(val as i16))
361 }
362 LogicalTypeID::Int8 => {
363 let val = cast_int_to_i64(array, row, column_name)?;
364 Ok(Value::Int8(val as i8))
365 }
366 LogicalTypeID::UInt64 => {
367 let val = cast_uint_to_u64(array, row, column_name)?;
368 Ok(Value::UInt64(val))
369 }
370 LogicalTypeID::UInt32 => {
371 let val = cast_uint_to_u64(array, row, column_name)?;
372 Ok(Value::UInt32(val as u32))
373 }
374 LogicalTypeID::UInt16 => {
375 let val = cast_uint_to_u64(array, row, column_name)?;
376 Ok(Value::UInt16(val as u16))
377 }
378 LogicalTypeID::UInt8 => {
379 let val = cast_uint_to_u64(array, row, column_name)?;
380 Ok(Value::UInt8(val as u8))
381 }
382 LogicalTypeID::Double => {
383 let val = cast_to_f64(array, row, column_name)?;
384 Ok(Value::Double(val))
385 }
386 LogicalTypeID::Float => {
387 let val = cast_to_f64(array, row, column_name)?;
388 Ok(Value::Float(val as f32))
389 }
390 LogicalTypeID::String => {
391 let s = array_to_string(array, row, column_name)?;
392 Ok(Value::String(s))
393 }
394 LogicalTypeID::Blob => {
395 let arr = downcast::<BinaryArray>(array, column_name)?;
396 Ok(Value::Blob(arr.value(row).to_vec()))
397 }
398 LogicalTypeID::Date => {
399 let val = cast_date_to_days(array, row, column_name)?;
400 Ok(Value::Date(Date::from_days_since_epoch(val)))
401 }
402 LogicalTypeID::Timestamp => {
403 let micros = cast_timestamp_to_micros(array, row, column_name)?;
404 Ok(Value::Timestamp(Timestamp::from_micros_since_epoch(micros)))
405 }
406 LogicalTypeID::TimestampMs => {
407 let micros = cast_timestamp_to_micros(array, row, column_name)?;
408 Ok(Value::TimestampMs(Timestamp::from_micros_since_epoch(micros)))
409 }
410 LogicalTypeID::TimestampSec => {
411 let micros = cast_timestamp_to_micros(array, row, column_name)?;
412 Ok(Value::TimestampSec(Timestamp(micros / 1_000_000)))
413 }
414 LogicalTypeID::TimestampNs => {
415 let micros = cast_timestamp_to_micros(array, row, column_name)?;
416 Ok(Value::TimestampNs(Timestamp(micros * 1000)))
417 }
418 LogicalTypeID::TimestampTz => {
419 let micros = cast_timestamp_to_micros(array, row, column_name)?;
420 Ok(Value::TimestampTz(akar_common::types::TimestampTZ(micros)))
421 }
422 LogicalTypeID::Interval => {
423 let arr = downcast::<DurationMicrosecondArray>(array, column_name)?;
424 Ok(Value::Interval(Interval::new(0, 0, arr.value(row))))
425 }
426 LogicalTypeID::List => {
427 let vals = array_list_to_values(array, row, column_name)?;
428 Ok(Value::List(vals))
429 }
430 LogicalTypeID::Struct => {
431 let vals = array_struct_to_values(array, row, column_name)?;
432 Ok(Value::Struct(vals))
433 }
434 LogicalTypeID::Map => {
435 let vals = array_map_to_values(array, row, column_name)?;
436 Ok(Value::Map(vals))
437 }
438 _ => {
440 let s = array_to_string(array, row, column_name)?;
441 Ok(Value::String(s))
442 }
443 }
444}
445
446fn downcast<'a, T: Array + 'static>(array: &'a dyn Array, column_name: &str) -> ParquetResult<&'a T> {
449 array
450 .as_any()
451 .downcast_ref::<T>()
452 .ok_or_else(|| ParquetReaderError::ConversionError {
453 column_name: column_name.to_string(),
454 row: 0,
455 message: format!(
456 "expected array type {} but got {:?}",
457 std::any::type_name::<T>(),
458 array.data_type()
459 ),
460 })
461}
462
463fn cast_int_to_i64(array: &dyn Array, row: usize, column_name: &str) -> ParquetResult<i64> {
467 if let Some(arr) = array.as_any().downcast_ref::<Int8Array>() {
468 return Ok(arr.value(row) as i64);
469 }
470 if let Some(arr) = array.as_any().downcast_ref::<Int16Array>() {
471 return Ok(arr.value(row) as i64);
472 }
473 if let Some(arr) = array.as_any().downcast_ref::<Int32Array>() {
474 return Ok(arr.value(row) as i64);
475 }
476 if let Some(arr) = array.as_any().downcast_ref::<Int64Array>() {
477 return Ok(arr.value(row));
478 }
479 Err(ParquetReaderError::ConversionError {
480 column_name: column_name.to_string(),
481 row,
482 message: format!("cannot cast {:?} to Int64", array.data_type()),
483 })
484}
485
486fn cast_uint_to_u64(array: &dyn Array, row: usize, column_name: &str) -> ParquetResult<u64> {
488 if let Some(arr) = array.as_any().downcast_ref::<UInt8Array>() {
489 return Ok(arr.value(row) as u64);
490 }
491 if let Some(arr) = array.as_any().downcast_ref::<UInt16Array>() {
492 return Ok(arr.value(row) as u64);
493 }
494 if let Some(arr) = array.as_any().downcast_ref::<UInt32Array>() {
495 return Ok(arr.value(row) as u64);
496 }
497 if let Some(arr) = array.as_any().downcast_ref::<UInt64Array>() {
498 return Ok(arr.value(row));
499 }
500 Err(ParquetReaderError::ConversionError {
501 column_name: column_name.to_string(),
502 row,
503 message: format!("cannot cast {:?} to UInt64", array.data_type()),
504 })
505}
506
507fn cast_to_f64(array: &dyn Array, row: usize, column_name: &str) -> ParquetResult<f64> {
509 if let Some(arr) = array.as_any().downcast_ref::<Float32Array>() {
510 return Ok(arr.value(row) as f64);
511 }
512 if let Some(arr) = array.as_any().downcast_ref::<Float64Array>() {
513 return Ok(arr.value(row));
514 }
515 if let Some(arr) = array.as_any().downcast_ref::<Int64Array>() {
516 return Ok(arr.value(row) as f64);
517 }
518 if let Some(arr) = array.as_any().downcast_ref::<Int32Array>() {
519 return Ok(arr.value(row) as f64);
520 }
521 Err(ParquetReaderError::ConversionError {
522 column_name: column_name.to_string(),
523 row,
524 message: format!("cannot cast {:?} to Float64", array.data_type()),
525 })
526}
527
528fn array_to_string(array: &dyn Array, row: usize, _column_name: &str) -> ParquetResult<String> {
530 if let Some(arr) = array.as_any().downcast_ref::<StringArray>() {
531 return Ok(arr.value(row).to_string());
532 }
533 if let Some(arr) = array.as_any().downcast_ref::<LargeStringArray>() {
534 return Ok(arr.value(row).to_string());
535 }
536 if let Some(arr) = array.as_any().downcast_ref::<BinaryArray>() {
537 return Ok(String::from_utf8_lossy(arr.value(row)).to_string());
538 }
539 if let Some(arr) = array.as_any().downcast_ref::<Int64Array>() {
540 return Ok(arr.value(row).to_string());
541 }
542 if let Some(arr) = array.as_any().downcast_ref::<Float64Array>() {
543 return Ok(arr.value(row).to_string());
544 }
545 if let Some(arr) = array.as_any().downcast_ref::<BooleanArray>() {
546 return Ok(arr.value(row).to_string());
547 }
548 Ok(format!("{:?}", array))
550}
551
552fn cast_date_to_days(array: &dyn Array, row: usize, column_name: &str) -> ParquetResult<i32> {
556 if let Some(arr) = array.as_any().downcast_ref::<Date32Array>() {
557 return Ok(arr.value(row));
558 }
559 if let Some(arr) = array.as_any().downcast_ref::<Date64Array>() {
560 return Ok((arr.value(row) / 86_400_000) as i32);
562 }
563 Err(ParquetReaderError::ConversionError {
564 column_name: column_name.to_string(),
565 row,
566 message: format!("cannot cast {:?} to Date", array.data_type()),
567 })
568}
569
570fn cast_timestamp_to_micros(array: &dyn Array, row: usize, column_name: &str) -> ParquetResult<i64> {
572 if let Some(arr) = array.as_any().downcast_ref::<TimestampSecondArray>() {
573 return Ok(arr.value(row) * 1_000_000);
574 }
575 if let Some(arr) = array.as_any().downcast_ref::<TimestampMillisecondArray>() {
576 return Ok(arr.value(row) * 1_000);
577 }
578 if let Some(arr) = array.as_any().downcast_ref::<TimestampMicrosecondArray>() {
579 return Ok(arr.value(row));
580 }
581 if let Some(arr) = array.as_any().downcast_ref::<TimestampNanosecondArray>() {
582 return Ok(arr.value(row) / 1_000);
583 }
584 Err(ParquetReaderError::ConversionError {
585 column_name: column_name.to_string(),
586 row,
587 message: format!("cannot cast {:?} to Timestamp", array.data_type()),
588 })
589}
590
591fn array_list_to_values(array: &dyn Array, row: usize, column_name: &str) -> ParquetResult<Vec<Value>> {
599 let list_arr = downcast::<ListArray>(array, column_name)?;
600 let values = list_arr.value(row);
601 let mut result = Vec::with_capacity(values.len());
602 for i in 0..values.len() {
603 if values.is_null(i) {
604 result.push(Value::Null);
605 } else if let Some(f) = values.as_any().downcast_ref::<Float64Array>() {
606 result.push(Value::Double(f.value(i)));
607 } else if let Some(f) = values.as_any().downcast_ref::<Float32Array>() {
608 result.push(Value::Float(f.value(i)));
609 } else if let Some(a) = values.as_any().downcast_ref::<Int64Array>() {
610 result.push(Value::Int64(a.value(i)));
611 } else if let Some(a) = values.as_any().downcast_ref::<Int32Array>() {
612 result.push(Value::Int32(a.value(i)));
613 } else if let Some(a) = values.as_any().downcast_ref::<Int16Array>() {
614 result.push(Value::Int16(a.value(i)));
615 } else if let Some(a) = values.as_any().downcast_ref::<Int8Array>() {
616 result.push(Value::Int8(a.value(i)));
617 } else if let Some(a) = values.as_any().downcast_ref::<BooleanArray>() {
618 result.push(Value::Bool(a.value(i)));
619 } else if let Some(s) = values.as_any().downcast_ref::<StringArray>() {
620 result.push(Value::String(s.value(i).to_string()));
621 } else {
622 let s = array_to_string(&values, i, column_name)?;
624 result.push(Value::String(s));
625 }
626 }
627 Ok(result)
628}
629
630fn array_struct_to_values(array: &dyn Array, row: usize, column_name: &str) -> ParquetResult<Vec<(String, Value)>> {
632 let struct_arr = downcast::<StructArray>(array, column_name)?;
633 let mut result = Vec::with_capacity(struct_arr.num_columns());
634 for col_idx in 0..struct_arr.num_columns() {
635 let field = struct_arr.column(col_idx);
636 let field_name = struct_arr
637 .fields()
638 .get(col_idx)
639 .map(|f| f.name().clone())
640 .unwrap_or_else(|| format!("_{col_idx}"));
641 let val = if field.is_null(row) {
642 Value::Null
643 } else {
644 Value::String(array_to_string(field.as_ref(), row, column_name)?)
645 };
646 result.push((field_name, val));
647 }
648 Ok(result)
649}
650
651fn array_map_to_values(array: &dyn Array, row: usize, column_name: &str) -> ParquetResult<Vec<(Value, Value)>> {
653 let map_arr = downcast::<MapArray>(array, column_name)?;
654 let entries = map_arr.value(row);
655 let keys = map_arr.keys();
656 let values = map_arr.values();
657
658 let mut result = Vec::new();
659 if let Some(_entries_struct) = entries.as_any().downcast_ref::<StructArray>() {
660 for i in 0..entries.len() {
662 let key_val = if keys.is_null(i) {
663 Value::Null
664 } else {
665 Value::String(array_to_string(keys, i, column_name)?)
666 };
667 let val_val = if values.is_null(i) {
668 Value::Null
669 } else {
670 Value::String(array_to_string(values, i, column_name)?)
671 };
672 result.push((key_val, val_val));
673 }
674 }
675 Ok(result)
676}
677
678#[cfg(test)]
681mod tests {
682 use super::*;
683 use arrow::datatypes::{DataType, Field, Schema};
684 use std::sync::Arc;
685
686 fn write_parquet_batch(dir: &tempfile::TempDir, filename: &str, batch: &RecordBatch) -> std::path::PathBuf {
688 let path = dir.path().join(filename);
689 let file = std::fs::File::create(&path).unwrap();
690 let schema = batch.schema();
691 let mut writer = parquet::arrow::ArrowWriter::try_new(file, schema, None).unwrap();
692 writer.write(batch).unwrap();
693 writer.close().unwrap();
694 path
695 }
696
697 fn test_schema() -> Vec<CatalogColumn> {
698 vec![
699 CatalogColumn {
700 compression: akar_common::enums::CompressionType::Uncompressed,
701 name: "name".into(),
702 logical_type: LogicalTypeID::String,
703 is_primary_key: true,
704 default_value: None,
705 },
706 CatalogColumn {
707 compression: akar_common::enums::CompressionType::Uncompressed,
708 name: "age".into(),
709 logical_type: LogicalTypeID::Int64,
710 is_primary_key: false,
711 default_value: None,
712 },
713 CatalogColumn {
714 compression: akar_common::enums::CompressionType::Uncompressed,
715 name: "score".into(),
716 logical_type: LogicalTypeID::Double,
717 is_primary_key: false,
718 default_value: None,
719 },
720 CatalogColumn {
721 compression: akar_common::enums::CompressionType::Uncompressed,
722 name: "active".into(),
723 logical_type: LogicalTypeID::Bool,
724 is_primary_key: false,
725 default_value: None,
726 },
727 ]
728 }
729
730 #[test]
731 fn test_validate_type_compatibility() {
732 assert!(validate_type_compatibility(&ArrowDataType::Int64, LogicalTypeID::Int64).is_ok());
733 assert!(validate_type_compatibility(&ArrowDataType::Utf8, LogicalTypeID::String).is_ok());
734 assert!(validate_type_compatibility(&ArrowDataType::Boolean, LogicalTypeID::Bool).is_ok());
735 assert!(validate_type_compatibility(&ArrowDataType::Float64, LogicalTypeID::Double).is_ok());
736 assert!(validate_type_compatibility(&ArrowDataType::Int32, LogicalTypeID::Int64).is_ok());
738 assert!(validate_type_compatibility(&ArrowDataType::Int8, LogicalTypeID::Int64).is_ok());
739 assert!(validate_type_compatibility(&ArrowDataType::Int64, LogicalTypeID::String).is_ok()); assert!(validate_type_compatibility(&ArrowDataType::Boolean, LogicalTypeID::Int64).is_err());
742 }
743
744 #[test]
745 fn test_cast_int_to_i64() {
746 let i32_arr = Int32Array::from(vec![42, -1]);
747 assert_eq!(cast_int_to_i64(&i32_arr, 0, "col").unwrap(), 42i64);
748 assert_eq!(cast_int_to_i64(&i32_arr, 1, "col").unwrap(), -1i64);
749
750 let i64_arr = Int64Array::from(vec![999_999_999_999i64]);
751 assert_eq!(cast_int_to_i64(&i64_arr, 0, "col").unwrap(), 999_999_999_999i64);
752 }
753
754 #[test]
755 fn test_cast_to_f64() {
756 let f64_arr = Float64Array::from(vec![3.15]);
757 assert!((cast_to_f64(&f64_arr, 0, "col").unwrap() - 3.15).abs() < 1e-10);
758
759 let f32_arr = Float32Array::from(vec![2.5f32]);
760 assert!((cast_to_f64(&f32_arr, 0, "col").unwrap() - 2.5).abs() < 1e-10);
761 }
762
763 #[test]
764 fn test_array_to_string() {
765 let str_arr = StringArray::from(vec!["hello"]);
766 assert_eq!(array_to_string(&str_arr, 0, "col").unwrap(), "hello");
767
768 let int_arr = Int64Array::from(vec![42]);
769 assert_eq!(array_to_string(&int_arr, 0, "col").unwrap(), "42");
770 }
771
772 #[test]
773 fn test_cast_date_to_days() {
774 let date_arr = Date32Array::from(vec![0i32, 19723i32]); assert_eq!(cast_date_to_days(&date_arr, 0, "col").unwrap(), 0);
776 assert_eq!(cast_date_to_days(&date_arr, 1, "col").unwrap(), 19723);
777 }
778
779 #[test]
780 fn test_cast_timestamp_to_micros() {
781 use arrow::array::TimestampMicrosecondArray;
782 let ts_arr = TimestampMicrosecondArray::from(vec![1_700_000_000_000_000i64]);
783 let micros = cast_timestamp_to_micros(&ts_arr, 0, "col").unwrap();
784 assert_eq!(micros, 1_700_000_000_000_000i64);
785 }
786
787 #[test]
788 fn test_column_not_found() {
789 let dir = tempfile::tempdir().unwrap();
790 let schema = Arc::new(Schema::new(vec![
792 Field::new("x", DataType::Int64, false),
793 Field::new("y", DataType::Utf8, false),
794 ]));
795 let batch = RecordBatch::try_new(
796 schema.clone(),
797 vec![
798 Arc::new(Int64Array::from(vec![1])),
799 Arc::new(StringArray::from(vec!["a"])),
800 ],
801 )
802 .unwrap();
803 let path = write_parquet_batch(&dir, "test.parquet", &batch);
804
805 let columns = vec![CatalogColumn {
806 compression: akar_common::enums::CompressionType::Uncompressed,
807 name: "missing_col".into(),
808 logical_type: LogicalTypeID::Int64,
809 is_primary_key: false,
810 default_value: None,
811 }];
812
813 let result = read_parquet(
814 path.to_str().unwrap(),
815 &akar_common::file_system::VirtualFileSystemRegistry::new(),
816 &columns,
817 );
818 assert!(result.is_err());
819 match result.unwrap_err() {
820 ParquetReaderError::ColumnNotFound { column_name, .. } => {
821 assert_eq!(column_name, "missing_col");
822 }
823 e => panic!("Expected ColumnNotFound, got: {e}"),
824 }
825 }
826
827 #[test]
828 fn test_type_mismatch() {
829 let dir = tempfile::tempdir().unwrap();
830 let schema = Arc::new(Schema::new(vec![Field::new("val", DataType::Boolean, false)]));
831 let batch = RecordBatch::try_new(schema.clone(), vec![Arc::new(BooleanArray::from(vec![true]))]).unwrap();
832 let path = write_parquet_batch(&dir, "mismatch.parquet", &batch);
833
834 let columns = vec![CatalogColumn {
835 compression: akar_common::enums::CompressionType::Uncompressed,
836 name: "val".into(),
837 logical_type: LogicalTypeID::Int64,
838 is_primary_key: false,
839 default_value: None,
840 }];
841
842 let result = read_parquet(
843 path.to_str().unwrap(),
844 &akar_common::file_system::VirtualFileSystemRegistry::new(),
845 &columns,
846 );
847 assert!(result.is_err());
848 match result.unwrap_err() {
849 ParquetReaderError::TypeMismatch { .. } => {} e => panic!("Expected TypeMismatch, got: {e}"),
851 }
852 }
853
854 #[test]
855 fn test_file_not_found() {
856 let result = read_parquet(
857 "nonexistent.parquet",
858 &akar_common::file_system::VirtualFileSystemRegistry::new(),
859 &test_schema(),
860 );
861 assert!(result.is_err());
862 match result.unwrap_err() {
863 ParquetReaderError::ParquetError(_) => {} _ => panic!("Expected ParquetError"),
865 }
866 }
867
868 #[test]
869 fn test_arrow_array_to_value_basic() {
870 let bool_arr = BooleanArray::from(vec![Some(true), None, Some(false)]);
872 assert_eq!(
873 arrow_array_to_value(&bool_arr, 0, "b", LogicalTypeID::Bool, 0).unwrap(),
874 Value::Bool(true)
875 );
876 assert_eq!(
877 arrow_array_to_value(&bool_arr, 1, "b", LogicalTypeID::Bool, 0).unwrap(),
878 Value::Null
879 );
880
881 let int_arr = Int64Array::from(vec![Some(42), None]);
883 assert_eq!(
884 arrow_array_to_value(&int_arr, 0, "i", LogicalTypeID::Int64, 0).unwrap(),
885 Value::Int64(42)
886 );
887 assert_eq!(
888 arrow_array_to_value(&int_arr, 1, "i", LogicalTypeID::Int64, 0).unwrap(),
889 Value::Null
890 );
891
892 let str_arr = StringArray::from(vec![Some("hello"), None]);
894 assert_eq!(
895 arrow_array_to_value(&str_arr, 0, "s", LogicalTypeID::String, 0).unwrap(),
896 Value::String("hello".into())
897 );
898 assert_eq!(
899 arrow_array_to_value(&str_arr, 1, "s", LogicalTypeID::String, 0).unwrap(),
900 Value::Null
901 );
902
903 let f64_arr = Float64Array::from(vec![Some(3.15), None]);
905 assert_eq!(
906 arrow_array_to_value(&f64_arr, 0, "d", LogicalTypeID::Double, 0).unwrap(),
907 Value::Double(3.15)
908 );
909 }
910
911 #[test]
912 fn test_arrow_array_widening() {
913 let i32_arr = Int32Array::from(vec![100]);
915 assert_eq!(
916 arrow_array_to_value(&i32_arr, 0, "i", LogicalTypeID::Int64, 0).unwrap(),
917 Value::Int64(100)
918 );
919
920 let f32_arr = Float32Array::from(vec![2.5f32]);
922 let val = arrow_array_to_value(&f32_arr, 0, "f", LogicalTypeID::Double, 0).unwrap();
923 if let Value::Double(d) = val {
924 assert!((d - 2.5).abs() < 1e-6);
925 } else {
926 panic!("Expected Double");
927 }
928 }
929
930 #[test]
931 fn test_unsigned_int_types() {
932 let u8_arr = UInt8Array::from(vec![200u8]);
933 assert_eq!(
934 arrow_array_to_value(&u8_arr, 0, "u", LogicalTypeID::UInt8, 0).unwrap(),
935 Value::UInt8(200)
936 );
937
938 let u32_arr = UInt32Array::from(vec![100000u32]);
939 assert_eq!(
940 arrow_array_to_value(&u32_arr, 0, "u", LogicalTypeID::UInt32, 0).unwrap(),
941 Value::UInt32(100000)
942 );
943
944 let u64_arr = UInt64Array::from(vec![u64::MAX]);
945 assert_eq!(
946 arrow_array_to_value(&u64_arr, 0, "u", LogicalTypeID::UInt64, 0).unwrap(),
947 Value::UInt64(u64::MAX)
948 );
949 }
950
951 #[test]
952 fn test_date_conversion() {
953 let date_arr = Date32Array::from(vec![7439i32]); let val = arrow_array_to_value(&date_arr, 0, "d", LogicalTypeID::Date, 0).unwrap();
955 if let Value::Date(d) = val {
956 assert_eq!(d.days_since_epoch(), 7439);
957 } else {
958 panic!("Expected Date");
959 }
960 }
961
962 #[test]
963 fn test_timestamp_conversion() {
964 use arrow::array::TimestampMicrosecondArray;
965 let ts_arr = TimestampMicrosecondArray::from(vec![1_704_198_600_000_000i64]); let val = arrow_array_to_value(&ts_arr, 0, "t", LogicalTypeID::Timestamp, 0).unwrap();
967 if let Value::Timestamp(ts) = val {
968 assert_eq!(ts.micros_since_epoch(), 1_704_198_600_000_000);
969 } else {
970 panic!("Expected Timestamp");
971 }
972 }
973
974 #[test]
975 fn test_blob_conversion() {
976 let blob_arr = BinaryArray::from(vec![&b"hello"[..]]);
977 let val = arrow_array_to_value(&blob_arr, 0, "b", LogicalTypeID::Blob, 0).unwrap();
978 assert_eq!(val, Value::Blob(b"hello".to_vec()));
979 }
980
981 #[test]
982 fn test_round_trip_parquet() {
983 let dir = tempfile::tempdir().unwrap();
984
985 let schema = Arc::new(Schema::new(vec![
987 Field::new("name", DataType::Utf8, false),
988 Field::new("age", DataType::Int64, false),
989 Field::new("score", DataType::Float64, false),
990 Field::new("active", DataType::Boolean, false),
991 ]));
992
993 let names = StringArray::from(vec!["Alice", "Bob", "Charlie"]);
994 let ages = Int64Array::from(vec![30, 25, 35]);
995 let scores = Float64Array::from(vec![95.5, 87.3, 91.2]);
996 let actives = BooleanArray::from(vec![true, false, true]);
997
998 let batch = RecordBatch::try_new(
999 schema,
1000 vec![Arc::new(names), Arc::new(ages), Arc::new(scores), Arc::new(actives)],
1001 )
1002 .unwrap();
1003
1004 let parquet_path = write_parquet_batch(&dir, "roundtrip.parquet", &batch);
1005
1006 let columns = test_schema();
1008 let rows = read_parquet(
1009 parquet_path.to_str().unwrap(),
1010 &akar_common::file_system::VirtualFileSystemRegistry::new(),
1011 &columns,
1012 )
1013 .unwrap();
1014
1015 assert_eq!(rows.len(), 3);
1016 assert_eq!(rows[0][0], Value::String("Alice".into()));
1017 assert_eq!(rows[0][1], Value::Int64(30));
1018 assert_eq!(rows[0][2], Value::Double(95.5));
1019 assert_eq!(rows[0][3], Value::Bool(true));
1020
1021 assert_eq!(rows[1][0], Value::String("Bob".into()));
1022 assert_eq!(rows[1][1], Value::Int64(25));
1023 assert_eq!(rows[1][2], Value::Double(87.3));
1024 assert_eq!(rows[1][3], Value::Bool(false));
1025
1026 assert_eq!(rows[2][0], Value::String("Charlie".into()));
1027 assert_eq!(rows[2][1], Value::Int64(35));
1028 assert_eq!(rows[2][2], Value::Double(91.2));
1029 assert_eq!(rows[2][3], Value::Bool(true));
1030 }
1031
1032 #[test]
1033 fn test_round_trip_with_nulls() {
1034 let dir = tempfile::tempdir().unwrap();
1035
1036 let schema = Arc::new(Schema::new(vec![
1037 Field::new("name", DataType::Utf8, true),
1038 Field::new("age", DataType::Int64, true),
1039 ]));
1040
1041 let names = StringArray::from(vec![Some("Alice"), None, Some("Charlie")]);
1042 let ages = Int64Array::from(vec![Some(30), Some(25), None]);
1043
1044 let batch = RecordBatch::try_new(schema, vec![Arc::new(names), Arc::new(ages)]).unwrap();
1045
1046 let parquet_path = write_parquet_batch(&dir, "nulls.parquet", &batch);
1047
1048 let columns = vec![
1049 CatalogColumn {
1050 compression: akar_common::enums::CompressionType::Uncompressed,
1051 name: "name".into(),
1052 logical_type: LogicalTypeID::String,
1053 is_primary_key: false,
1054 default_value: None,
1055 },
1056 CatalogColumn {
1057 compression: akar_common::enums::CompressionType::Uncompressed,
1058 name: "age".into(),
1059 logical_type: LogicalTypeID::Int64,
1060 is_primary_key: false,
1061 default_value: None,
1062 },
1063 ];
1064
1065 let rows = read_parquet(
1066 parquet_path.to_str().unwrap(),
1067 &akar_common::file_system::VirtualFileSystemRegistry::new(),
1068 &columns,
1069 )
1070 .unwrap();
1071 assert_eq!(rows.len(), 3);
1072 assert_eq!(rows[0][0], Value::String("Alice".into()));
1073 assert_eq!(rows[0][1], Value::Int64(30));
1074 assert_eq!(rows[1][0], Value::Null);
1075 assert_eq!(rows[1][1], Value::Int64(25));
1076 assert_eq!(rows[2][0], Value::String("Charlie".into()));
1077 assert_eq!(rows[2][1], Value::Null);
1078 }
1079
1080 #[test]
1081 fn test_round_trip_unsigned_and_floats() {
1082 let dir = tempfile::tempdir().unwrap();
1083
1084 let schema = Arc::new(Schema::new(vec![
1085 Field::new("small", DataType::UInt8, false),
1086 Field::new("medium", DataType::UInt32, false),
1087 Field::new("large", DataType::UInt64, false),
1088 Field::new("temp", DataType::Float32, false),
1089 ]));
1090
1091 let small = UInt8Array::from(vec![100u8, 200u8]);
1092 let medium = UInt32Array::from(vec![1000u32, 50000u32]);
1093 let large = UInt64Array::from(vec![100000u64, u64::MAX]);
1094 let temp = Float32Array::from(vec![36.5f32, 98.6f32]);
1095
1096 let batch = RecordBatch::try_new(
1097 schema,
1098 vec![Arc::new(small), Arc::new(medium), Arc::new(large), Arc::new(temp)],
1099 )
1100 .unwrap();
1101
1102 let parquet_path = write_parquet_batch(&dir, "uints.parquet", &batch);
1103
1104 let columns = vec![
1105 CatalogColumn {
1106 compression: akar_common::enums::CompressionType::Uncompressed,
1107 name: "small".into(),
1108 logical_type: LogicalTypeID::UInt8,
1109 is_primary_key: false,
1110 default_value: None,
1111 },
1112 CatalogColumn {
1113 compression: akar_common::enums::CompressionType::Uncompressed,
1114 name: "medium".into(),
1115 logical_type: LogicalTypeID::UInt32,
1116 is_primary_key: false,
1117 default_value: None,
1118 },
1119 CatalogColumn {
1120 compression: akar_common::enums::CompressionType::Uncompressed,
1121 name: "large".into(),
1122 logical_type: LogicalTypeID::UInt64,
1123 is_primary_key: false,
1124 default_value: None,
1125 },
1126 CatalogColumn {
1127 compression: akar_common::enums::CompressionType::Uncompressed,
1128 name: "temp".into(),
1129 logical_type: LogicalTypeID::Float,
1130 is_primary_key: false,
1131 default_value: None,
1132 },
1133 ];
1134
1135 let rows = read_parquet(
1136 parquet_path.to_str().unwrap(),
1137 &akar_common::file_system::VirtualFileSystemRegistry::new(),
1138 &columns,
1139 )
1140 .unwrap();
1141 assert_eq!(rows.len(), 2);
1142 assert_eq!(rows[0][0], Value::UInt8(100));
1143 assert_eq!(rows[0][1], Value::UInt32(1000));
1144 assert_eq!(rows[0][2], Value::UInt64(100000));
1145 if let Value::Float(f) = rows[0][3] {
1146 assert!((f - 36.5).abs() < 1e-5);
1147 } else {
1148 panic!("Expected Float");
1149 }
1150 assert_eq!(rows[1][0], Value::UInt8(200));
1151 assert_eq!(rows[1][2], Value::UInt64(u64::MAX));
1152 }
1153}