1use 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 #[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" => {
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" => {
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" => {
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" => {
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" => {
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" => {
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() .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, #[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 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 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 fn check_and_correct_column_type(
466 fields_type: &DataType,
467 ck_column: &SystemColumn,
468 ) -> Result<()> {
469 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 | 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 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 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 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 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 fn build_column_correct_map(
681 ck_column: &SystemColumn,
682 ) -> Result<(String, ClickHouseSchemaFeature)> {
683 let can_null = ck_column.r#type.contains("Nullable");
684 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#[derive(ClickHouseRow, Debug)]
887struct ClickHouseColumn {
888 row: Vec<ClickHouseFieldWithNull>,
889}
890
891#[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#[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
1157pub 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}