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