1mod 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
32pub 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 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 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 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 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]);