Skip to main content

risingwave_connector/sink/encoder/
json.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
15#[cfg(test)]
16use std::collections::{HashMap, HashSet};
17use std::sync::Arc;
18
19use anyhow::Context;
20use base64::Engine as _;
21use base64::engine::general_purpose;
22use chrono::{DateTime, Datelike, Timelike};
23use chrono_tz::Tz;
24use indexmap::IndexMap;
25use itertools::Itertools;
26use risingwave_common::array::{ArrayError, ArrayResult};
27use risingwave_common::catalog::{Field, Schema};
28use risingwave_common::row::Row;
29use risingwave_common::types::{DataType, DatumRef, JsonbVal, ScalarRefImpl, ToText};
30use risingwave_common::util::iter_util::ZipEqDebug;
31use serde_json::{Map, Value, json};
32use thiserror_ext::AsReport;
33
34use super::{
35    CustomJsonType, DateHandlingMode, DorisJsonConfig, JsonbHandlingMode, KafkaConnectParams,
36    KafkaConnectParamsRef, Result, RowEncoder, SerTo, TimeHandlingMode, TimestampHandlingMode,
37    TimestamptzHandlingMode,
38};
39use crate::sink::SinkError;
40
41pub struct JsonEncoderConfig {
42    time_handling_mode: TimeHandlingMode,
43    date_handling_mode: DateHandlingMode,
44    timestamp_handling_mode: TimestampHandlingMode,
45    timestamptz_handling_mode: TimestamptzHandlingMode,
46    custom_json_type: CustomJsonType,
47    jsonb_handling_mode: JsonbHandlingMode,
48}
49
50pub struct JsonEncoder {
51    schema: Schema,
52    col_indices: Option<Vec<usize>>,
53    kafka_connect: Option<KafkaConnectParamsRef>,
54    config: JsonEncoderConfig,
55}
56
57impl JsonEncoder {
58    pub fn new(
59        schema: Schema,
60        col_indices: Option<Vec<usize>>,
61        date_handling_mode: DateHandlingMode,
62        timestamp_handling_mode: TimestampHandlingMode,
63        timestamptz_handling_mode: TimestamptzHandlingMode,
64        time_handling_mode: TimeHandlingMode,
65        jsonb_handling_mode: JsonbHandlingMode,
66    ) -> Self {
67        let config = JsonEncoderConfig {
68            time_handling_mode,
69            date_handling_mode,
70            timestamp_handling_mode,
71            timestamptz_handling_mode,
72            custom_json_type: CustomJsonType::None,
73            jsonb_handling_mode,
74        };
75        Self {
76            schema,
77            col_indices,
78            kafka_connect: None,
79            config,
80        }
81    }
82
83    pub fn new_with_es(schema: Schema, col_indices: Option<Vec<usize>>) -> Self {
84        let config = JsonEncoderConfig {
85            time_handling_mode: TimeHandlingMode::String,
86            date_handling_mode: DateHandlingMode::String,
87            timestamp_handling_mode: TimestampHandlingMode::String,
88            timestamptz_handling_mode: TimestamptzHandlingMode::UtcWithoutSuffix,
89            custom_json_type: CustomJsonType::Es,
90            jsonb_handling_mode: JsonbHandlingMode::Dynamic,
91        };
92        Self {
93            schema,
94            col_indices,
95            kafka_connect: None,
96            config,
97        }
98    }
99
100    pub fn new_with_doris(
101        schema: Schema,
102        col_indices: Option<Vec<usize>>,
103        doris_config: DorisJsonConfig,
104    ) -> Self {
105        let config = JsonEncoderConfig {
106            time_handling_mode: TimeHandlingMode::Milli,
107            date_handling_mode: DateHandlingMode::String,
108            timestamp_handling_mode: TimestampHandlingMode::String,
109            timestamptz_handling_mode: TimestamptzHandlingMode::UtcWithoutSuffix,
110            custom_json_type: CustomJsonType::Doris(doris_config),
111            jsonb_handling_mode: JsonbHandlingMode::String,
112        };
113        Self {
114            schema,
115            col_indices,
116            kafka_connect: None,
117            config,
118        }
119    }
120
121    pub fn new_with_starrocks(
122        schema: Schema,
123        col_indices: Option<Vec<usize>>,
124        time_zone: Tz,
125    ) -> Self {
126        let config = JsonEncoderConfig {
127            time_handling_mode: TimeHandlingMode::Milli,
128            date_handling_mode: DateHandlingMode::String,
129            timestamp_handling_mode: TimestampHandlingMode::String,
130            timestamptz_handling_mode: TimestamptzHandlingMode::SpecifiedTimezoneWithoutSuffix(
131                time_zone,
132            ),
133            custom_json_type: CustomJsonType::StarRocks,
134            jsonb_handling_mode: JsonbHandlingMode::Dynamic,
135        };
136        Self {
137            schema,
138            col_indices,
139            kafka_connect: None,
140            config,
141        }
142    }
143
144    pub fn new_with_turbopuffer(schema: Schema, col_indices: Option<Vec<usize>>) -> Self {
145        let config = JsonEncoderConfig {
146            time_handling_mode: TimeHandlingMode::String,
147            date_handling_mode: DateHandlingMode::String,
148            timestamp_handling_mode: TimestampHandlingMode::Iso8601String,
149            timestamptz_handling_mode: TimestamptzHandlingMode::UtcString,
150            custom_json_type: CustomJsonType::Turbopuffer,
151            jsonb_handling_mode: JsonbHandlingMode::Dynamic,
152        };
153        Self {
154            schema,
155            col_indices,
156            kafka_connect: None,
157            config,
158        }
159    }
160
161    pub fn with_kafka_connect(self, kafka_connect: KafkaConnectParams) -> Self {
162        Self {
163            kafka_connect: Some(Arc::new(kafka_connect)),
164            ..self
165        }
166    }
167}
168
169impl JsonEncoderConfig {
170    fn jsonb_handling_mode_for_field(&self, field_name: &str) -> JsonbHandlingMode {
171        match &self.custom_json_type {
172            CustomJsonType::Doris(config) if config.variant_columns.contains(field_name) => {
173                JsonbHandlingMode::Dynamic
174            }
175            _ => self.jsonb_handling_mode,
176        }
177    }
178}
179
180impl RowEncoder for JsonEncoder {
181    type Output = Map<String, Value>;
182
183    fn schema(&self) -> &Schema {
184        &self.schema
185    }
186
187    fn col_indices(&self) -> Option<&[usize]> {
188        self.col_indices.as_ref().map(Vec::as_ref)
189    }
190
191    fn encode_cols(
192        &self,
193        row: impl Row,
194        col_indices: impl Iterator<Item = usize>,
195    ) -> Result<Self::Output> {
196        let mut mappings = Map::with_capacity(self.schema.len());
197        let col_indices = col_indices.collect_vec();
198        for idx in &col_indices {
199            let field = &self.schema[*idx];
200            let key = field.name.clone();
201            let value = datum_to_json_object(field, row.datum_at(*idx), &self.config)
202                .map_err(|e| SinkError::Encode(e.to_report_string()))?;
203            mappings.insert(key, value);
204        }
205
206        Ok(if let Some(param) = &self.kafka_connect {
207            json_converter_with_schema(
208                Value::Object(mappings),
209                param.schema_name.clone(),
210                col_indices.into_iter().map(|i| &self.schema[i]),
211            )
212        } else {
213            mappings
214        })
215    }
216}
217
218impl SerTo<String> for Map<String, Value> {
219    fn ser_to(self) -> Result<String> {
220        Value::Object(self).ser_to()
221    }
222}
223
224impl SerTo<String> for Value {
225    fn ser_to(self) -> Result<String> {
226        Ok(self.to_string())
227    }
228}
229
230fn datum_to_json_object(
231    field: &Field,
232    datum: DatumRef<'_>,
233    config: &JsonEncoderConfig,
234) -> ArrayResult<Value> {
235    let scalar_ref = match datum {
236        None => {
237            return Ok(Value::Null);
238        }
239        Some(datum) => datum,
240    };
241
242    let data_type = field.data_type();
243
244    tracing::trace!("datum_to_json_object: {:?}, {:?}", data_type, scalar_ref);
245
246    let value = match (data_type, scalar_ref) {
247        (DataType::Boolean, ScalarRefImpl::Bool(v)) => {
248            json!(v)
249        }
250        (DataType::Int16, ScalarRefImpl::Int16(v)) => {
251            json!(v)
252        }
253        (DataType::Int32, ScalarRefImpl::Int32(v)) => {
254            json!(v)
255        }
256        (DataType::Int64, ScalarRefImpl::Int64(v)) => {
257            json!(v)
258        }
259        (DataType::Serial, ScalarRefImpl::Serial(v)) => {
260            if matches!(&config.custom_json_type, CustomJsonType::Turbopuffer) {
261                json!(v.into_inner())
262            } else {
263                // The serial type needs to be handled as a string to prevent primary key conflicts caused by the precision issues of JSON numbers.
264                json!(format!("{:#018x}", v.into_inner()))
265            }
266        }
267        (DataType::Float32, ScalarRefImpl::Float32(v)) => {
268            json!(f32::from(v))
269        }
270        (DataType::Float64, ScalarRefImpl::Float64(v)) => {
271            json!(f64::from(v))
272        }
273        (DataType::Varchar, ScalarRefImpl::Utf8(v)) => {
274            json!(v)
275        }
276        // Doris/Starrocks will convert out-of-bounds decimal and -INF, INF, NAN to NULL
277        (DataType::Decimal, ScalarRefImpl::Decimal(mut v)) => match &config.custom_json_type {
278            CustomJsonType::Doris(config) => {
279                let s = config.decimal_scale.get(&field.name).unwrap();
280                v.rescale(*s as u32);
281                json!(v.to_text())
282            }
283            CustomJsonType::Turbopuffer => {
284                let value = f64::try_from(v).map_err(|err| {
285                    ArrayError::internal(format!(
286                        "failed to convert decimal to f64: {}",
287                        err.as_report()
288                    ))
289                })?;
290                serde_json::Number::from_f64(value)
291                    .map(Value::Number)
292                    .ok_or_else(|| {
293                        ArrayError::internal(
294                            "failed to encode non-finite decimal as JSON number".to_owned(),
295                        )
296                    })?
297            }
298            CustomJsonType::Es | CustomJsonType::None | CustomJsonType::StarRocks => {
299                json!(v.to_text())
300            }
301        },
302        (DataType::Timestamptz, ScalarRefImpl::Timestamptz(v)) => {
303            match config.timestamptz_handling_mode {
304                TimestamptzHandlingMode::UtcString => {
305                    let parsed = v.to_datetime_utc();
306                    let v = parsed.to_rfc3339_opts(chrono::SecondsFormat::Micros, true);
307                    json!(v)
308                }
309                TimestamptzHandlingMode::UtcWithoutSuffix => {
310                    let parsed = v.to_datetime_utc().naive_utc();
311                    let v = parsed.format("%Y-%m-%d %H:%M:%S%.6f").to_string();
312                    json!(v)
313                }
314                TimestamptzHandlingMode::SpecifiedTimezoneWithoutSuffix(time_zone) => {
315                    let parsed = v.to_datetime_in_zone(time_zone).naive_local();
316                    let v = parsed.format("%Y-%m-%d %H:%M:%S%.6f").to_string();
317                    json!(v)
318                }
319                TimestamptzHandlingMode::Micro => json!(v.timestamp_micros()),
320                TimestamptzHandlingMode::Milli => json!(v.timestamp_millis()),
321            }
322        }
323        (DataType::Time, ScalarRefImpl::Time(v)) => match config.time_handling_mode {
324            TimeHandlingMode::Milli => {
325                // todo: just ignore the nanos part to avoid leap second complex
326                json!(v.0.num_seconds_from_midnight() as i64 * 1000)
327            }
328            TimeHandlingMode::String => {
329                let a = v.0.format("%H:%M:%S%.6f").to_string();
330                json!(a)
331            }
332        },
333        (DataType::Date, ScalarRefImpl::Date(v)) => match config.date_handling_mode {
334            DateHandlingMode::FromCe => json!(v.0.num_days_from_ce()),
335            DateHandlingMode::FromEpoch => {
336                let duration = v.0 - DateTime::UNIX_EPOCH.date_naive();
337                json!(duration.num_days())
338            }
339            DateHandlingMode::String => {
340                let a = v.0.format("%Y-%m-%d").to_string();
341                json!(a)
342            }
343        },
344        (DataType::Timestamp, ScalarRefImpl::Timestamp(v)) => {
345            match config.timestamp_handling_mode {
346                TimestampHandlingMode::Milli => json!(v.0.and_utc().timestamp_millis()),
347                TimestampHandlingMode::String => {
348                    json!(v.0.format("%Y-%m-%d %H:%M:%S%.6f").to_string())
349                }
350                TimestampHandlingMode::Iso8601String => {
351                    json!(v.0.format("%Y-%m-%dT%H:%M:%S%.6f").to_string())
352                }
353            }
354        }
355        (DataType::Bytea, ScalarRefImpl::Bytea(v)) => {
356            json!(general_purpose::STANDARD.encode(v))
357        }
358        // P<years>Y<months>M<days>DT<hours>H<minutes>M<seconds>S
359        (DataType::Interval, ScalarRefImpl::Interval(v)) => {
360            json!(v.as_iso_8601())
361        }
362
363        (DataType::Jsonb, ScalarRefImpl::Jsonb(jsonb_ref)) => {
364            match config.jsonb_handling_mode_for_field(&field.name) {
365                JsonbHandlingMode::String => {
366                    json!(jsonb_ref.to_string())
367                }
368                JsonbHandlingMode::Dynamic => JsonbVal::from(jsonb_ref).take(),
369            }
370        }
371        (DataType::List(lt), ScalarRefImpl::List(list_ref)) => {
372            let elems = list_ref.iter();
373            let mut vec = Vec::with_capacity(elems.len());
374            let inner_field = Field::unnamed(lt.into_elem());
375            for sub_datum_ref in elems {
376                let value = datum_to_json_object(&inner_field, sub_datum_ref, config)?;
377                vec.push(value);
378            }
379            json!(vec)
380        }
381        (DataType::Vector(_), ScalarRefImpl::Vector(vector)) => {
382            let elems = vector.as_raw_slice();
383            let mut vec = Vec::with_capacity(elems.len());
384            for v in elems {
385                let value = serde_json::Number::from_f64(*v as _)
386                    .map(Value::Number)
387                    .unwrap_or(Value::Null);
388                vec.push(value);
389            }
390            json!(vec)
391        }
392        (DataType::Struct(st), ScalarRefImpl::Struct(struct_ref)) => {
393            match config.custom_json_type {
394                CustomJsonType::Doris(_) => {
395                    // We need to ensure that the order of elements in the json matches the insertion order.
396                    let mut map = IndexMap::with_capacity(st.len());
397                    for (sub_datum_ref, sub_field) in struct_ref.iter_fields_ref().zip_eq_debug(
398                        st.iter()
399                            .map(|(name, dt)| Field::with_name(dt.clone(), name)),
400                    ) {
401                        let value = datum_to_json_object(&sub_field, sub_datum_ref, config)?;
402                        map.insert(sub_field.name.clone(), value);
403                    }
404                    Value::String(
405                        serde_json::to_string(&map).context("failed to serialize into JSON")?,
406                    )
407                }
408                CustomJsonType::StarRocks => {
409                    return Err(ArrayError::internal(
410                        "starrocks can't support struct".to_owned(),
411                    ));
412                }
413                CustomJsonType::Es | CustomJsonType::None | CustomJsonType::Turbopuffer => {
414                    let mut map = Map::with_capacity(st.len());
415                    for (sub_datum_ref, sub_field) in struct_ref.iter_fields_ref().zip_eq_debug(
416                        st.iter()
417                            .map(|(name, dt)| Field::with_name(dt.clone(), name)),
418                    ) {
419                        let value = datum_to_json_object(&sub_field, sub_datum_ref, config)?;
420                        map.insert(sub_field.name.clone(), value);
421                    }
422                    json!(map)
423                }
424            }
425        }
426        // TODO(map): support map
427        (data_type, scalar_ref) => {
428            return Err(ArrayError::internal(format!(
429                "datum_to_json_object: unsupported data type: field name: {:?}, logical type: {:?}, physical type: {:?}",
430                field.name, data_type, scalar_ref
431            )));
432        }
433    };
434
435    Ok(value)
436}
437
438fn json_converter_with_schema<'a>(
439    object: Value,
440    name: String,
441    fields: impl Iterator<Item = &'a Field>,
442) -> Map<String, Value> {
443    let mut mapping = Map::with_capacity(2);
444    mapping.insert(
445        "schema".to_owned(),
446        json!({
447            "type": "struct",
448            "fields": fields.map(|field| {
449                let mut mapping = type_as_json_schema(&field.data_type);
450                mapping.insert("field".to_owned(), json!(field.name));
451                mapping
452            }).collect_vec(),
453            "optional": false,
454            "name": name,
455        }),
456    );
457    mapping.insert("payload".to_owned(), object);
458    mapping
459}
460
461// reference: https://github.com/apache/kafka/blob/80982c4ae3fe6be127b48ec09caff11ab5f87c69/connect/json/src/main/java/org/apache/kafka/connect/json/JsonSchema.java#L39
462pub(crate) fn schema_type_mapping(rw_type: &DataType) -> &'static str {
463    match rw_type {
464        DataType::Boolean => "boolean",
465        DataType::Int16 => "int16",
466        DataType::Int32 => "int32",
467        DataType::Int64 => "int64",
468        DataType::Float32 => "float",
469        DataType::Float64 => "double",
470        DataType::Decimal => "string",
471        DataType::Date => "int32",
472        DataType::Varchar => "string",
473        DataType::Time => "int64",
474        DataType::Timestamp => "int64",
475        DataType::Timestamptz => "string",
476        DataType::Interval => "string",
477        DataType::Struct(_) => "struct",
478        DataType::List(_) => "array",
479        DataType::Vector(_) => "array",
480        DataType::Bytea => "bytes",
481        DataType::Jsonb => "string",
482        DataType::Variant => "string",
483        DataType::Serial => "string",
484        DataType::Int256 => "string",
485        DataType::Map(_) => "map",
486    }
487}
488
489fn type_as_json_schema(rw_type: &DataType) -> Map<String, Value> {
490    let mut mapping = Map::with_capacity(4); // type + optional + fields/items + field
491    mapping.insert("type".to_owned(), json!(schema_type_mapping(rw_type)));
492    mapping.insert("optional".to_owned(), json!(true));
493    match rw_type {
494        DataType::Struct(struct_type) => {
495            let sub_fields = struct_type
496                .iter()
497                .map(|(sub_name, sub_type)| {
498                    let mut sub_mapping = type_as_json_schema(sub_type);
499                    sub_mapping.insert("field".to_owned(), json!(sub_name));
500                    sub_mapping
501                })
502                .collect_vec();
503            mapping.insert("fields".to_owned(), json!(sub_fields));
504        }
505        DataType::List(list_type) => {
506            mapping.insert(
507                "items".to_owned(),
508                json!(type_as_json_schema(list_type.elem())),
509            );
510        }
511        _ => {}
512    }
513
514    mapping
515}
516
517#[cfg(test)]
518mod tests {
519    use risingwave_common::row::OwnedRow;
520    use risingwave_common::types::{
521        Date, Decimal, Interval, Scalar, ScalarImpl, StructRef, StructType, StructValue, Time,
522        Timestamp,
523    };
524
525    use super::*;
526
527    #[test]
528    fn test_starrocks_timestamptz_encoding() {
529        let schema = Schema::new(vec![Field::with_name(DataType::Timestamptz, "ts")]);
530        let encoder = JsonEncoder::new_with_starrocks(schema, None, chrono_tz::Asia::Shanghai);
531        let tstz = "2018-01-26T18:30:09.453Z".parse().unwrap();
532        let row = OwnedRow::new(vec![Some(ScalarImpl::Timestamptz(tstz))]);
533
534        let encoded = encoder.encode(row).unwrap();
535
536        assert_eq!(
537            encoded.get("ts"),
538            Some(&json!("2018-01-27 02:30:09.453000"))
539        );
540    }
541
542    #[test]
543    fn test_to_json_basic_type() {
544        let mock_field = Field {
545            data_type: DataType::Boolean,
546            name: Default::default(),
547        };
548
549        let config = JsonEncoderConfig {
550            time_handling_mode: TimeHandlingMode::Milli,
551            date_handling_mode: DateHandlingMode::FromCe,
552            timestamp_handling_mode: TimestampHandlingMode::String,
553            timestamptz_handling_mode: TimestamptzHandlingMode::UtcString,
554            custom_json_type: CustomJsonType::None,
555            jsonb_handling_mode: JsonbHandlingMode::String,
556        };
557
558        let boolean_value = datum_to_json_object(
559            &Field {
560                data_type: DataType::Boolean,
561                ..mock_field.clone()
562            },
563            Some(ScalarImpl::Bool(false).as_scalar_ref_impl()),
564            &config,
565        )
566        .unwrap();
567        assert_eq!(boolean_value, json!(false));
568
569        let int16_value = datum_to_json_object(
570            &Field {
571                data_type: DataType::Int16,
572                ..mock_field.clone()
573            },
574            Some(ScalarImpl::Int16(16).as_scalar_ref_impl()),
575            &config,
576        )
577        .unwrap();
578        assert_eq!(int16_value, json!(16));
579
580        let int64_value = datum_to_json_object(
581            &Field {
582                data_type: DataType::Int64,
583                ..mock_field.clone()
584            },
585            Some(ScalarImpl::Int64(i64::MAX).as_scalar_ref_impl()),
586            &config,
587        )
588        .unwrap();
589        assert_eq!(
590            serde_json::to_string(&int64_value).unwrap(),
591            i64::MAX.to_string()
592        );
593
594        let serial_value = datum_to_json_object(
595            &Field {
596                data_type: DataType::Serial,
597                ..mock_field.clone()
598            },
599            Some(ScalarImpl::Serial(i64::MAX.into()).as_scalar_ref_impl()),
600            &config,
601        )
602        .unwrap();
603        assert_eq!(
604            serde_json::to_string(&serial_value).unwrap(),
605            format!("\"{:#018x}\"", i64::MAX)
606        );
607
608        let turbopuffer_config = JsonEncoderConfig {
609            time_handling_mode: TimeHandlingMode::String,
610            date_handling_mode: DateHandlingMode::String,
611            timestamp_handling_mode: TimestampHandlingMode::String,
612            timestamptz_handling_mode: TimestamptzHandlingMode::UtcString,
613            custom_json_type: CustomJsonType::Turbopuffer,
614            jsonb_handling_mode: JsonbHandlingMode::Dynamic,
615        };
616        let turbopuffer_serial_value = datum_to_json_object(
617            &Field {
618                data_type: DataType::Serial,
619                ..mock_field.clone()
620            },
621            Some(ScalarImpl::Serial(i64::MAX.into()).as_scalar_ref_impl()),
622            &turbopuffer_config,
623        )
624        .unwrap();
625        assert_eq!(turbopuffer_serial_value, json!(i64::MAX));
626        let turbopuffer_decimal_value = datum_to_json_object(
627            &Field {
628                data_type: DataType::Decimal,
629                ..mock_field.clone()
630            },
631            Some(ScalarImpl::Decimal(Decimal::try_from(1.25).unwrap()).as_scalar_ref_impl()),
632            &turbopuffer_config,
633        )
634        .unwrap();
635        assert_eq!(turbopuffer_decimal_value, json!(1.25));
636        assert!(
637            datum_to_json_object(
638                &Field {
639                    data_type: DataType::Decimal,
640                    ..mock_field.clone()
641                },
642                Some(ScalarImpl::Decimal(Decimal::NaN).as_scalar_ref_impl()),
643                &turbopuffer_config,
644            )
645            .unwrap_err()
646            .to_string()
647            .contains("non-finite decimal")
648        );
649
650        // https://github.com/debezium/debezium/blob/main/debezium-core/src/main/java/io/debezium/time/ZonedTimestamp.java
651        let tstz_inner = "2018-01-26T18:30:09.453Z".parse().unwrap();
652        let tstz_value = datum_to_json_object(
653            &Field {
654                data_type: DataType::Timestamptz,
655                ..mock_field.clone()
656            },
657            Some(ScalarImpl::Timestamptz(tstz_inner).as_scalar_ref_impl()),
658            &config,
659        )
660        .unwrap();
661        assert_eq!(tstz_value, "2018-01-26T18:30:09.453000Z");
662
663        let unix_wo_suffix_config = JsonEncoderConfig {
664            time_handling_mode: TimeHandlingMode::Milli,
665            date_handling_mode: DateHandlingMode::FromCe,
666            timestamp_handling_mode: TimestampHandlingMode::String,
667            timestamptz_handling_mode: TimestamptzHandlingMode::UtcWithoutSuffix,
668            custom_json_type: CustomJsonType::None,
669            jsonb_handling_mode: JsonbHandlingMode::String,
670        };
671
672        let tstz_inner = "2018-01-26T18:30:09.453Z".parse().unwrap();
673        let tstz_value = datum_to_json_object(
674            &Field {
675                data_type: DataType::Timestamptz,
676                ..mock_field.clone()
677            },
678            Some(ScalarImpl::Timestamptz(tstz_inner).as_scalar_ref_impl()),
679            &unix_wo_suffix_config,
680        )
681        .unwrap();
682        assert_eq!(tstz_value, "2018-01-26 18:30:09.453000");
683
684        let timestamp_milli_config = JsonEncoderConfig {
685            time_handling_mode: TimeHandlingMode::String,
686            date_handling_mode: DateHandlingMode::FromCe,
687            timestamp_handling_mode: TimestampHandlingMode::Milli,
688            timestamptz_handling_mode: TimestamptzHandlingMode::UtcString,
689            custom_json_type: CustomJsonType::None,
690            jsonb_handling_mode: JsonbHandlingMode::String,
691        };
692        let ts_value = datum_to_json_object(
693            &Field {
694                data_type: DataType::Timestamp,
695                ..mock_field.clone()
696            },
697            Some(
698                ScalarImpl::Timestamp(Timestamp::from_timestamp_uncheck(1000, 0))
699                    .as_scalar_ref_impl(),
700            ),
701            &timestamp_milli_config,
702        )
703        .unwrap();
704        assert_eq!(ts_value, json!(1000 * 1000));
705
706        let ts_value = datum_to_json_object(
707            &Field {
708                data_type: DataType::Timestamp,
709                ..mock_field.clone()
710            },
711            Some(
712                ScalarImpl::Timestamp(Timestamp::from_timestamp_uncheck(1000, 0))
713                    .as_scalar_ref_impl(),
714            ),
715            &config,
716        )
717        .unwrap();
718        assert_eq!(ts_value, json!("1970-01-01 00:16:40.000000".to_owned()));
719
720        // Represents the number of milliseconds past midnigh, org.apache.kafka.connect.data.Time
721        let time_value = datum_to_json_object(
722            &Field {
723                data_type: DataType::Time,
724                ..mock_field.clone()
725            },
726            Some(
727                ScalarImpl::Time(Time::from_num_seconds_from_midnight_uncheck(1000, 0))
728                    .as_scalar_ref_impl(),
729            ),
730            &config,
731        )
732        .unwrap();
733        assert_eq!(time_value, json!(1000 * 1000));
734
735        let interval_value = datum_to_json_object(
736            &Field {
737                data_type: DataType::Interval,
738                ..mock_field.clone()
739            },
740            Some(
741                ScalarImpl::Interval(Interval::from_month_day_usec(13, 2, 1000000))
742                    .as_scalar_ref_impl(),
743            ),
744            &config,
745        )
746        .unwrap();
747        assert_eq!(interval_value, json!("P1Y1M2DT0H0M1S"));
748
749        let mut decimal_scale = HashMap::default();
750        decimal_scale.insert("aaa".to_owned(), 5_u8);
751        let doris_config = JsonEncoderConfig {
752            time_handling_mode: TimeHandlingMode::String,
753            date_handling_mode: DateHandlingMode::String,
754            timestamp_handling_mode: TimestampHandlingMode::String,
755            timestamptz_handling_mode: TimestamptzHandlingMode::UtcString,
756            custom_json_type: CustomJsonType::Doris(DorisJsonConfig {
757                decimal_scale,
758                variant_columns: HashSet::default(),
759            }),
760            jsonb_handling_mode: JsonbHandlingMode::String,
761        };
762        let decimal = datum_to_json_object(
763            &Field {
764                data_type: DataType::Decimal,
765                name: "aaa".to_owned(),
766            },
767            Some(ScalarImpl::Decimal(Decimal::try_from(1.1111111).unwrap()).as_scalar_ref_impl()),
768            &doris_config,
769        )
770        .unwrap();
771        assert_eq!(decimal, json!("1.11111"));
772
773        let date_value = datum_to_json_object(
774            &Field {
775                data_type: DataType::Date,
776                ..mock_field.clone()
777            },
778            Some(ScalarImpl::Date(Date::from_ymd_uncheck(1970, 1, 1)).as_scalar_ref_impl()),
779            &config,
780        )
781        .unwrap();
782        assert_eq!(date_value, json!(719163));
783
784        let from_epoch_config = JsonEncoderConfig {
785            time_handling_mode: TimeHandlingMode::String,
786            date_handling_mode: DateHandlingMode::FromEpoch,
787            timestamp_handling_mode: TimestampHandlingMode::String,
788            timestamptz_handling_mode: TimestamptzHandlingMode::UtcString,
789            custom_json_type: CustomJsonType::None,
790            jsonb_handling_mode: JsonbHandlingMode::String,
791        };
792        let date_value = datum_to_json_object(
793            &Field {
794                data_type: DataType::Date,
795                ..mock_field.clone()
796            },
797            Some(ScalarImpl::Date(Date::from_ymd_uncheck(1970, 1, 1)).as_scalar_ref_impl()),
798            &from_epoch_config,
799        )
800        .unwrap();
801        assert_eq!(date_value, json!(0));
802
803        let doris_config = JsonEncoderConfig {
804            time_handling_mode: TimeHandlingMode::String,
805            date_handling_mode: DateHandlingMode::String,
806            timestamp_handling_mode: TimestampHandlingMode::String,
807            timestamptz_handling_mode: TimestamptzHandlingMode::UtcString,
808            custom_json_type: CustomJsonType::Doris(DorisJsonConfig {
809                decimal_scale: HashMap::default(),
810                variant_columns: HashSet::default(),
811            }),
812            jsonb_handling_mode: JsonbHandlingMode::String,
813        };
814        let date_value = datum_to_json_object(
815            &Field {
816                data_type: DataType::Date,
817                ..mock_field.clone()
818            },
819            Some(ScalarImpl::Date(Date::from_ymd_uncheck(2010, 10, 10)).as_scalar_ref_impl()),
820            &doris_config,
821        )
822        .unwrap();
823        assert_eq!(date_value, json!("2010-10-10"));
824
825        let value = StructValue::new(vec![
826            Some(3_i32.to_scalar_value()),
827            Some(2_i32.to_scalar_value()),
828            Some(1_i32.to_scalar_value()),
829        ]);
830
831        let interval_value = datum_to_json_object(
832            &Field {
833                data_type: DataType::Struct(StructType::new(vec![
834                    ("v3", DataType::Int32),
835                    ("v2", DataType::Int32),
836                    ("v1", DataType::Int32),
837                ])),
838                ..mock_field.clone()
839            },
840            Some(ScalarRefImpl::Struct(StructRef::ValueRef { val: &value })),
841            &doris_config,
842        )
843        .unwrap();
844        assert_eq!(interval_value, json!("{\"v3\":3,\"v2\":2,\"v1\":1}"));
845
846        let encode_jsonb_obj_config = JsonEncoderConfig {
847            time_handling_mode: TimeHandlingMode::String,
848            date_handling_mode: DateHandlingMode::String,
849            timestamp_handling_mode: TimestampHandlingMode::String,
850            timestamptz_handling_mode: TimestamptzHandlingMode::UtcString,
851            custom_json_type: CustomJsonType::None,
852            jsonb_handling_mode: JsonbHandlingMode::Dynamic,
853        };
854        let json_value = datum_to_json_object(
855            &Field {
856                data_type: DataType::Jsonb,
857                ..mock_field
858            },
859            Some(ScalarImpl::Jsonb(JsonbVal::from(json!([1, 2, 3]))).as_scalar_ref_impl()),
860            &encode_jsonb_obj_config,
861        )
862        .unwrap();
863        assert_eq!(json_value, json!([1, 2, 3]));
864
865        let variant_json_value = datum_to_json_object(
866            &Field {
867                data_type: DataType::Jsonb,
868                name: "variant_col".into(),
869            },
870            Some(
871                ScalarImpl::Jsonb(JsonbVal::from(json!({"nested": [1, 2, 3]})))
872                    .as_scalar_ref_impl(),
873            ),
874            &JsonEncoderConfig {
875                time_handling_mode: TimeHandlingMode::String,
876                date_handling_mode: DateHandlingMode::String,
877                timestamp_handling_mode: TimestampHandlingMode::String,
878                timestamptz_handling_mode: TimestamptzHandlingMode::UtcString,
879                custom_json_type: CustomJsonType::Doris(DorisJsonConfig {
880                    decimal_scale: HashMap::default(),
881                    variant_columns: HashSet::from(["variant_col".to_owned()]),
882                }),
883                jsonb_handling_mode: JsonbHandlingMode::String,
884            },
885        )
886        .unwrap();
887        assert_eq!(variant_json_value, json!({"nested": [1, 2, 3]}));
888    }
889
890    #[test]
891    fn test_generate_json_converter_schema() {
892        let fields = vec![
893            Field {
894                data_type: DataType::Boolean,
895                name: "v1".into(),
896            },
897            Field {
898                data_type: DataType::Int16,
899                name: "v2".into(),
900            },
901            Field {
902                data_type: DataType::Int32,
903                name: "v3".into(),
904            },
905            Field {
906                data_type: DataType::Float32,
907                name: "v4".into(),
908            },
909            Field {
910                data_type: DataType::Decimal,
911                name: "v5".into(),
912            },
913            Field {
914                data_type: DataType::Date,
915                name: "v6".into(),
916            },
917            Field {
918                data_type: DataType::Varchar,
919                name: "v7".into(),
920            },
921            Field {
922                data_type: DataType::Time,
923                name: "v8".into(),
924            },
925            Field {
926                data_type: DataType::Interval,
927                name: "v9".into(),
928            },
929            Field {
930                data_type: DataType::Struct(StructType::new(vec![
931                    ("a", DataType::Timestamp),
932                    ("b", DataType::Timestamptz),
933                    (
934                        "c",
935                        DataType::Struct(StructType::new(vec![
936                            ("aa", DataType::Int64),
937                            ("bb", DataType::Float64),
938                        ])),
939                    ),
940                ])),
941                name: "v10".into(),
942            },
943            Field {
944                data_type: DataType::list(DataType::list(DataType::Struct(StructType::new(vec![
945                    ("aa", DataType::Int64),
946                    ("bb", DataType::Float64),
947                ])))),
948                name: "v11".into(),
949            },
950            Field {
951                data_type: DataType::Jsonb,
952                name: "12".into(),
953            },
954            Field {
955                data_type: DataType::Serial,
956                name: "13".into(),
957            },
958            Field {
959                data_type: DataType::Int256,
960                name: "14".into(),
961            },
962        ];
963        let schema =
964            json_converter_with_schema(json!({}), "test".to_owned(), fields.iter())["schema"]
965                .to_string();
966        let ans = r#"{"fields":[{"field":"v1","optional":true,"type":"boolean"},{"field":"v2","optional":true,"type":"int16"},{"field":"v3","optional":true,"type":"int32"},{"field":"v4","optional":true,"type":"float"},{"field":"v5","optional":true,"type":"string"},{"field":"v6","optional":true,"type":"int32"},{"field":"v7","optional":true,"type":"string"},{"field":"v8","optional":true,"type":"int64"},{"field":"v9","optional":true,"type":"string"},{"field":"v10","fields":[{"field":"a","optional":true,"type":"int64"},{"field":"b","optional":true,"type":"string"},{"field":"c","fields":[{"field":"aa","optional":true,"type":"int64"},{"field":"bb","optional":true,"type":"double"}],"optional":true,"type":"struct"}],"optional":true,"type":"struct"},{"field":"v11","items":{"items":{"fields":[{"field":"aa","optional":true,"type":"int64"},{"field":"bb","optional":true,"type":"double"}],"optional":true,"type":"struct"},"optional":true,"type":"array"},"optional":true,"type":"array"},{"field":"12","optional":true,"type":"string"},{"field":"13","optional":true,"type":"string"},{"field":"14","optional":true,"type":"string"}],"name":"test","optional":false,"type":"struct"}"#;
967        assert_eq!(schema, ans);
968    }
969}