Skip to main content

risingwave_connector/sink/
clickhouse.rs

1// Copyright 2023 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use core::fmt::Debug;
16use core::num::NonZeroU64;
17use std::collections::{BTreeMap, HashMap, HashSet};
18
19use anyhow::anyhow;
20use clickhouse::insert::Insert;
21use clickhouse::{Client as ClickHouseClient, Row as ClickHouseRow};
22use itertools::Itertools;
23use phf::{Set, phf_set};
24use risingwave_common::array::{Op, StreamChunk};
25use risingwave_common::catalog::{FieldLike, Schema};
26use risingwave_common::row::Row;
27use risingwave_common::types::{DataType, Decimal, ScalarRefImpl, Serial};
28use risingwave_common::util::iter_util::ZipEqDebug;
29use serde::ser::{SerializeSeq, SerializeStruct};
30use serde::{Deserialize, Serialize};
31use serde_with::{DisplayFromStr, serde_as};
32use thiserror_ext::AsReport;
33use tonic::async_trait;
34use tracing::warn;
35use with_options::WithOptions;
36
37use super::decouple_checkpoint_log_sink::{
38    DecoupleCheckpointLogSinkerOf, default_commit_checkpoint_interval,
39};
40use super::writer::SinkWriter;
41use super::{SinkWriterMetrics, SinkWriterParam};
42use crate::enforce_secret::EnforceSecret;
43use crate::error::ConnectorResult;
44use crate::sink::{
45    Result, SINK_TYPE_APPEND_ONLY, SINK_TYPE_OPTION, SINK_TYPE_UPSERT, Sink, SinkError, SinkParam,
46};
47
48const QUERY_ENGINE: &str =
49    "select distinct ?fields from system.tables where database = ? and name = ?";
50const QUERY_COLUMN: &str =
51    "select distinct ?fields from system.columns where database = ? and table = ? order by ?";
52pub const CLICKHOUSE_SINK: &str = "clickhouse";
53
54const ALLOW_EXPERIMENTAL_JSON_TYPE: &str = "allow_experimental_json_type";
55const INPUT_FORMAT_BINARY_READ_JSON_AS_STRING: &str = "input_format_binary_read_json_as_string";
56const OUTPUT_FORMAT_BINARY_WRITE_JSON_AS_STRING: &str = "output_format_binary_write_json_as_string";
57
58/// `ClickHouse` backs `Decimal(P, S)` with `Decimal256` once `P > 38`, which takes 32 bytes on the
59/// wire, while [`ClickHouseDecimal`] never writes more than 16.
60const MAX_ENCODABLE_DECIMAL_PRECISION: u8 = 38;
61
62#[serde_as]
63#[derive(Deserialize, Debug, Clone, WithOptions)]
64pub struct ClickHouseCommon {
65    #[serde(rename = "clickhouse.url")]
66    pub url: String,
67    #[serde(rename = "clickhouse.user")]
68    pub user: String,
69    #[serde(rename = "clickhouse.password")]
70    pub password: String,
71    #[serde(rename = "clickhouse.database")]
72    pub database: String,
73    #[serde(rename = "clickhouse.table")]
74    pub table: String,
75    #[serde(rename = "clickhouse.delete.column")]
76    pub delete_column: Option<String>,
77    /// Commit every n(>0) checkpoints, default is 10.
78    #[serde(default = "default_commit_checkpoint_interval")]
79    #[serde_as(as = "DisplayFromStr")]
80    #[with_option(allow_alter_on_fly)]
81    pub commit_checkpoint_interval: u64,
82}
83
84impl EnforceSecret for ClickHouseCommon {
85    const ENFORCE_SECRET_PROPERTIES: Set<&'static str> = phf_set! {
86        "clickhouse.password", "clickhouse.user"
87    };
88}
89
90#[derive(Debug)]
91enum ClickHouseEngine {
92    MergeTree,
93    ReplacingMergeTree(Option<String>),
94    SummingMergeTree,
95    AggregatingMergeTree,
96    CollapsingMergeTree(String),
97    VersionedCollapsingMergeTree(String),
98    GraphiteMergeTree,
99    ReplicatedMergeTree,
100    ReplicatedReplacingMergeTree(Option<String>),
101    ReplicatedSummingMergeTree,
102    ReplicatedAggregatingMergeTree,
103    ReplicatedCollapsingMergeTree(String),
104    ReplicatedVersionedCollapsingMergeTree(String),
105    ReplicatedGraphiteMergeTree,
106    SharedMergeTree,
107    SharedReplacingMergeTree(Option<String>),
108    SharedSummingMergeTree,
109    SharedAggregatingMergeTree,
110    SharedCollapsingMergeTree(String),
111    SharedVersionedCollapsingMergeTree(String),
112    SharedGraphiteMergeTree,
113    Null,
114}
115impl ClickHouseEngine {
116    pub fn is_collapsing_engine(&self) -> bool {
117        matches!(
118            self,
119            ClickHouseEngine::CollapsingMergeTree(_)
120                | ClickHouseEngine::VersionedCollapsingMergeTree(_)
121                | ClickHouseEngine::ReplicatedCollapsingMergeTree(_)
122                | ClickHouseEngine::ReplicatedVersionedCollapsingMergeTree(_)
123                | ClickHouseEngine::SharedCollapsingMergeTree(_)
124                | ClickHouseEngine::SharedVersionedCollapsingMergeTree(_)
125        )
126    }
127
128    pub fn is_delete_replacing_engine(&self) -> bool {
129        match self {
130            ClickHouseEngine::ReplacingMergeTree(delete_col) => delete_col.is_some(),
131            ClickHouseEngine::ReplicatedReplacingMergeTree(delete_col) => delete_col.is_some(),
132            ClickHouseEngine::SharedReplacingMergeTree(delete_col) => delete_col.is_some(),
133            _ => false,
134        }
135    }
136
137    pub fn get_delete_col(&self) -> Option<String> {
138        match self {
139            ClickHouseEngine::ReplacingMergeTree(Some(delete_col)) => Some(delete_col.clone()),
140            ClickHouseEngine::ReplicatedReplacingMergeTree(Some(delete_col)) => {
141                Some(delete_col.clone())
142            }
143            ClickHouseEngine::SharedReplacingMergeTree(Some(delete_col)) => {
144                Some(delete_col.clone())
145            }
146            _ => None,
147        }
148    }
149
150    pub fn get_sign_name(&self) -> Option<String> {
151        match self {
152            ClickHouseEngine::CollapsingMergeTree(sign_name) => Some(sign_name.clone()),
153            ClickHouseEngine::VersionedCollapsingMergeTree(sign_name) => Some(sign_name.clone()),
154            ClickHouseEngine::ReplicatedCollapsingMergeTree(sign_name) => Some(sign_name.clone()),
155            ClickHouseEngine::ReplicatedVersionedCollapsingMergeTree(sign_name) => {
156                Some(sign_name.clone())
157            }
158            ClickHouseEngine::SharedCollapsingMergeTree(sign_name) => Some(sign_name.clone()),
159            ClickHouseEngine::SharedVersionedCollapsingMergeTree(sign_name) => {
160                Some(sign_name.clone())
161            }
162            _ => None,
163        }
164    }
165
166    pub fn is_shared_tree(&self) -> bool {
167        matches!(
168            self,
169            ClickHouseEngine::SharedMergeTree
170                | ClickHouseEngine::SharedReplacingMergeTree(_)
171                | ClickHouseEngine::SharedSummingMergeTree
172                | ClickHouseEngine::SharedAggregatingMergeTree
173                | ClickHouseEngine::SharedCollapsingMergeTree(_)
174                | ClickHouseEngine::SharedVersionedCollapsingMergeTree(_)
175                | ClickHouseEngine::SharedGraphiteMergeTree
176        )
177    }
178
179    pub fn from_query_engine(
180        engine_name: &ClickhouseQueryEngine,
181        config: &ClickHouseConfig,
182    ) -> Result<Self> {
183        match engine_name.engine.as_str() {
184            "MergeTree" => Ok(ClickHouseEngine::MergeTree),
185            "Null" => Ok(ClickHouseEngine::Null),
186            "ReplacingMergeTree" => {
187                let delete_column = config.common.delete_column.clone();
188                Ok(ClickHouseEngine::ReplacingMergeTree(delete_column))
189            }
190            "SummingMergeTree" => Ok(ClickHouseEngine::SummingMergeTree),
191            "AggregatingMergeTree" => Ok(ClickHouseEngine::AggregatingMergeTree),
192            // VersionedCollapsingMergeTree(sign_name,"a")
193            "VersionedCollapsingMergeTree" => {
194                let sign_name = engine_name
195                    .create_table_query
196                    .split("VersionedCollapsingMergeTree(")
197                    .last()
198                    .ok_or_else(|| SinkError::ClickHouse("must have last".to_owned()))?
199                    .split(',')
200                    .next()
201                    .ok_or_else(|| SinkError::ClickHouse("must have next".to_owned()))?
202                    .trim()
203                    .to_owned();
204                Ok(ClickHouseEngine::VersionedCollapsingMergeTree(sign_name))
205            }
206            // CollapsingMergeTree(sign_name)
207            "CollapsingMergeTree" => {
208                let sign_name = engine_name
209                    .create_table_query
210                    .split("CollapsingMergeTree(")
211                    .last()
212                    .ok_or_else(|| SinkError::ClickHouse("must have last".to_owned()))?
213                    .split(')')
214                    .next()
215                    .ok_or_else(|| SinkError::ClickHouse("must have next".to_owned()))?
216                    .trim()
217                    .to_owned();
218                Ok(ClickHouseEngine::CollapsingMergeTree(sign_name))
219            }
220            "GraphiteMergeTree" => Ok(ClickHouseEngine::GraphiteMergeTree),
221            "ReplicatedMergeTree" => Ok(ClickHouseEngine::ReplicatedMergeTree),
222            "ReplicatedReplacingMergeTree" => {
223                let delete_column = config.common.delete_column.clone();
224                Ok(ClickHouseEngine::ReplicatedReplacingMergeTree(
225                    delete_column,
226                ))
227            }
228            "ReplicatedSummingMergeTree" => Ok(ClickHouseEngine::ReplicatedSummingMergeTree),
229            "ReplicatedAggregatingMergeTree" => {
230                Ok(ClickHouseEngine::ReplicatedAggregatingMergeTree)
231            }
232            // ReplicatedVersionedCollapsingMergeTree("a","b",sign_name,"c")
233            "ReplicatedVersionedCollapsingMergeTree" => {
234                let sign_name = engine_name
235                    .create_table_query
236                    .split("ReplicatedVersionedCollapsingMergeTree(")
237                    .last()
238                    .ok_or_else(|| SinkError::ClickHouse("must have last".to_owned()))?
239                    .split(',')
240                    .rev()
241                    .nth(1)
242                    .ok_or_else(|| SinkError::ClickHouse("must have index 1".to_owned()))?
243                    .trim()
244                    .to_owned();
245                Ok(ClickHouseEngine::ReplicatedVersionedCollapsingMergeTree(
246                    sign_name,
247                ))
248            }
249            // ReplicatedCollapsingMergeTree("a","b",sign_name)
250            "ReplicatedCollapsingMergeTree" => {
251                let sign_name = engine_name
252                    .create_table_query
253                    .split("ReplicatedCollapsingMergeTree(")
254                    .last()
255                    .ok_or_else(|| SinkError::ClickHouse("must have last".to_owned()))?
256                    .split(')')
257                    .next()
258                    .ok_or_else(|| SinkError::ClickHouse("must have next".to_owned()))?
259                    .split(',')
260                    .next_back()
261                    .ok_or_else(|| SinkError::ClickHouse("must have last".to_owned()))?
262                    .trim()
263                    .to_owned();
264                Ok(ClickHouseEngine::ReplicatedCollapsingMergeTree(sign_name))
265            }
266            "ReplicatedGraphiteMergeTree" => Ok(ClickHouseEngine::ReplicatedGraphiteMergeTree),
267            "SharedMergeTree" => Ok(ClickHouseEngine::SharedMergeTree),
268            "SharedReplacingMergeTree" => {
269                let delete_column = config.common.delete_column.clone();
270                Ok(ClickHouseEngine::SharedReplacingMergeTree(delete_column))
271            }
272            "SharedSummingMergeTree" => Ok(ClickHouseEngine::SharedSummingMergeTree),
273            "SharedAggregatingMergeTree" => Ok(ClickHouseEngine::SharedAggregatingMergeTree),
274            // SharedVersionedCollapsingMergeTree("a","b",sign_name,"c")
275            "SharedVersionedCollapsingMergeTree" => {
276                let sign_name = engine_name
277                    .create_table_query
278                    .split("SharedVersionedCollapsingMergeTree(")
279                    .last()
280                    .ok_or_else(|| SinkError::ClickHouse("must have last".to_owned()))?
281                    .split(',')
282                    .rev()
283                    .nth(1)
284                    .ok_or_else(|| SinkError::ClickHouse("must have index 1".to_owned()))?
285                    .trim()
286                    .to_owned();
287                Ok(ClickHouseEngine::SharedVersionedCollapsingMergeTree(
288                    sign_name,
289                ))
290            }
291            // SharedCollapsingMergeTree("a","b",sign_name)
292            "SharedCollapsingMergeTree" => {
293                let sign_name = engine_name
294                    .create_table_query
295                    .split("SharedCollapsingMergeTree(")
296                    .last()
297                    .ok_or_else(|| SinkError::ClickHouse("must have last".to_owned()))?
298                    .split(')')
299                    .next()
300                    .ok_or_else(|| SinkError::ClickHouse("must have next".to_owned()))?
301                    .split(',')
302                    .next_back()
303                    .ok_or_else(|| SinkError::ClickHouse("must have last".to_owned()))?
304                    .trim()
305                    .to_owned();
306                Ok(ClickHouseEngine::SharedCollapsingMergeTree(sign_name))
307            }
308            "SharedGraphiteMergeTree" => Ok(ClickHouseEngine::SharedGraphiteMergeTree),
309            _ => Err(SinkError::ClickHouse(format!(
310                "Cannot find clickhouse engine {:?}",
311                engine_name.engine
312            ))),
313        }
314    }
315}
316
317impl ClickHouseCommon {
318    pub(crate) fn build_client(&self) -> ConnectorResult<ClickHouseClient> {
319        let client = ClickHouseClient::default() // hyper(0.14) client inside
320            .with_url(&self.url)
321            .with_user(&self.user)
322            .with_password(&self.password)
323            .with_database(&self.database)
324            .with_option(ALLOW_EXPERIMENTAL_JSON_TYPE, "1")
325            .with_option(INPUT_FORMAT_BINARY_READ_JSON_AS_STRING, "1")
326            .with_option(OUTPUT_FORMAT_BINARY_WRITE_JSON_AS_STRING, "1");
327        Ok(client)
328    }
329}
330
331#[serde_as]
332#[derive(Clone, Debug, Deserialize, WithOptions)]
333pub struct ClickHouseConfig {
334    #[serde(flatten)]
335    pub common: ClickHouseCommon,
336
337    pub r#type: String, // accept "append-only" or "upsert"
338
339    #[serde(flatten)]
340    pub unknown_fields: std::collections::HashMap<String, String>,
341}
342
343crate::impl_sink_unknown_fields!(ClickHouseConfig);
344
345impl EnforceSecret for ClickHouseConfig {
346    fn enforce_one(prop: &str) -> crate::error::ConnectorResult<()> {
347        ClickHouseCommon::enforce_one(prop)
348    }
349
350    fn enforce_secret<'a>(
351        prop_iter: impl Iterator<Item = &'a str>,
352    ) -> crate::error::ConnectorResult<()> {
353        for prop in prop_iter {
354            ClickHouseCommon::enforce_one(prop)?;
355        }
356        Ok(())
357    }
358}
359
360#[derive(Clone, Debug)]
361pub struct ClickHouseSink {
362    pub config: ClickHouseConfig,
363    schema: Schema,
364    pk_indices: Vec<usize>,
365    is_append_only: bool,
366}
367
368impl EnforceSecret for ClickHouseSink {
369    fn enforce_secret<'a>(
370        prop_iter: impl Iterator<Item = &'a str>,
371    ) -> crate::error::ConnectorResult<()> {
372        for prop in prop_iter {
373            ClickHouseConfig::enforce_one(prop)?;
374        }
375        Ok(())
376    }
377}
378
379impl ClickHouseConfig {
380    pub fn from_btreemap(properties: BTreeMap<String, String>) -> Result<Self> {
381        let config =
382            serde_json::from_value::<ClickHouseConfig>(serde_json::to_value(properties).unwrap())
383                .map_err(|e| SinkError::Config(anyhow!(e)))?;
384        if config.r#type != SINK_TYPE_APPEND_ONLY && config.r#type != SINK_TYPE_UPSERT {
385            return Err(SinkError::Config(anyhow!(
386                "`{}` must be {}, or {}",
387                SINK_TYPE_OPTION,
388                SINK_TYPE_APPEND_ONLY,
389                SINK_TYPE_UPSERT
390            )));
391        }
392        Ok(config)
393    }
394}
395
396impl TryFrom<SinkParam> for ClickHouseSink {
397    type Error = SinkError;
398
399    fn try_from(param: SinkParam) -> std::result::Result<Self, Self::Error> {
400        let schema = param.schema();
401        let pk_indices = param.downstream_pk_or_empty();
402        let config = ClickHouseConfig::from_btreemap(param.properties)?;
403        Ok(Self {
404            config,
405            schema,
406            pk_indices,
407            is_append_only: param.sink_type.is_append_only(),
408        })
409    }
410}
411
412impl ClickHouseSink {
413    /// Check that the column names and types of risingwave and clickhouse are identical
414    fn check_column_name_and_type(&self, clickhouse_columns_desc: &[SystemColumn]) -> Result<()> {
415        let rw_fields_name = build_fields_name_type_from_schema(&self.schema)?;
416        let clickhouse_columns_desc: HashMap<String, SystemColumn> = clickhouse_columns_desc
417            .iter()
418            .map(|s| (s.name.clone(), s.clone()))
419            .collect();
420
421        if rw_fields_name.len().gt(&clickhouse_columns_desc.len()) {
422            return Err(SinkError::ClickHouse("The columns of the sink must be equal to or a superset of the target table's columns.".to_owned()));
423        }
424
425        for i in rw_fields_name {
426            let value = clickhouse_columns_desc.get(&i.0).ok_or_else(|| {
427                SinkError::ClickHouse(format!(
428                    "Column name don't find in clickhouse, risingwave is {:?} ",
429                    i.0
430                ))
431            })?;
432
433            Self::check_and_correct_column_type(&i.1, value)?;
434        }
435        Ok(())
436    }
437
438    /// Check that the column names and types of risingwave and clickhouse are identical
439    fn check_pk_match(&self, clickhouse_columns_desc: &[SystemColumn]) -> Result<()> {
440        let mut clickhouse_pks: HashSet<String> = clickhouse_columns_desc
441            .iter()
442            .filter(|s| s.is_in_primary_key == 1)
443            .map(|s| s.name.clone())
444            .collect();
445
446        for (_, field) in self
447            .schema
448            .fields()
449            .iter()
450            .enumerate()
451            .filter(|(index, _)| self.pk_indices.contains(index))
452        {
453            if !clickhouse_pks.remove(&field.name) {
454                return Err(SinkError::ClickHouse(
455                    "Clicklhouse and RisingWave pk is not match".to_owned(),
456                ));
457            }
458        }
459
460        if !clickhouse_pks.is_empty() {
461            return Err(SinkError::ClickHouse(
462                "Clicklhouse and RisingWave pk is not match".to_owned(),
463            ));
464        }
465        Ok(())
466    }
467
468    /// Check that the column types of risingwave and clickhouse are identical
469    fn check_and_correct_column_type(
470        fields_type: &DataType,
471        ck_column: &SystemColumn,
472    ) -> Result<()> {
473        // FIXME: the "contains" based implementation is wrong
474        let is_match = match fields_type {
475            risingwave_common::types::DataType::Boolean => Ok(ck_column.r#type.contains("Bool")),
476            risingwave_common::types::DataType::Int16 => Ok(ck_column.r#type.contains("UInt16")
477                | ck_column.r#type.contains("Int16")
478                // Allow Int16 to be pushed to Enum16, they share an encoding and value range
479                // No special care is taken to ensure values are valid.
480                | ck_column.r#type.contains("Enum16")),
481            risingwave_common::types::DataType::Int32 => {
482                Ok(ck_column.r#type.contains("UInt32") | ck_column.r#type.contains("Int32"))
483            }
484            risingwave_common::types::DataType::Int64 => {
485                Ok(ck_column.r#type.contains("UInt64") | ck_column.r#type.contains("Int64"))
486            }
487            risingwave_common::types::DataType::Float32 => Ok(ck_column.r#type.contains("Float32")),
488            risingwave_common::types::DataType::Float64 => Ok(ck_column.r#type.contains("Float64")),
489            risingwave_common::types::DataType::Decimal => {
490                match parse_decimal_accuracy(&ck_column.r#type)? {
491                    Some((precision, _)) if precision > MAX_ENCODABLE_DECIMAL_PRECISION => {
492                        Err(SinkError::ClickHouse(format!(
493                            "column {:?} has ClickHouse type `{}`: precision {} is backed by Decimal256, \
494                             but this sink can only encode decimals up to precision {}. Lower the column \
495                             precision (recreating the table if it is a key column), or leave the column \
496                             out of the sink.",
497                            ck_column.name,
498                            ck_column.r#type,
499                            precision,
500                            MAX_ENCODABLE_DECIMAL_PRECISION,
501                        )))
502                    }
503                    accuracy => Ok(accuracy.is_some()),
504                }
505            }
506            risingwave_common::types::DataType::Date => Ok(ck_column.r#type.contains("Date32")),
507            risingwave_common::types::DataType::Varchar => Ok(ck_column.r#type.contains("String")),
508            risingwave_common::types::DataType::Time => Err(SinkError::ClickHouse(
509                "TIME is not supported for ClickHouse sink. Please convert to VARCHAR or other supported types.".to_owned(),
510            )),
511            risingwave_common::types::DataType::Timestamp => Err(SinkError::ClickHouse(
512                "clickhouse does not have a type corresponding to naive timestamp".to_owned(),
513            )),
514            risingwave_common::types::DataType::Timestamptz => {
515                Ok(ck_column.r#type.contains("DateTime64"))
516            }
517            risingwave_common::types::DataType::Interval => Err(SinkError::ClickHouse(
518                "INTERVAL is not supported for ClickHouse sink. Please convert to VARCHAR or other supported types.".to_owned(),
519            )),
520            risingwave_common::types::DataType::Struct(_) => Err(SinkError::ClickHouse(
521                "struct needs to be converted into a list".to_owned(),
522            )),
523            risingwave_common::types::DataType::List(list) => {
524                Self::check_and_correct_column_type(list.elem(), ck_column)?;
525                Ok(ck_column.r#type.contains("Array"))
526            }
527            risingwave_common::types::DataType::Bytea => Err(SinkError::ClickHouse(
528                "BYTEA is not supported for ClickHouse sink. Please convert to VARCHAR or other supported types.".to_owned(),
529            )),
530            risingwave_common::types::DataType::Jsonb => Ok(ck_column.r#type.contains("JSON")),
531            risingwave_common::types::DataType::Variant => Err(SinkError::ClickHouse(
532                "VARIANT is not supported for ClickHouse sink.".to_owned(),
533            )),
534            risingwave_common::types::DataType::Serial => {
535                Ok(ck_column.r#type.contains("UInt64") | ck_column.r#type.contains("Int64"))
536            }
537            risingwave_common::types::DataType::Int256 => Err(SinkError::ClickHouse(
538                "INT256 is not supported for ClickHouse sink.".to_owned(),
539            )),
540            risingwave_common::types::DataType::Map(_) => Err(SinkError::ClickHouse(
541                "MAP is not supported for ClickHouse sink.".to_owned(),
542            )),
543            DataType::Vector(_) => Err(SinkError::ClickHouse(
544                "VECTOR is not supported for ClickHouse sink.".to_owned(),
545            )),
546        };
547        if !is_match? {
548            return Err(SinkError::ClickHouse(format!(
549                "Column type mismatch for column {:?}: RisingWave type is {:?}, ClickHouse type is {:?}",
550                ck_column.name, fields_type, ck_column.r#type
551            )));
552        }
553
554        Ok(())
555    }
556}
557
558impl Sink for ClickHouseSink {
559    type LogSinker = DecoupleCheckpointLogSinkerOf<ClickHouseSinkWriter>;
560
561    const SINK_NAME: &'static str = CLICKHOUSE_SINK;
562
563    crate::impl_validate_sink_unknown_fields!();
564
565    async fn validate(&self) -> Result<()> {
566        // For upsert clickhouse sink, the primary key must be defined.
567        if !self.is_append_only && self.pk_indices.is_empty() {
568            return Err(SinkError::Config(anyhow!(
569                "Primary key not defined for upsert clickhouse sink (please define in `primary_key` field)"
570            )));
571        }
572
573        // check reachability
574        let client = self.config.common.build_client()?;
575
576        let (clickhouse_column, clickhouse_engine) =
577            query_column_engine_from_ck(client, &self.config).await?;
578        if clickhouse_engine.is_shared_tree() {
579            risingwave_common::license::Feature::ClickHouseSharedEngine
580                .check_available()
581                .map_err(|e| anyhow::anyhow!(e))?;
582        }
583
584        if !self.is_append_only
585            && !clickhouse_engine.is_collapsing_engine()
586            && !clickhouse_engine.is_delete_replacing_engine()
587        {
588            return match clickhouse_engine {
589                ClickHouseEngine::ReplicatedReplacingMergeTree(None) | ClickHouseEngine::ReplacingMergeTree(None) | ClickHouseEngine::SharedReplacingMergeTree(None) =>  {
590                    Err(SinkError::ClickHouse("To enable upsert with a `ReplacingMergeTree`, you must set a `clickhouse.delete.column` to the UInt8 column in ClickHouse used to signify deletes. See https://clickhouse.com/docs/en/engines/table-engines/mergetree-family/replacingmergetree#is_deleted for more information".to_owned()))
591                }
592                _ => Err(SinkError::ClickHouse("If you want to use upsert, please use either `VersionedCollapsingMergeTree`, `CollapsingMergeTree` or the `ReplacingMergeTree` in ClickHouse".to_owned()))
593            };
594        }
595
596        self.check_column_name_and_type(&clickhouse_column)?;
597        if !self.is_append_only {
598            self.check_pk_match(&clickhouse_column)?;
599        }
600
601        if self.config.common.commit_checkpoint_interval == 0 {
602            return Err(SinkError::Config(anyhow!(
603                "`commit_checkpoint_interval` must be greater than 0"
604            )));
605        }
606        Ok(())
607    }
608
609    fn validate_alter_config(config: &BTreeMap<String, String>) -> Result<()> {
610        ClickHouseConfig::from_btreemap(config.clone())?;
611        Ok(())
612    }
613
614    async fn new_log_sinker(&self, writer_param: SinkWriterParam) -> Result<Self::LogSinker> {
615        let writer = ClickHouseSinkWriter::new(
616            self.config.clone(),
617            self.schema.clone(),
618            self.pk_indices.clone(),
619            self.is_append_only,
620        )
621        .await?;
622        let commit_checkpoint_interval =
623    NonZeroU64::new(self.config.common.commit_checkpoint_interval).expect(
624        "commit_checkpoint_interval should be greater than 0, and it should be checked in config validation",
625    );
626
627        Ok(DecoupleCheckpointLogSinkerOf::new(
628            writer,
629            SinkWriterMetrics::new(&writer_param),
630            commit_checkpoint_interval,
631        ))
632    }
633}
634pub struct ClickHouseSinkWriter {
635    pub config: ClickHouseConfig,
636    schema: Schema,
637    #[expect(dead_code)]
638    pk_indices: Vec<usize>,
639    client: ClickHouseClient,
640    #[expect(dead_code)]
641    is_append_only: bool,
642    // Save some features of the clickhouse column type
643    column_correct_map: HashMap<String, ClickHouseSchemaFeature>,
644    rw_fields_name_after_calibration: Vec<String>,
645    clickhouse_engine: ClickHouseEngine,
646    inserter: Option<Insert<ClickHouseColumn>>,
647}
648#[derive(Debug)]
649struct ClickHouseSchemaFeature {
650    can_null: bool,
651    // Time accuracy in clickhouse for rw and ck conversions
652    accuracy_time: u8,
653
654    accuracy_decimal: (u8, u8),
655}
656
657impl ClickHouseSinkWriter {
658    pub async fn new(
659        config: ClickHouseConfig,
660        schema: Schema,
661        pk_indices: Vec<usize>,
662        is_append_only: bool,
663    ) -> Result<Self> {
664        let client = config.common.build_client()?;
665
666        let (clickhouse_column, clickhouse_engine) =
667            query_column_engine_from_ck(client.clone(), &config).await?;
668
669        let column_correct_map: Result<HashMap<String, ClickHouseSchemaFeature>> =
670            clickhouse_column
671                .iter()
672                .map(Self::build_column_correct_map)
673                .collect();
674        let sink_fields = build_fields_name_type_from_schema(&schema)?;
675
676        // The table may have been altered since `CREATE SINK` validated it, and an under-sized
677        // decimal encoding would desync every following column of the same RowBinary row. Only
678        // columns the writer encodes as decimals are held to the decimal rule.
679        let ck_columns: HashMap<&str, &SystemColumn> = clickhouse_column
680            .iter()
681            .map(|c| (c.name.as_str(), c))
682            .collect();
683        for (name, data_type) in sink_fields.iter().filter(|(_, t)| contains_decimal(t)) {
684            if let Some(ck_column) = ck_columns.get(name.as_str()) {
685                ClickHouseSink::check_and_correct_column_type(data_type, ck_column)?;
686            }
687        }
688
689        let mut rw_fields_name_after_calibration =
690            sink_fields.iter().map(|(a, _)| a.clone()).collect_vec();
691
692        if let Some(sign) = clickhouse_engine.get_sign_name() {
693            rw_fields_name_after_calibration.push(sign);
694        }
695        if let Some(delete_col) = clickhouse_engine.get_delete_col() {
696            rw_fields_name_after_calibration.push(delete_col);
697        }
698        Ok(Self {
699            config,
700            schema,
701            pk_indices,
702            client,
703            is_append_only,
704            column_correct_map: column_correct_map?,
705            rw_fields_name_after_calibration,
706            clickhouse_engine,
707            inserter: None,
708        })
709    }
710
711    /// Check if clickhouse's column is 'Nullable', valid bits of `DateTime64`. And save it in
712    /// `column_correct_map`
713    fn build_column_correct_map(
714        ck_column: &SystemColumn,
715    ) -> Result<(String, ClickHouseSchemaFeature)> {
716        let can_null = ck_column.r#type.contains("Nullable");
717        // `DateTime64` without precision is already displayed as `DateTime(3)` in `system.columns`.
718        let accuracy_time = if ck_column.r#type.contains("DateTime64(") {
719            ck_column
720                .r#type
721                .split("DateTime64(")
722                .last()
723                .ok_or_else(|| SinkError::ClickHouse("must have last".to_owned()))?
724                .split(')')
725                .next()
726                .ok_or_else(|| SinkError::ClickHouse("must have next".to_owned()))?
727                .split(',')
728                .next()
729                .ok_or_else(|| SinkError::ClickHouse("must have next".to_owned()))?
730                .parse::<u8>()
731                .map_err(|e| SinkError::ClickHouse(e.to_report_string()))?
732        } else {
733            0_u8
734        };
735        let accuracy_decimal = parse_decimal_accuracy(&ck_column.r#type)?.unwrap_or((0, 0));
736        Ok((
737            ck_column.name.clone(),
738            ClickHouseSchemaFeature {
739                can_null,
740                accuracy_time,
741                accuracy_decimal,
742            },
743        ))
744    }
745
746    async fn write(&mut self, chunk: StreamChunk) -> Result<()> {
747        if self.inserter.is_none() {
748            self.inserter = Some(self.client.insert_with_fields_name(
749                &self.config.common.table,
750                self.rw_fields_name_after_calibration.clone(),
751            )?);
752        }
753        for (op, row) in chunk.rows() {
754            let mut clickhouse_filed_vec = vec![];
755            for (data, field) in row.iter().zip_eq_debug(self.schema.fields()) {
756                clickhouse_filed_vec.extend(ClickHouseFieldWithNull::from_scalar_ref(
757                    data,
758                    &self.column_correct_map,
759                    field.name(),
760                    &field.data_type(),
761                )?);
762            }
763            match op {
764                Op::Insert | Op::UpdateInsert => {
765                    if self.clickhouse_engine.is_collapsing_engine() {
766                        clickhouse_filed_vec.push(ClickHouseFieldWithNull::WithoutSome(
767                            ClickHouseField::Int8(1),
768                        ));
769                    }
770                    if self.clickhouse_engine.is_delete_replacing_engine() {
771                        clickhouse_filed_vec.push(ClickHouseFieldWithNull::WithoutSome(
772                            ClickHouseField::Int8(0),
773                        ))
774                    }
775                }
776                Op::Delete | Op::UpdateDelete => {
777                    if !self.clickhouse_engine.is_collapsing_engine()
778                        && !self.clickhouse_engine.is_delete_replacing_engine()
779                    {
780                        return Err(SinkError::ClickHouse(
781                            "Clickhouse engine don't support upsert".to_owned(),
782                        ));
783                    }
784                    if self.clickhouse_engine.is_collapsing_engine() {
785                        clickhouse_filed_vec.push(ClickHouseFieldWithNull::WithoutSome(
786                            ClickHouseField::Int8(-1),
787                        ));
788                    }
789                    if self.clickhouse_engine.is_delete_replacing_engine() {
790                        clickhouse_filed_vec.push(ClickHouseFieldWithNull::WithoutSome(
791                            ClickHouseField::Int8(1),
792                        ))
793                    }
794                }
795            }
796            let clickhouse_column = ClickHouseColumn {
797                row: clickhouse_filed_vec,
798            };
799            self.inserter
800                .as_mut()
801                .unwrap()
802                .write(&clickhouse_column)
803                .await?;
804        }
805        Ok(())
806    }
807}
808
809#[async_trait]
810impl SinkWriter for ClickHouseSinkWriter {
811    async fn write_batch(&mut self, chunk: StreamChunk) -> Result<()> {
812        self.write(chunk).await
813    }
814
815    async fn begin_epoch(&mut self, _epoch: u64) -> Result<()> {
816        Ok(())
817    }
818
819    async fn abort(&mut self) -> Result<()> {
820        Ok(())
821    }
822
823    async fn barrier(&mut self, is_checkpoint: bool) -> Result<()> {
824        if is_checkpoint && let Some(inserter) = self.inserter.take() {
825            inserter.end().await?;
826        }
827        Ok(())
828    }
829}
830
831#[derive(ClickHouseRow, Deserialize, Clone)]
832struct SystemColumn {
833    name: String,
834    r#type: String,
835    is_in_primary_key: u8,
836}
837
838/// Parse `(precision, scale)` out of a decimal type, or `None` if it is not a decimal.
839/// `system.columns` renders decimals canonically, so `Decimal128(10)` arrives as `Decimal(38, 10)`.
840fn parse_decimal_accuracy(ck_type: &str) -> Result<Option<(u8, u8)>> {
841    let Some((_, rest)) = ck_type.split_once("Decimal(") else {
842        return Ok(None);
843    };
844    let (args, _) = rest
845        .split_once(')')
846        .ok_or_else(|| SinkError::ClickHouse(format!("unterminated decimal type {ck_type:?}")))?;
847    let (precision, scale) = args
848        .split_once(',')
849        .ok_or_else(|| SinkError::ClickHouse(format!("decimal type {ck_type:?} has no scale")))?;
850    let parse = |field: &str| {
851        field.trim().parse::<u8>().map_err(|e| {
852            SinkError::ClickHouse(format!(
853                "cannot parse decimal type {:?}: {}",
854                ck_type,
855                e.to_report_string()
856            ))
857        })
858    };
859    Ok(Some((parse(precision)?, parse(scale)?)))
860}
861
862fn contains_decimal(data_type: &DataType) -> bool {
863    match data_type {
864        DataType::Decimal => true,
865        DataType::List(list) => contains_decimal(list.elem()),
866        _ => false,
867    }
868}
869
870#[derive(ClickHouseRow, Deserialize)]
871struct ClickhouseQueryEngine {
872    #[expect(dead_code)]
873    name: String,
874    engine: String,
875    create_table_query: String,
876}
877
878async fn query_column_engine_from_ck(
879    client: ClickHouseClient,
880    config: &ClickHouseConfig,
881) -> Result<(Vec<SystemColumn>, ClickHouseEngine)> {
882    let query_engine = QUERY_ENGINE;
883    let query_column = QUERY_COLUMN;
884
885    let clickhouse_engine = client
886        .query(query_engine)
887        .bind(config.common.database.clone())
888        .bind(config.common.table.clone())
889        .fetch_all::<ClickhouseQueryEngine>()
890        .await?;
891    let mut clickhouse_column = client
892        .query(query_column)
893        .bind(config.common.database.clone())
894        .bind(config.common.table.clone())
895        .bind("position")
896        .fetch_all::<SystemColumn>()
897        .await?;
898    if clickhouse_engine.is_empty() || clickhouse_column.is_empty() {
899        return Err(SinkError::ClickHouse(format!(
900            "table {:?}.{:?} is not find in clickhouse",
901            config.common.database, config.common.table
902        )));
903    }
904
905    let clickhouse_engine =
906        ClickHouseEngine::from_query_engine(clickhouse_engine.first().unwrap(), config)?;
907
908    if let Some(sign) = &clickhouse_engine.get_sign_name() {
909        clickhouse_column.retain(|a| sign.ne(&a.name))
910    }
911
912    if let Some(delete_col) = &clickhouse_engine.get_delete_col() {
913        clickhouse_column.retain(|a| delete_col.ne(&a.name))
914    }
915
916    Ok((clickhouse_column, clickhouse_engine))
917}
918
919/// Serialize this structure to simulate the `struct` call clickhouse interface
920#[derive(ClickHouseRow, Debug)]
921struct ClickHouseColumn {
922    row: Vec<ClickHouseFieldWithNull>,
923}
924
925/// Basic data types for use with the clickhouse interface
926#[derive(Debug)]
927enum ClickHouseField {
928    Int16(i16),
929    Int32(i32),
930    Int64(i64),
931    Serial(Serial),
932    Float32(f32),
933    Float64(f64),
934    String(String),
935    Bool(bool),
936    List(Vec<ClickHouseFieldWithNull>),
937    Int8(i8),
938    Decimal(ClickHouseDecimal),
939}
940#[derive(Debug)]
941enum ClickHouseDecimal {
942    Decimal32(i32),
943    Decimal64(i64),
944    Decimal128(i128),
945}
946impl Serialize for ClickHouseDecimal {
947    fn serialize<S>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error>
948    where
949        S: serde::Serializer,
950    {
951        match self {
952            ClickHouseDecimal::Decimal32(v) => serializer.serialize_i32(*v),
953            ClickHouseDecimal::Decimal64(v) => serializer.serialize_i64(*v),
954            ClickHouseDecimal::Decimal128(v) => serializer.serialize_i128(*v),
955        }
956    }
957}
958
959/// Enum that support clickhouse nullable
960#[derive(Debug)]
961enum ClickHouseFieldWithNull {
962    WithSome(ClickHouseField),
963    WithoutSome(ClickHouseField),
964    None,
965}
966
967impl ClickHouseFieldWithNull {
968    pub fn from_scalar_ref(
969        data: Option<ScalarRefImpl<'_>>,
970        clickhouse_schema_feature_map: &HashMap<String, ClickHouseSchemaFeature>,
971        rw_name: &str,
972        rw_type: &DataType,
973    ) -> Result<Vec<ClickHouseFieldWithNull>> {
974        let clickhouse_schema_feature = if !matches!(rw_type, DataType::Struct(_)) {
975            Some(clickhouse_schema_feature_map.get(rw_name).ok_or_else(|| {
976                SinkError::ClickHouse(format!(
977                    "No column found from clickhouse table schema, name is {rw_name}"
978                ))
979            })?)
980        } else {
981            None
982        };
983
984        let Some(data) = data else {
985            let Some(clickhouse_schema_feature) = clickhouse_schema_feature else {
986                return Err(SinkError::ClickHouse(
987                    "clickhouse's nested can not insert null".to_owned(),
988                ));
989            };
990            if !clickhouse_schema_feature.can_null {
991                return Err(SinkError::ClickHouse(
992                    "Cannot insert null value into non-nullable ClickHouse column".to_owned(),
993                ));
994            } else {
995                return Ok(vec![ClickHouseFieldWithNull::None]);
996            }
997        };
998        let data = match data {
999            ScalarRefImpl::Int16(v) => ClickHouseField::Int16(v),
1000            ScalarRefImpl::Int32(v) => ClickHouseField::Int32(v),
1001            ScalarRefImpl::Int64(v) => ClickHouseField::Int64(v),
1002            ScalarRefImpl::Int256(_) => {
1003                return Err(SinkError::ClickHouse(
1004                    "INT256 is not supported for ClickHouse sink.".to_owned(),
1005                ));
1006            }
1007            ScalarRefImpl::Serial(v) => ClickHouseField::Serial(v),
1008            ScalarRefImpl::Float32(v) => ClickHouseField::Float32(v.into_inner()),
1009            ScalarRefImpl::Float64(v) => ClickHouseField::Float64(v.into_inner()),
1010            ScalarRefImpl::Utf8(v) => ClickHouseField::String(v.to_owned()),
1011            ScalarRefImpl::Bool(v) => ClickHouseField::Bool(v),
1012            ScalarRefImpl::Decimal(d) => {
1013                let d = if let Decimal::Normalized(d) = d {
1014                    let scale = clickhouse_schema_feature.unwrap().accuracy_decimal.1 as i32
1015                        - d.scale() as i32;
1016                    if scale < 0 {
1017                        d.mantissa() / 10_i128.pow(scale.unsigned_abs())
1018                    } else {
1019                        d.mantissa() * 10_i128.pow(scale as u32)
1020                    }
1021                } else if clickhouse_schema_feature.unwrap().can_null {
1022                    warn!("Inf, -Inf, Nan in RW decimal is converted into clickhouse null!");
1023                    return Ok(vec![ClickHouseFieldWithNull::None]);
1024                } else {
1025                    warn!("Inf, -Inf, Nan in RW decimal is converted into clickhouse 0!");
1026                    0_i128
1027                };
1028                if clickhouse_schema_feature.unwrap().accuracy_decimal.0 <= 9 {
1029                    ClickHouseField::Decimal(ClickHouseDecimal::Decimal32(d as i32))
1030                } else if clickhouse_schema_feature.unwrap().accuracy_decimal.0 <= 18 {
1031                    ClickHouseField::Decimal(ClickHouseDecimal::Decimal64(d as i64))
1032                } else {
1033                    ClickHouseField::Decimal(ClickHouseDecimal::Decimal128(d))
1034                }
1035            }
1036            ScalarRefImpl::Interval(_) => {
1037                return Err(SinkError::ClickHouse(
1038                    "INTERVAL is not supported for ClickHouse sink. Please convert to VARCHAR or other supported types.".to_owned(),
1039                ));
1040            }
1041            ScalarRefImpl::Date(v) => {
1042                let days = v.get_nums_days_unix_epoch();
1043                ClickHouseField::Int32(days)
1044            }
1045            ScalarRefImpl::Time(_) => {
1046                return Err(SinkError::ClickHouse(
1047                    "TIME is not supported for ClickHouse sink. Please convert to VARCHAR or other supported types.".to_owned(),
1048                ));
1049            }
1050            ScalarRefImpl::Timestamp(_) => {
1051                return Err(SinkError::ClickHouse(
1052                    "clickhouse does not have a type corresponding to naive timestamp".to_owned(),
1053                ));
1054            }
1055            ScalarRefImpl::Timestamptz(v) => {
1056                let micros = v.timestamp_micros();
1057                let ticks = match clickhouse_schema_feature.unwrap().accuracy_time <= 6 {
1058                    true => {
1059                        micros
1060                            / 10_i64
1061                                .pow((6 - clickhouse_schema_feature.unwrap().accuracy_time).into())
1062                    }
1063                    false => micros
1064                        .checked_mul(
1065                            10_i64
1066                                .pow((clickhouse_schema_feature.unwrap().accuracy_time - 6).into()),
1067                        )
1068                        .ok_or_else(|| SinkError::ClickHouse("DateTime64 overflow".to_owned()))?,
1069                };
1070                ClickHouseField::Int64(ticks)
1071            }
1072            ScalarRefImpl::Jsonb(v) => {
1073                let json_str = v.to_string();
1074                ClickHouseField::String(json_str)
1075            }
1076            ScalarRefImpl::Variant(_) => {
1077                return Err(SinkError::ClickHouse(
1078                    "VARIANT is not supported for ClickHouse sink.".to_owned(),
1079                ));
1080            }
1081            ScalarRefImpl::Struct(v) => {
1082                let mut struct_vec = vec![];
1083                for (field, (struct_field_name, struct_field_type)) in
1084                    v.iter_fields_ref().zip_eq_debug(rw_type.as_struct().iter())
1085                {
1086                    let name = format!("{}.{}", rw_name, struct_field_name);
1087                    let a = Self::from_scalar_ref(
1088                        field,
1089                        clickhouse_schema_feature_map,
1090                        &name,
1091                        struct_field_type,
1092                    )?;
1093                    struct_vec.push(ClickHouseFieldWithNull::WithoutSome(ClickHouseField::List(
1094                        a,
1095                    )));
1096                }
1097                return Ok(struct_vec);
1098            }
1099            ScalarRefImpl::List(v) => {
1100                let mut vec = vec![];
1101                for i in v.iter() {
1102                    vec.extend(Self::from_scalar_ref(
1103                        i,
1104                        clickhouse_schema_feature_map,
1105                        rw_name,
1106                        rw_type,
1107                    )?)
1108                }
1109                return Ok(vec![ClickHouseFieldWithNull::WithoutSome(
1110                    ClickHouseField::List(vec),
1111                )]);
1112            }
1113            ScalarRefImpl::Bytea(_) => {
1114                return Err(SinkError::ClickHouse(
1115                    "BYTEA is not supported for ClickHouse sink. Please convert to VARCHAR or other supported types.".to_owned(),
1116                ));
1117            }
1118            ScalarRefImpl::Map(_) => {
1119                return Err(SinkError::ClickHouse(
1120                    "MAP is not supported for ClickHouse sink.".to_owned(),
1121                ));
1122            }
1123            ScalarRefImpl::Vector(_) => {
1124                return Err(SinkError::ClickHouse(
1125                    "VECTOR is not supported for ClickHouse sink.".to_owned(),
1126                ));
1127            }
1128        };
1129        let data = if let Some(clickhouse_schema_feature) = clickhouse_schema_feature
1130            && clickhouse_schema_feature.can_null
1131        {
1132            vec![ClickHouseFieldWithNull::WithSome(data)]
1133        } else {
1134            vec![ClickHouseFieldWithNull::WithoutSome(data)]
1135        };
1136        Ok(data)
1137    }
1138}
1139
1140impl Serialize for ClickHouseField {
1141    fn serialize<S>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error>
1142    where
1143        S: serde::Serializer,
1144    {
1145        match self {
1146            ClickHouseField::Int16(v) => serializer.serialize_i16(*v),
1147            ClickHouseField::Int32(v) => serializer.serialize_i32(*v),
1148            ClickHouseField::Int64(v) => serializer.serialize_i64(*v),
1149            ClickHouseField::Serial(v) => v.serialize(serializer),
1150            ClickHouseField::Float32(v) => serializer.serialize_f32(*v),
1151            ClickHouseField::Float64(v) => serializer.serialize_f64(*v),
1152            ClickHouseField::String(v) => serializer.serialize_str(v),
1153            ClickHouseField::Bool(v) => serializer.serialize_bool(*v),
1154            ClickHouseField::List(v) => {
1155                let mut s = serializer.serialize_seq(Some(v.len()))?;
1156                for i in v {
1157                    s.serialize_element(i)?;
1158                }
1159                s.end()
1160            }
1161            ClickHouseField::Decimal(v) => v.serialize(serializer),
1162            ClickHouseField::Int8(v) => serializer.serialize_i8(*v),
1163        }
1164    }
1165}
1166impl Serialize for ClickHouseFieldWithNull {
1167    fn serialize<S>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error>
1168    where
1169        S: serde::Serializer,
1170    {
1171        match self {
1172            ClickHouseFieldWithNull::WithSome(v) => serializer.serialize_some(v),
1173            ClickHouseFieldWithNull::WithoutSome(v) => v.serialize(serializer),
1174            ClickHouseFieldWithNull::None => serializer.serialize_none(),
1175        }
1176    }
1177}
1178impl Serialize for ClickHouseColumn {
1179    fn serialize<S>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error>
1180    where
1181        S: serde::Serializer,
1182    {
1183        let mut s = serializer.serialize_struct("useless", self.row.len())?;
1184        for data in &self.row {
1185            s.serialize_field("useless", &data)?
1186        }
1187        s.end()
1188    }
1189}
1190
1191/// 'Struct'(clickhouse type name is nested) will be converted into some arrays by clickhouse. So we
1192/// need to make some conversions
1193pub fn build_fields_name_type_from_schema(schema: &Schema) -> Result<Vec<(String, DataType)>> {
1194    let mut vec = vec![];
1195    for field in schema.fields() {
1196        if let DataType::Struct(st) = &field.data_type {
1197            for (name, data_type) in st.iter() {
1198                if matches!(data_type, DataType::Struct(_)) {
1199                    return Err(SinkError::ClickHouse(
1200                        "Only one level of nesting is supported for struct".to_owned(),
1201                    ));
1202                } else {
1203                    vec.push((
1204                        format!("{}.{}", field.name, name),
1205                        DataType::list(data_type.clone()),
1206                    ))
1207                }
1208            }
1209        } else {
1210            vec.push((field.name.clone(), field.data_type()));
1211        }
1212    }
1213    Ok(vec)
1214}
1215
1216impl From<::clickhouse::error::Error> for SinkError {
1217    fn from(value: ::clickhouse::error::Error) -> Self {
1218        SinkError::ClickHouse(value.to_report_string())
1219    }
1220}
1221
1222#[cfg(test)]
1223mod tests {
1224    use super::*;
1225
1226    fn system_column(ck_type: &str) -> SystemColumn {
1227        SystemColumn {
1228            name: "amount".to_owned(),
1229            r#type: ck_type.to_owned(),
1230            is_in_primary_key: 0,
1231        }
1232    }
1233
1234    #[test]
1235    fn test_parse_decimal_accuracy() {
1236        let parse = |ck_type: &str| parse_decimal_accuracy(ck_type).unwrap();
1237
1238        assert_eq!(parse("Decimal(38, 10)"), Some((38, 10)));
1239        assert_eq!(parse("Decimal(38,10)"), Some((38, 10)));
1240        assert_eq!(parse("Nullable(Decimal(76, 4))"), Some((76, 4)));
1241        assert_eq!(parse("DateTime64(3)"), None);
1242
1243        assert!(parse_decimal_accuracy("Decimal(38)").is_err());
1244    }
1245
1246    #[test]
1247    fn test_decimal_column_check() {
1248        let check = |rw_type: &DataType, ck_type: &str| {
1249            ClickHouseSink::check_and_correct_column_type(rw_type, &system_column(ck_type))
1250        };
1251        let decimal_list = DataType::list(DataType::Decimal);
1252
1253        check(&DataType::Decimal, "Nullable(Decimal(38, 10))").unwrap();
1254        check(&decimal_list, "Array(Decimal(9, 2))").unwrap();
1255        // Only columns encoded as decimals are held to the precision rule.
1256        check(&DataType::Jsonb, "JSON(a Decimal(76, 4), b Decimal(9, 2))").unwrap();
1257
1258        check(&DataType::Decimal, "String").unwrap_err();
1259        let err = check(&DataType::Decimal, "Decimal(76, 4)")
1260            .unwrap_err()
1261            .to_string();
1262        assert!(err.contains("amount"), "{err}");
1263        assert!(err.contains("Decimal(76, 4)"), "{err}");
1264        assert!(err.contains("precision 76"), "{err}");
1265    }
1266}