Skip to main content

connectorx/transports/
bigquery_arrowstream.rs

1//! Transport from BigQuery Source to Arrow Destination.
2
3use crate::{
4    destinations::arrowstream::{
5        typesystem::{
6            ArrowTypeSystem, DateTimeWrapperMicro, NaiveDateTimeWrapperMicro, NaiveTimeWrapperMicro,
7        },
8        ArrowDestination, ArrowDestinationError,
9    },
10    sources::bigquery::{BigQuerySource, BigQuerySourceError, BigQueryTypeSystem},
11    typesystem::TypeConversion,
12};
13use chrono::{DateTime, NaiveDate, NaiveDateTime, NaiveTime, Utc};
14use thiserror::Error;
15
16#[derive(Error, Debug)]
17pub enum BigQueryArrowTransportError {
18    #[error(transparent)]
19    Source(#[from] BigQuerySourceError),
20
21    #[error(transparent)]
22    Destination(#[from] ArrowDestinationError),
23
24    #[error(transparent)]
25    ConnectorX(#[from] crate::errors::ConnectorXError),
26}
27
28/// Convert BigQuery data types to Arrow data types.
29pub struct BigQueryArrowTransport;
30
31impl_transport!(
32    name = BigQueryArrowTransport,
33    error = BigQueryArrowTransportError,
34    systems = BigQueryTypeSystem => ArrowTypeSystem,
35    route = BigQuerySource => ArrowDestination,
36    mappings = {
37        { Bool[bool]                 => Boolean[bool]             | conversion auto }
38        { Boolean[bool]              => Boolean[bool]             | conversion none }
39        { Int64[i64]                 => Int64[i64]                | conversion auto }
40        { Integer[i64]               => Int64[i64]                | conversion none }
41        { Float64[f64]               => Float64[f64]              | conversion auto }
42        { Float[f64]                 => Float64[f64]              | conversion none }
43        { Numeric[f64]               => Float64[f64]              | conversion none }
44        { Bignumeric[f64]            => Float64[f64]              | conversion none }
45        { String[String]             => LargeUtf8[String]         | conversion auto }
46        { Bytes[String]              => LargeUtf8[String]         | conversion none }
47        { Date[NaiveDate]            => Date32[NaiveDate]         | conversion auto }
48        { Datetime[NaiveDateTime]    => Date64Micro[NaiveDateTimeWrapperMicro] | conversion option }
49        { Time[NaiveTime]            => Time64Micro[NaiveTimeWrapperMicro]     | conversion option }
50        { Timestamp[DateTime<Utc>]   => DateTimeTzMicro[DateTimeWrapperMicro]  | conversion option }
51    }
52);
53
54impl TypeConversion<NaiveDateTime, NaiveDateTimeWrapperMicro> for BigQueryArrowTransport {
55    fn convert(val: NaiveDateTime) -> NaiveDateTimeWrapperMicro {
56        NaiveDateTimeWrapperMicro(val)
57    }
58}
59
60impl TypeConversion<NaiveTime, NaiveTimeWrapperMicro> for BigQueryArrowTransport {
61    fn convert(val: NaiveTime) -> NaiveTimeWrapperMicro {
62        NaiveTimeWrapperMicro(val)
63    }
64}
65
66impl TypeConversion<DateTime<Utc>, DateTimeWrapperMicro> for BigQueryArrowTransport {
67    fn convert(val: DateTime<Utc>) -> DateTimeWrapperMicro {
68        DateTimeWrapperMicro(val)
69    }
70}