connectorx/transports/
trino_arrow.rs1use crate::{
4 destinations::arrow::{
5 typesystem::{ArrowTypeSystem, NaiveDateTimeWrapperMicro, NaiveTimeWrapperMicro},
6 ArrowDestination, ArrowDestinationError,
7 },
8 sources::trino::{TrinoSource, TrinoSourceError, TrinoTypeSystem},
9 typesystem::TypeConversion,
10};
11use chrono::{NaiveDate, NaiveDateTime, NaiveTime};
12use num_traits::ToPrimitive;
13use rust_decimal::Decimal;
14use serde_json::{to_string, Value};
15use thiserror::Error;
16
17#[derive(Error, Debug)]
18pub enum TrinoArrowTransportError {
19 #[error(transparent)]
20 Source(#[from] TrinoSourceError),
21
22 #[error(transparent)]
23 Destination(#[from] ArrowDestinationError),
24
25 #[error(transparent)]
26 ConnectorX(#[from] crate::errors::ConnectorXError),
27}
28
29pub struct TrinoArrowTransport();
31
32impl_transport!(
33 name = TrinoArrowTransport,
34 error = TrinoArrowTransportError,
35 systems = TrinoTypeSystem => ArrowTypeSystem,
36 route = TrinoSource => ArrowDestination,
37 mappings = {
38 { Date[NaiveDate] => Date32[NaiveDate] | conversion auto }
39 { Time[NaiveTime] => Time64Micro[NaiveTimeWrapperMicro] | conversion option }
40 { Timestamp[NaiveDateTime] => Date64Micro[NaiveDateTimeWrapperMicro] | conversion option }
41 { Boolean[bool] => Boolean[bool] | conversion auto }
42 { Bigint[i64] => Int64[i64] | conversion auto }
43 { Integer[i32] => Int64[i64] | conversion auto }
44 { Smallint[i16] => Int64[i64] | conversion auto }
45 { Tinyint[i8] => Int64[i64] | conversion auto }
46 { Double[f64] => Float64[f64] | conversion auto }
47 { Real[f32] => Float64[f64] | conversion auto }
48 { Varchar[String] => LargeUtf8[String] | conversion auto }
49 { Char[String] => LargeUtf8[String] | conversion none }
50 }
51);
52
53impl TypeConversion<Decimal, f64> for TrinoArrowTransport {
54 fn convert(val: Decimal) -> f64 {
55 val.to_f64()
56 .unwrap_or_else(|| panic!("cannot convert decimal {:?} to float64", val))
57 }
58}
59
60impl TypeConversion<Value, String> for TrinoArrowTransport {
61 fn convert(val: Value) -> String {
62 to_string(&val).unwrap()
63 }
64}
65
66impl TypeConversion<NaiveTime, NaiveTimeWrapperMicro> for TrinoArrowTransport {
67 fn convert(val: NaiveTime) -> NaiveTimeWrapperMicro {
68 NaiveTimeWrapperMicro(val)
69 }
70}
71
72impl TypeConversion<NaiveDateTime, NaiveDateTimeWrapperMicro> for TrinoArrowTransport {
73 fn convert(val: NaiveDateTime) -> NaiveDateTimeWrapperMicro {
74 NaiveDateTimeWrapperMicro(val)
75 }
76}
77
78#[cfg(test)]
79mod tests {
80 use super::*;
81
82 #[test]
83 fn bigint_conversion_preserves_full_i64_range() {
84 for value in [
85 i64::MIN,
86 -(1_i64 << 53) - 1,
87 (1_i64 << 53) + 1,
88 2_518_422_941_645_303_032,
89 i64::MAX,
90 ] {
91 assert_eq!(
92 <TrinoArrowTransport as TypeConversion<i64, i64>>::convert(value),
93 value
94 );
95 }
96 }
97}