Skip to main content

connectorx/sources/clickhouse/
mod.rs

1//! Source implementation for ClickHouse using the native protocol.
2
3mod errors;
4mod typesystem;
5
6pub use self::errors::ClickHouseSourceError;
7pub use self::typesystem::{ClickHouseTypeSystem, TypeMetadata};
8
9use crate::{
10    data_order::DataOrder,
11    errors::ConnectorXError,
12    sources::{
13        clickhouse::typesystem::DataType, PartitionParser, Produce, Source, SourcePartition,
14    },
15    sql::{count_query, limit0_query, CXQuery},
16};
17use anyhow::anyhow;
18use chrono::{DateTime, Duration, NaiveDate, NaiveTime, Utc};
19use chrono_tz::Tz;
20use clickhouse::Client;
21use fehler::{throw, throws};
22use rust_decimal::Decimal;
23use serde::Deserialize;
24use serde_json::Value as JsonValue;
25use sqlparser::dialect::{ClickHouseDialect, GenericDialect};
26use std::io::{Cursor, Read};
27use std::net::{IpAddr, Ipv4Addr, Ipv6Addr};
28use std::sync::Arc;
29use tokio::runtime::Runtime;
30use uuid::Uuid;
31
32/// ClickHouse source that uses the HTTP protocol.
33pub struct ClickHouseSource {
34    rt: Arc<Runtime>,
35    pub client: Client,
36    origin_query: Option<String>,
37    queries: Vec<CXQuery<String>>,
38    names: Vec<String>,
39    schema: Vec<ClickHouseTypeSystem>,
40    metadata: Vec<TypeMetadata>,
41}
42
43impl ClickHouseSource {
44    #[throws(ClickHouseSourceError)]
45    pub fn new(rt: Arc<Runtime>, conn: &str) -> Self {
46        let url = url::Url::parse(conn)?;
47
48        let use_https = url
49            .query_pairs()
50            .find(|(k, v)| k == "protocol" && v == "https")
51            .is_some();
52
53        let base_url = format!(
54            "{}://{}:{}",
55            if use_https { "https" } else { "http" },
56            url.host_str().unwrap_or("localhost"),
57            url.port().unwrap_or(8123)
58        );
59
60        // clickhouse 0.13 decodes LZ4 only; newer servers default to ZSTD.
61        let mut client = Client::default()
62            .with_url(&base_url)
63            .with_option("network_compression_method", "lz4");
64
65        let database = url.path().trim_start_matches('/');
66        if !database.is_empty() {
67            client = client.with_database(database);
68        }
69
70        let username = url.username();
71        if !username.is_empty() {
72            client = client.with_user(username);
73        }
74
75        let password = url.password().unwrap_or("");
76        if !password.is_empty() {
77            client = client.with_password(password);
78        }
79
80        Self {
81            rt,
82            client,
83            origin_query: None,
84            queries: vec![],
85            names: vec![],
86            schema: vec![],
87            metadata: vec![],
88        }
89    }
90}
91
92impl Source for ClickHouseSource
93where
94    ClickHouseSourcePartition:
95        SourcePartition<TypeSystem = ClickHouseTypeSystem, Error = ClickHouseSourceError>,
96{
97    const DATA_ORDERS: &'static [DataOrder] = &[DataOrder::RowMajor];
98    type Partition = ClickHouseSourcePartition;
99    type TypeSystem = ClickHouseTypeSystem;
100    type Error = ClickHouseSourceError;
101
102    #[throws(ClickHouseSourceError)]
103    fn set_data_order(&mut self, data_order: DataOrder) {
104        if !matches!(data_order, DataOrder::RowMajor) {
105            throw!(ConnectorXError::UnsupportedDataOrder(data_order));
106        }
107    }
108
109    fn set_queries<Q: ToString>(&mut self, queries: &[CXQuery<Q>]) {
110        self.queries = queries.iter().map(|q| q.map(Q::to_string)).collect();
111    }
112
113    fn set_origin_query(&mut self, query: Option<String>) {
114        self.origin_query = query;
115    }
116
117    #[throws(ClickHouseSourceError)]
118    fn fetch_metadata(&mut self) {
119        assert!(!self.queries.is_empty());
120
121        let first_query = &self.queries[0];
122        let l1query = limit0_query(first_query, &ClickHouseDialect {})?;
123
124        let describe_query = format!("DESCRIBE ({})", l1query.as_str());
125
126        let response = self.rt.block_on(async {
127            let mut cursor = self
128                .client
129                .query(&describe_query)
130                .fetch_bytes("JSONCompact")
131                .map_err(|e| anyhow!("ClickHouse error: {}", e))?;
132            let bytes = cursor
133                .collect()
134                .await
135                .map_err(|e| anyhow!("ClickHouse error: {}", e))?;
136            Ok::<_, ClickHouseSourceError>(bytes)
137        })?;
138
139        #[derive(Debug, Deserialize)]
140        struct DescribeResponse {
141            data: Vec<Vec<JsonValue>>,
142        }
143
144        let parsed: DescribeResponse = serde_json::from_slice(&response)
145            .map_err(|e| anyhow!("Failed to parse DESCRIBE response: {}", e))?;
146
147        let mut names = Vec::new();
148        let mut types = Vec::new();
149        let mut metadata = Vec::new();
150
151        for row in parsed.data {
152            if row.len() >= 2 {
153                let name = row[0].as_str().unwrap_or("").to_string();
154                let type_str = row[1].as_str().unwrap_or("String");
155                let (ts, meta) = ClickHouseTypeSystem::from_type_str_with_metadata(type_str);
156                names.push(name);
157                types.push(ts);
158                metadata.push(meta);
159            }
160        }
161
162        self.names = names;
163        self.schema = types;
164        self.metadata = metadata;
165    }
166
167    #[throws(ClickHouseSourceError)]
168    fn result_rows(&mut self) -> Option<usize> {
169        match &self.origin_query {
170            Some(q) => {
171                let cxq = CXQuery::Naked(q.clone());
172                let cquery = count_query(&cxq, &ClickHouseDialect {})?;
173
174                let response = self.rt.block_on(async {
175                    let mut cursor = self
176                        .client
177                        .query(cquery.as_str())
178                        .fetch_bytes("JSONCompact")
179                        .map_err(|e| anyhow!("ClickHouse error: {}", e))?;
180                    let bytes = cursor
181                        .collect()
182                        .await
183                        .map_err(|e| anyhow!("ClickHouse error: {}", e))?;
184                    Ok::<_, ClickHouseSourceError>(bytes)
185                })?;
186
187                #[derive(Debug, Deserialize)]
188                struct CountResponse {
189                    data: Vec<Vec<JsonValue>>,
190                }
191
192                let parsed: CountResponse = serde_json::from_slice(&response)
193                    .map_err(|e| anyhow!("Failed to parse count response: {}", e))?;
194
195                if let Some(row) = parsed.data.first() {
196                    if let Some(count_val) = row.first() {
197                        let count = match count_val {
198                            JsonValue::Number(n) => n.as_u64().unwrap_or(0),
199                            JsonValue::String(s) => s.parse().unwrap_or(0),
200                            _ => 0,
201                        };
202                        return Some(count as usize);
203                    }
204                }
205                None
206            }
207            None => None,
208        }
209    }
210
211    fn names(&self) -> Vec<String> {
212        self.names.clone()
213    }
214
215    fn schema(&self) -> Vec<Self::TypeSystem> {
216        self.schema.clone()
217    }
218
219    #[throws(ClickHouseSourceError)]
220    fn partition(self) -> Vec<Self::Partition> {
221        let mut ret = vec![];
222        for query in self.queries {
223            ret.push(ClickHouseSourcePartition::new(
224                self.rt.clone(),
225                self.client.clone(),
226                &query,
227                &self.schema,
228                &self.metadata,
229            ));
230        }
231        ret
232    }
233}
234
235pub struct ClickHouseSourcePartition {
236    rt: Arc<Runtime>,
237    client: Client,
238    query: CXQuery<String>,
239    schema: Vec<ClickHouseTypeSystem>,
240    metadata: Vec<TypeMetadata>,
241    nrows: usize,
242    ncols: usize,
243}
244
245impl ClickHouseSourcePartition {
246    pub fn new(
247        rt: Arc<Runtime>,
248        client: Client,
249        query: &CXQuery<String>,
250        schema: &[ClickHouseTypeSystem],
251        metadata: &[TypeMetadata],
252    ) -> Self {
253        Self {
254            rt,
255            client,
256            query: query.clone(),
257            schema: schema.to_vec(),
258            metadata: metadata.to_vec(),
259            nrows: 0,
260            ncols: schema.len(),
261        }
262    }
263}
264
265impl SourcePartition for ClickHouseSourcePartition {
266    type TypeSystem = ClickHouseTypeSystem;
267    type Parser<'a> = ClickHouseSourceParser<'a>;
268    type Error = ClickHouseSourceError;
269
270    #[throws(ClickHouseSourceError)]
271    fn result_rows(&mut self) {
272        let cquery = count_query(&self.query, &GenericDialect {})?;
273
274        let response = self.rt.block_on(async {
275            let mut cursor = self
276                .client
277                .query(cquery.as_str())
278                .fetch_bytes("JSONCompact")
279                .map_err(|e| anyhow!("ClickHouse error: {}", e))?;
280            let bytes = cursor
281                .collect()
282                .await
283                .map_err(|e| anyhow!("ClickHouse error: {}", e))?;
284            Ok::<_, ClickHouseSourceError>(bytes)
285        })?;
286
287        #[derive(Debug, Deserialize)]
288        struct CountResponse {
289            data: Vec<Vec<JsonValue>>,
290        }
291
292        let parsed: CountResponse = serde_json::from_slice(&response)
293            .map_err(|e| anyhow!("Failed to parse count response: {}", e))?;
294
295        if let Some(row) = parsed.data.first() {
296            if let Some(count_val) = row.first() {
297                let count = match count_val {
298                    JsonValue::Number(n) => n.as_u64().unwrap_or(0),
299                    JsonValue::String(s) => s.parse().unwrap_or(0),
300                    _ => 0,
301                };
302                self.nrows = count as usize;
303            }
304        }
305    }
306
307    #[throws(ClickHouseSourceError)]
308    fn parser(&mut self) -> Self::Parser<'_> {
309        ClickHouseSourceParser::new(
310            self.rt.clone(),
311            self.client.clone(),
312            self.query.clone(),
313            &self.schema,
314            &self.metadata,
315        )?
316    }
317
318    fn nrows(&self) -> usize {
319        self.nrows
320    }
321
322    fn ncols(&self) -> usize {
323        self.ncols
324    }
325}
326
327struct BinaryReader<'a> {
328    cursor: Cursor<&'a [u8]>,
329}
330
331impl<'a> BinaryReader<'a> {
332    fn new(data: &'a [u8]) -> Self {
333        Self {
334            cursor: Cursor::new(data),
335        }
336    }
337
338    fn read_bytes<const N: usize>(&mut self) -> Result<[u8; N], ClickHouseSourceError> {
339        let mut buf = [0u8; N];
340        self.cursor
341            .read_exact(&mut buf)
342            .map_err(|e| anyhow!("Failed to read {} bytes: {}", N, e))?;
343        Ok(buf)
344    }
345
346    fn read_u8(&mut self) -> Result<u8, ClickHouseSourceError> {
347        Ok(self.read_bytes::<1>()?[0])
348    }
349
350    fn read_i8(&mut self) -> Result<i8, ClickHouseSourceError> {
351        Ok(self.read_u8()? as i8)
352    }
353
354    fn read_u16(&mut self) -> Result<u16, ClickHouseSourceError> {
355        Ok(u16::from_le_bytes(self.read_bytes()?))
356    }
357
358    fn read_i16(&mut self) -> Result<i16, ClickHouseSourceError> {
359        Ok(i16::from_le_bytes(self.read_bytes()?))
360    }
361
362    fn read_u32(&mut self) -> Result<u32, ClickHouseSourceError> {
363        Ok(u32::from_le_bytes(self.read_bytes()?))
364    }
365
366    fn read_i32(&mut self) -> Result<i32, ClickHouseSourceError> {
367        Ok(i32::from_le_bytes(self.read_bytes()?))
368    }
369
370    fn read_u64(&mut self) -> Result<u64, ClickHouseSourceError> {
371        Ok(u64::from_le_bytes(self.read_bytes()?))
372    }
373
374    fn read_i64(&mut self) -> Result<i64, ClickHouseSourceError> {
375        Ok(i64::from_le_bytes(self.read_bytes()?))
376    }
377
378    fn read_f32(&mut self) -> Result<f32, ClickHouseSourceError> {
379        Ok(f32::from_le_bytes(self.read_bytes()?))
380    }
381
382    fn read_f64(&mut self) -> Result<f64, ClickHouseSourceError> {
383        Ok(f64::from_le_bytes(self.read_bytes()?))
384    }
385
386    fn read_varint(&mut self) -> Result<u64, ClickHouseSourceError> {
387        let mut result: u64 = 0;
388        let mut shift = 0;
389        loop {
390            let byte = self.read_u8()?;
391            result |= ((byte & 0x7f) as u64) << shift;
392            if byte & 0x80 == 0 {
393                break;
394            }
395            shift += 7;
396            if shift >= 64 {
397                return Err(anyhow!("Varint too long").into());
398            }
399        }
400        Ok(result)
401    }
402
403    fn read_fixed_string(&mut self, len: usize) -> Result<Vec<u8>, ClickHouseSourceError> {
404        let mut buf = vec![0u8; len];
405        self.cursor
406            .read_exact(&mut buf)
407            .map_err(|e| anyhow!("Failed to read FixedString: {}", e))?;
408        Ok(buf)
409    }
410
411    fn read_string(&mut self) -> Result<String, ClickHouseSourceError> {
412        let len = self.read_varint()? as usize;
413        let mut buf = vec![0u8; len];
414        self.cursor
415            .read_exact(&mut buf)
416            .map_err(|e| anyhow!("Failed to read string: {}", e))?;
417        String::from_utf8(buf).map_err(|e| anyhow!("Invalid UTF-8: {}", e).into())
418    }
419
420    fn read_uuid(&mut self) -> Result<Uuid, ClickHouseSourceError> {
421        // ClickHouse stores UUID as two UInt64 in big-endian order
422        let high = self.read_u64()?;
423        let low = self.read_u64()?;
424
425        Ok(Uuid::from_u64_pair(high, low))
426    }
427
428    fn read_bool(&mut self) -> Result<bool, ClickHouseSourceError> {
429        Ok(self.read_u8()? != 0)
430    }
431
432    fn read_date(&mut self) -> Result<NaiveDate, ClickHouseSourceError> {
433        // ClickHouse stores Date as UInt16 representing days since 1970-01-01
434        let days = self.read_u16()? as i64;
435        let epoch = NaiveDate::from_ymd_opt(1970, 1, 1).unwrap();
436        epoch
437            .checked_add_signed(Duration::days(days))
438            .ok_or_else(|| anyhow!("Invalid date value: {} days since epoch", days).into())
439    }
440
441    fn read_date32(&mut self) -> Result<NaiveDate, ClickHouseSourceError> {
442        // ClickHouse stores Date32 as Int32 representing days since 1970-01-01
443        // negative values represent dates before 1970-01-01
444        let days = self.read_i32()? as i64;
445        let epoch = NaiveDate::from_ymd_opt(1970, 1, 1).unwrap();
446        epoch
447            .checked_add_signed(Duration::days(days))
448            .ok_or_else(|| anyhow!("Invalid date32 value: {} days since epoch", days).into())
449    }
450
451    fn read_datetime(&mut self, tz: Option<&Tz>) -> Result<DateTime<Utc>, ClickHouseSourceError> {
452        let seconds = self.read_u32()? as i64;
453        chrono::DateTime::from_timestamp(seconds, 0)
454            .map(|dt| dt.with_timezone(tz.unwrap_or(&Tz::UTC)).to_utc())
455            .ok_or_else(|| {
456                anyhow!("Invalid datetime value: {} seconds since epoch", seconds).into()
457            })
458    }
459
460    fn read_datetime64(
461        &mut self,
462        precision: u8,
463        tz: Option<&Tz>,
464    ) -> Result<DateTime<Utc>, ClickHouseSourceError> {
465        if precision > 9 {
466            return Err(anyhow!("Unsupported DateTime64 precision: {}", precision).into());
467        }
468        let ticks = self.read_i64()?;
469        let scale = 10_i64.pow(precision as u32);
470        let seconds = ticks.div_euclid(scale);
471        let fractional_ticks = ticks.rem_euclid(scale) as u32;
472        let nanos = fractional_ticks * 10_u32.pow(9 - precision as u32);
473
474        DateTime::from_timestamp(seconds, nanos)
475            .map(|dt| dt.with_timezone(tz.unwrap_or(&Tz::UTC)).to_utc())
476            .ok_or_else(|| anyhow!("Invalid datetime64 value: {} ticks", ticks).into())
477    }
478
479    fn read_decimal(&mut self, precision: u8, scale: u8) -> Result<Decimal, ClickHouseSourceError> {
480        match precision {
481            1..=9 => {
482                let value = self.read_i32()?;
483                Ok(Decimal::new(value as i64, scale as u32))
484            }
485            10..=18 => {
486                let value = self.read_i64()?;
487                Ok(Decimal::new(value, scale as u32))
488            }
489            _ => Err(anyhow!("Unsupported Decimal precision: {}", precision).into()),
490        }
491    }
492
493    fn read_time(&mut self) -> Result<NaiveTime, ClickHouseSourceError> {
494        let seconds = self.read_u32()? as i64;
495        Ok(
496            NaiveTime::from_num_seconds_from_midnight_opt(seconds as u32, 0)
497                .ok_or_else(|| anyhow!("Invalid time value: {} seconds since midnight", seconds))?,
498        )
499    }
500
501    fn read_time64(&mut self, precision: u8) -> Result<NaiveTime, ClickHouseSourceError> {
502        if precision > 9 {
503            return Err(anyhow!("Unsupported Time64 precision: {}", precision).into());
504        }
505        let ticks = self.read_i64()?;
506        let nanos = ticks * 10_i64.pow(9 - precision as u32);
507        Ok(NaiveTime::from_num_seconds_from_midnight_opt(
508            (nanos / 1_000_000_000) as u32,
509            (nanos % 1_000_000_000) as u32,
510        )
511        .ok_or_else(|| anyhow!("Invalid time64 value: {} ticks since midnight", ticks))?)
512    }
513
514    fn read_ipv4(&mut self) -> Result<IpAddr, ClickHouseSourceError> {
515        let bytes = self.read_u32()?;
516
517        Ok(IpAddr::V4(Ipv4Addr::from_bits(bytes)))
518    }
519
520    fn read_ipv6(&mut self) -> Result<IpAddr, ClickHouseSourceError> {
521        let seg1 = u16::from_be(self.read_u16()?);
522        let seg2 = u16::from_be(self.read_u16()?);
523        let seg3 = u16::from_be(self.read_u16()?);
524        let seg4 = u16::from_be(self.read_u16()?);
525        let seg5 = u16::from_be(self.read_u16()?);
526        let seg6 = u16::from_be(self.read_u16()?);
527        let seg7 = u16::from_be(self.read_u16()?);
528        let seg8 = u16::from_be(self.read_u16()?);
529
530        Ok(IpAddr::V6(Ipv6Addr::from_segments([
531            seg1, seg2, seg3, seg4, seg5, seg6, seg7, seg8,
532        ])))
533    }
534
535    fn read_enum8(&mut self) -> Result<i8, ClickHouseSourceError> {
536        self.read_i8()
537    }
538
539    fn read_enum16(&mut self) -> Result<i16, ClickHouseSourceError> {
540        self.read_i16()
541    }
542
543    fn read_array<T, F>(&mut self, read_elem: F) -> Result<Vec<Option<T>>, ClickHouseSourceError>
544    where
545        F: Fn(&mut Self) -> Result<T, ClickHouseSourceError>,
546    {
547        let len = self.read_varint()? as usize;
548        let mut result = Vec::with_capacity(len);
549        for _ in 0..len {
550            result.push(Some(read_elem(self)?));
551        }
552        Ok(result)
553    }
554
555    fn is_empty(&self) -> bool {
556        self.cursor.position() as usize >= self.cursor.get_ref().len()
557    }
558}
559
560pub struct ClickHouseSourceParser<'a> {
561    rt: Arc<Runtime>,
562    client: Client,
563    query: CXQuery<String>,
564    schema: Vec<ClickHouseTypeSystem>,
565    metadata: Vec<TypeMetadata>,
566    rowbuf: Vec<Vec<DataType>>,
567    ncols: usize,
568    current_row: usize,
569    current_col: usize,
570    is_finished: bool,
571    _phantom: std::marker::PhantomData<&'a ()>,
572}
573
574impl<'a> ClickHouseSourceParser<'a> {
575    #[throws(ClickHouseSourceError)]
576    pub fn new(
577        rt: Arc<Runtime>,
578        client: Client,
579        query: CXQuery<String>,
580        schema: &[ClickHouseTypeSystem],
581        metadata: &[TypeMetadata],
582    ) -> Self {
583        Self {
584            rt,
585            client,
586            query,
587            schema: schema.to_vec(),
588            metadata: metadata.to_vec(),
589            rowbuf: Vec::new(),
590            ncols: schema.len(),
591            current_row: 0,
592            current_col: 0,
593            is_finished: false,
594            _phantom: std::marker::PhantomData,
595        }
596    }
597
598    fn parse_row_binary(
599        &self,
600        reader: &mut BinaryReader,
601    ) -> Result<Vec<DataType>, ClickHouseSourceError> {
602        let mut row = Vec::with_capacity(self.ncols);
603
604        for (col_idx, col_type) in self.schema.iter().enumerate() {
605            let is_nullable = col_type.is_nullable();
606            let meta = &self.metadata[col_idx];
607
608            if is_nullable {
609                let null_flag = reader.read_u8()?;
610                if null_flag == 1 {
611                    row.push(DataType::Null);
612                    continue;
613                }
614            }
615
616            let value = match col_type {
617                ClickHouseTypeSystem::Int8(_) => DataType::Int8(reader.read_i8()?),
618                ClickHouseTypeSystem::Int16(_) => DataType::Int16(reader.read_i16()?),
619                ClickHouseTypeSystem::Int32(_) => DataType::Int32(reader.read_i32()?),
620                ClickHouseTypeSystem::Int64(_) => DataType::Int64(reader.read_i64()?),
621                ClickHouseTypeSystem::UInt8(_) => DataType::UInt8(reader.read_u8()?),
622                ClickHouseTypeSystem::UInt16(_) => DataType::UInt16(reader.read_u16()?),
623                ClickHouseTypeSystem::UInt32(_) => DataType::UInt32(reader.read_u32()?),
624                ClickHouseTypeSystem::UInt64(_) => DataType::UInt64(reader.read_u64()?),
625
626                ClickHouseTypeSystem::Float32(_) => DataType::Float32(reader.read_f32()?),
627                ClickHouseTypeSystem::Float64(_) => DataType::Float64(reader.read_f64()?),
628
629                ClickHouseTypeSystem::Decimal(_) => {
630                    DataType::Decimal(reader.read_decimal(meta.precision, meta.scale)?)
631                }
632
633                ClickHouseTypeSystem::String(_) => DataType::String(reader.read_string()?),
634
635                ClickHouseTypeSystem::FixedString(_) => {
636                    DataType::FixedString(reader.read_fixed_string(meta.length)?)
637                }
638
639                ClickHouseTypeSystem::Date(_) => DataType::Date(reader.read_date()?),
640                ClickHouseTypeSystem::Date32(_) => DataType::Date32(reader.read_date32()?),
641
642                ClickHouseTypeSystem::DateTime(_) => {
643                    DataType::DateTime(reader.read_datetime(meta.timezone.as_ref())?)
644                }
645                ClickHouseTypeSystem::DateTime64(_) => DataType::DateTime64(
646                    reader.read_datetime64(meta.precision, meta.timezone.as_ref())?,
647                ),
648
649                ClickHouseTypeSystem::Time(_) => DataType::Time(reader.read_time()?),
650                ClickHouseTypeSystem::Time64(_) => {
651                    DataType::Time64(reader.read_time64(meta.precision)?)
652                }
653
654                ClickHouseTypeSystem::Enum8(_) => {
655                    let enum_value = reader.read_enum8()?;
656                    let enum_str = meta
657                        .named_values
658                        .as_ref()
659                        .and_then(|h| h.get(&(enum_value as i16)));
660                    DataType::Enum8(enum_str.cloned().unwrap_or_default())
661                }
662                ClickHouseTypeSystem::Enum16(_) => {
663                    let enum_value = reader.read_enum16()?;
664                    let enum_str = meta.named_values.as_ref().and_then(|h| h.get(&enum_value));
665                    DataType::Enum16(enum_str.cloned().unwrap_or_default())
666                }
667
668                ClickHouseTypeSystem::UUID(_) => DataType::UUID(reader.read_uuid()?),
669
670                ClickHouseTypeSystem::IPv4(_) => DataType::IPv4(reader.read_ipv4()?),
671                ClickHouseTypeSystem::IPv6(_) => DataType::IPv6(reader.read_ipv6()?),
672
673                ClickHouseTypeSystem::Bool(_) => DataType::Bool(reader.read_bool()?),
674
675                ClickHouseTypeSystem::ArrayBool(_) => {
676                    DataType::ArrayBool(reader.read_array(BinaryReader::read_bool)?)
677                }
678                ClickHouseTypeSystem::ArrayString(_) => {
679                    DataType::ArrayString(reader.read_array(BinaryReader::read_string)?)
680                }
681                ClickHouseTypeSystem::ArrayInt8(_) => {
682                    DataType::ArrayInt8(reader.read_array(BinaryReader::read_i8)?)
683                }
684                ClickHouseTypeSystem::ArrayInt16(_) => {
685                    DataType::ArrayInt16(reader.read_array(BinaryReader::read_i16)?)
686                }
687                ClickHouseTypeSystem::ArrayInt32(_) => {
688                    DataType::ArrayInt32(reader.read_array(BinaryReader::read_i32)?)
689                }
690                ClickHouseTypeSystem::ArrayInt64(_) => {
691                    DataType::ArrayInt64(reader.read_array(BinaryReader::read_i64)?)
692                }
693                ClickHouseTypeSystem::ArrayUInt8(_) => {
694                    DataType::ArrayUInt8(reader.read_array(BinaryReader::read_u8)?)
695                }
696                ClickHouseTypeSystem::ArrayUInt16(_) => {
697                    DataType::ArrayUInt16(reader.read_array(BinaryReader::read_u16)?)
698                }
699                ClickHouseTypeSystem::ArrayUInt32(_) => {
700                    DataType::ArrayUInt32(reader.read_array(BinaryReader::read_u32)?)
701                }
702                ClickHouseTypeSystem::ArrayUInt64(_) => {
703                    DataType::ArrayUInt64(reader.read_array(BinaryReader::read_u64)?)
704                }
705                ClickHouseTypeSystem::ArrayFloat32(_) => {
706                    DataType::ArrayFloat32(reader.read_array(BinaryReader::read_f32)?)
707                }
708                ClickHouseTypeSystem::ArrayFloat64(_) => {
709                    DataType::ArrayFloat64(reader.read_array(BinaryReader::read_f64)?)
710                }
711                ClickHouseTypeSystem::ArrayDecimal(_) => {
712                    let precision = meta.precision;
713                    let scale = meta.scale;
714                    DataType::ArrayDecimal(reader.read_array(|r| r.read_decimal(precision, scale))?)
715                }
716            };
717
718            row.push(value);
719        }
720
721        Ok(row)
722    }
723
724    #[throws(ClickHouseSourceError)]
725    fn next_loc(&mut self) -> (usize, usize) {
726        let ret = (self.current_row, self.current_col);
727        self.current_row += (self.current_col + 1) / self.ncols;
728        self.current_col = (self.current_col + 1) % self.ncols;
729        ret
730    }
731}
732
733impl<'a> PartitionParser<'a> for ClickHouseSourceParser<'a> {
734    type TypeSystem = ClickHouseTypeSystem;
735    type Error = ClickHouseSourceError;
736
737    #[throws(ClickHouseSourceError)]
738    fn fetch_next(&mut self) -> (usize, bool) {
739        assert!(self.current_col == 0);
740
741        if self.is_finished {
742            return (0, true);
743        }
744
745        let response = self.rt.block_on(async {
746            let mut cursor = self
747                .client
748                .query(self.query.as_str())
749                .fetch_bytes("RowBinary")
750                .map_err(|e| anyhow!("ClickHouse error: {}", e))?;
751            let bytes = cursor
752                .collect()
753                .await
754                .map_err(|e| anyhow!("ClickHouse error: {}", e))?;
755            Ok::<_, ClickHouseSourceError>(bytes)
756        })?;
757        let mut reader = BinaryReader::new(&response);
758        let mut rows = Vec::new();
759
760        while !reader.is_empty() {
761            match self.parse_row_binary(&mut reader) {
762                Ok(row) => rows.push(row),
763                Err(_) => break,
764            }
765        }
766
767        self.rowbuf = rows;
768        self.current_row = 0;
769        self.is_finished = true;
770
771        (self.rowbuf.len(), true)
772    }
773}
774
775macro_rules! impl_produce {
776    ($rust_type:ty, [$($variant:ident),+]) => {
777        impl<'r, 'a> Produce<'r, $rust_type> for ClickHouseSourceParser<'a> {
778            type Error = ClickHouseSourceError;
779
780            #[throws(ClickHouseSourceError)]
781            fn produce(&'r mut self) -> $rust_type {
782                let (ridx, cidx) = self.next_loc()?;
783                let value = &self.rowbuf[ridx][cidx];
784
785                match value {
786                    $(DataType::$variant(v) => *v as $rust_type,)+
787                    _ => throw!(ConnectorXError::cannot_produce::<$rust_type>(Some(
788                        format!("{:?}", value)
789                    ))),
790                }
791            }
792        }
793
794        impl<'r, 'a> Produce<'r, Option<$rust_type>> for ClickHouseSourceParser<'a> {
795            type Error = ClickHouseSourceError;
796
797            #[throws(ClickHouseSourceError)]
798            fn produce(&'r mut self) -> Option<$rust_type> {
799                let (ridx, cidx) = self.next_loc()?;
800                let value = &self.rowbuf[ridx][cidx];
801
802                match value {
803                    DataType::Null => None,
804                    $(DataType::$variant(v) => Some(*v as $rust_type),)+
805                    _ => throw!(ConnectorXError::cannot_produce::<$rust_type>(Some(
806                        format!("{:?}", value)
807                    ))),
808                }
809            }
810        }
811    };
812}
813
814macro_rules! impl_produce_with_clone {
815    ($rust_type:ty, [$($variant:ident),+]) => {
816        impl<'r, 'a> Produce<'r, $rust_type> for ClickHouseSourceParser<'a> {
817            type Error = ClickHouseSourceError;
818
819            #[throws(ClickHouseSourceError)]
820            fn produce(&'r mut self) -> $rust_type {
821                let (ridx, cidx) = self.next_loc()?;
822                let value = &self.rowbuf[ridx][cidx];
823
824                match value {
825                    $(DataType::$variant(v) => v.clone() as $rust_type,)+
826                    _ => throw!(ConnectorXError::cannot_produce::<$rust_type>(Some(
827                        format!("{:?}", value)
828                    ))),
829                }
830            }
831        }
832
833        impl<'r, 'a> Produce<'r, Option<$rust_type>> for ClickHouseSourceParser<'a> {
834            type Error = ClickHouseSourceError;
835
836            #[throws(ClickHouseSourceError)]
837            fn produce(&'r mut self) -> Option<$rust_type> {
838                let (ridx, cidx) = self.next_loc()?;
839                let value = &self.rowbuf[ridx][cidx];
840
841                match value {
842                    DataType::Null => None,
843                    $(DataType::$variant(v) => Some(v.clone() as $rust_type),)+
844                    _ => throw!(ConnectorXError::cannot_produce::<$rust_type>(Some(
845                        format!("{:?}", value)
846                    ))),
847                }
848            }
849        }
850    };
851}
852
853macro_rules! impl_produce_vec {
854    ($rust_type:ty, [$($variant:ident),+]) => {
855        impl<'r, 'a> Produce<'r, Vec<Option<$rust_type>>> for ClickHouseSourceParser<'a> {
856            type Error = ClickHouseSourceError;
857
858            #[throws(ClickHouseSourceError)]
859            fn produce(&'r mut self) -> Vec<Option<$rust_type>> {
860                let (ridx, cidx) = self.next_loc()?;
861                let value = &self.rowbuf[ridx][cidx];
862
863                match value {
864                    $(DataType::$variant(v) => v.clone(),)+
865                    _ => throw!(ConnectorXError::cannot_produce::<Vec<Option<$rust_type>>>(Some(
866                        format!("{:?}", value)
867                    ))),
868                }
869            }
870        }
871
872        impl<'r, 'a> Produce<'r, Option<Vec<Option<$rust_type>>>> for ClickHouseSourceParser<'a> {
873            type Error = ClickHouseSourceError;
874
875            #[throws(ClickHouseSourceError)]
876            fn produce(&'r mut self) -> Option<Vec<Option<$rust_type>>> {
877                let (ridx, cidx) = self.next_loc()?;
878                let value = &self.rowbuf[ridx][cidx];
879
880                match value {
881                    DataType::Null => None,
882                    $(DataType::$variant(v) => Some(v.clone()),)+
883                    _ => throw!(ConnectorXError::cannot_produce::<Option<Vec<Option<$rust_type>>>>(Some(
884                        format!("{:?}", value)
885                    ))),
886                }
887            }
888        }
889    };
890}
891
892impl_produce!(i8, [Int8]);
893impl_produce!(i16, [Int16]);
894impl_produce!(i32, [Int32]);
895impl_produce!(i64, [Int64]);
896impl_produce!(u8, [UInt8]);
897impl_produce!(u16, [UInt16]);
898impl_produce!(u32, [UInt32]);
899impl_produce!(u64, [UInt64]);
900impl_produce!(f32, [Float32]);
901impl_produce!(f64, [Float64]);
902impl_produce!(Decimal, [Decimal]);
903impl_produce_with_clone!(String, [String, Enum8, Enum16]);
904impl_produce_with_clone!(Vec<u8>, [FixedString]);
905impl_produce!(NaiveDate, [Date, Date32]);
906impl_produce!(DateTime<Utc>, [DateTime, DateTime64]);
907impl_produce!(NaiveTime, [Time, Time64]);
908impl_produce!(Uuid, [UUID]);
909impl_produce!(IpAddr, [IPv4, IPv6]);
910impl_produce!(bool, [Bool]);
911impl_produce_vec!(bool, [ArrayBool]);
912impl_produce_vec!(String, [ArrayString]);
913impl_produce_vec!(i8, [ArrayInt8]);
914impl_produce_vec!(i16, [ArrayInt16]);
915impl_produce_vec!(i32, [ArrayInt32]);
916impl_produce_vec!(i64, [ArrayInt64]);
917impl_produce_vec!(u8, [ArrayUInt8]);
918impl_produce_vec!(u16, [ArrayUInt16]);
919impl_produce_vec!(u32, [ArrayUInt32]);
920impl_produce_vec!(u64, [ArrayUInt64]);
921impl_produce_vec!(f32, [ArrayFloat32]);
922impl_produce_vec!(f64, [ArrayFloat64]);
923impl_produce_vec!(Decimal, [ArrayDecimal]);