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
58const 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 #[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" => {
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" => {
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" => {
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" => {
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" => {
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" => {
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() .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, #[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 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 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 fn check_and_correct_column_type(
470 fields_type: &DataType,
471 ck_column: &SystemColumn,
472 ) -> Result<()> {
473 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 | 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 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 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 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 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 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 fn build_column_correct_map(
714 ck_column: &SystemColumn,
715 ) -> Result<(String, ClickHouseSchemaFeature)> {
716 let can_null = ck_column.r#type.contains("Nullable");
717 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
838fn 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#[derive(ClickHouseRow, Debug)]
921struct ClickHouseColumn {
922 row: Vec<ClickHouseFieldWithNull>,
923}
924
925#[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#[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
1191pub 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 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}