Skip to main content

connectorx/transports/
trino_arrow.rs

1//! Transport from Trino Source to Arrow Destination.
2
3use 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
29/// Convert Trino data types to Arrow data types.
30pub 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}