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::Serial => "string",
483        DataType::Int256 => "string",
484        DataType::Map(_) => "map",
485    }
486}
487
488fn type_as_json_schema(rw_type: &DataType) -> Map<String, Value> {
489    let mut mapping = Map::with_capacity(4); // type + optional + fields/items + field
490    mapping.insert("type".to_owned(), json!(schema_type_mapping(rw_type)));
491    mapping.insert("optional".to_owned(), json!(true));
492    match rw_type {
493        DataType::Struct(struct_type) => {
494            let sub_fields = struct_type
495                .iter()
496                .map(|(sub_name, sub_type)| {
497                    let mut sub_mapping = type_as_json_schema(sub_type);
498                    sub_mapping.insert("field".to_owned(), json!(sub_name));
499                    sub_mapping
500                })
501                .collect_vec();
502            mapping.insert("fields".to_owned(), json!(sub_fields));
503        }
504        DataType::List(list_type) => {
505            mapping.insert(
506                "items".to_owned(),
507                json!(type_as_json_schema(list_type.elem())),
508            );
509        }
510        _ => {}
511    }
512
513    mapping
514}
515
516#[cfg(test)]
517mod tests {
518    use risingwave_common::row::OwnedRow;
519    use risingwave_common::types::{
520        Date, Decimal, Interval, Scalar, ScalarImpl, StructRef, StructType, StructValue, Time,
521        Timestamp,
522    };
523
524    use super::*;
525
526    #[test]
527    fn test_starrocks_timestamptz_encoding() {
528        let schema = Schema::new(vec![Field::with_name(DataType::Timestamptz, "ts")]);
529        let encoder = JsonEncoder::new_with_starrocks(schema, None, chrono_tz::Asia::Shanghai);
530        let tstz = "2018-01-26T18:30:09.453Z".parse().unwrap();
531        let row = OwnedRow::new(vec![Some(ScalarImpl::Timestamptz(tstz))]);
532
533        let encoded = encoder.encode(row).unwrap();
534
535        assert_eq!(
536            encoded.get("ts"),
537            Some(&json!("2018-01-27 02:30:09.453000"))
538        );
539    }
540
541    #[test]
542    fn test_to_json_basic_type() {
543        let mock_field = Field {
544            data_type: DataType::Boolean,
545            name: Default::default(),
546        };
547
548        let config = JsonEncoderConfig {
549            time_handling_mode: TimeHandlingMode::Milli,
550            date_handling_mode: DateHandlingMode::FromCe,
551            timestamp_handling_mode: TimestampHandlingMode::String,
552            timestamptz_handling_mode: TimestamptzHandlingMode::UtcString,
553            custom_json_type: CustomJsonType::None,
554            jsonb_handling_mode: JsonbHandlingMode::String,
555        };
556
557        let boolean_value = datum_to_json_object(
558            &Field {
559                data_type: DataType::Boolean,
560                ..mock_field.clone()
561            },
562            Some(ScalarImpl::Bool(false).as_scalar_ref_impl()),
563            &config,
564        )
565        .unwrap();
566        assert_eq!(boolean_value, json!(false));
567
568        let int16_value = datum_to_json_object(
569            &Field {
570                data_type: DataType::Int16,
571                ..mock_field.clone()
572            },
573            Some(ScalarImpl::Int16(16).as_scalar_ref_impl()),
574            &config,
575        )
576        .unwrap();
577        assert_eq!(int16_value, json!(16));
578
579        let int64_value = datum_to_json_object(
580            &Field {
581                data_type: DataType::Int64,
582                ..mock_field.clone()
583            },
584            Some(ScalarImpl::Int64(i64::MAX).as_scalar_ref_impl()),
585            &config,
586        )
587        .unwrap();
588        assert_eq!(
589            serde_json::to_string(&int64_value).unwrap(),
590            i64::MAX.to_string()
591        );
592
593        let serial_value = datum_to_json_object(
594            &Field {
595                data_type: DataType::Serial,
596                ..mock_field.clone()
597            },
598            Some(ScalarImpl::Serial(i64::MAX.into()).as_scalar_ref_impl()),
599            &config,
600        )
601        .unwrap();
602        assert_eq!(
603            serde_json::to_string(&serial_value).unwrap(),
604            format!("\"{:#018x}\"", i64::MAX)
605        );
606
607        let turbopuffer_config = JsonEncoderConfig {
608            time_handling_mode: TimeHandlingMode::String,
609            date_handling_mode: DateHandlingMode::String,
610            timestamp_handling_mode: TimestampHandlingMode::String,
611            timestamptz_handling_mode: TimestamptzHandlingMode::UtcString,
612            custom_json_type: CustomJsonType::Turbopuffer,
613            jsonb_handling_mode: JsonbHandlingMode::Dynamic,
614        };
615        let turbopuffer_serial_value = datum_to_json_object(
616            &Field {
617                data_type: DataType::Serial,
618                ..mock_field.clone()
619            },
620            Some(ScalarImpl::Serial(i64::MAX.into()).as_scalar_ref_impl()),
621            &turbopuffer_config,
622        )
623        .unwrap();
624        assert_eq!(turbopuffer_serial_value, json!(i64::MAX));
625        let turbopuffer_decimal_value = datum_to_json_object(
626            &Field {
627                data_type: DataType::Decimal,
628                ..mock_field.clone()
629            },
630            Some(ScalarImpl::Decimal(Decimal::try_from(1.25).unwrap()).as_scalar_ref_impl()),
631            &turbopuffer_config,
632        )
633        .unwrap();
634        assert_eq!(turbopuffer_decimal_value, json!(1.25));
635        assert!(
636            datum_to_json_object(
637                &Field {
638                    data_type: DataType::Decimal,
639                    ..mock_field.clone()
640                },
641                Some(ScalarImpl::Decimal(Decimal::NaN).as_scalar_ref_impl()),
642                &turbopuffer_config,
643            )
644            .unwrap_err()
645            .to_string()
646            .contains("non-finite decimal")
647        );
648
649        // https://github.com/debezium/debezium/blob/main/debezium-core/src/main/java/io/debezium/time/ZonedTimestamp.java
650        let tstz_inner = "2018-01-26T18:30:09.453Z".parse().unwrap();
651        let tstz_value = datum_to_json_object(
652            &Field {
653                data_type: DataType::Timestamptz,
654                ..mock_field.clone()
655            },
656            Some(ScalarImpl::Timestamptz(tstz_inner).as_scalar_ref_impl()),
657            &config,
658        )
659        .unwrap();
660        assert_eq!(tstz_value, "2018-01-26T18:30:09.453000Z");
661
662        let unix_wo_suffix_config = JsonEncoderConfig {
663            time_handling_mode: TimeHandlingMode::Milli,
664            date_handling_mode: DateHandlingMode::FromCe,
665            timestamp_handling_mode: TimestampHandlingMode::String,
666            timestamptz_handling_mode: TimestamptzHandlingMode::UtcWithoutSuffix,
667            custom_json_type: CustomJsonType::None,
668            jsonb_handling_mode: JsonbHandlingMode::String,
669        };
670
671        let tstz_inner = "2018-01-26T18:30:09.453Z".parse().unwrap();
672        let tstz_value = datum_to_json_object(
673            &Field {
674                data_type: DataType::Timestamptz,
675                ..mock_field.clone()
676            },
677            Some(ScalarImpl::Timestamptz(tstz_inner).as_scalar_ref_impl()),
678            &unix_wo_suffix_config,
679        )
680        .unwrap();
681        assert_eq!(tstz_value, "2018-01-26 18:30:09.453000");
682
683        let timestamp_milli_config = JsonEncoderConfig {
684            time_handling_mode: TimeHandlingMode::String,
685            date_handling_mode: DateHandlingMode::FromCe,
686            timestamp_handling_mode: TimestampHandlingMode::Milli,
687            timestamptz_handling_mode: TimestamptzHandlingMode::UtcString,
688            custom_json_type: CustomJsonType::None,
689            jsonb_handling_mode: JsonbHandlingMode::String,
690        };
691        let ts_value = datum_to_json_object(
692            &Field {
693                data_type: DataType::Timestamp,
694                ..mock_field.clone()
695            },
696            Some(
697                ScalarImpl::Timestamp(Timestamp::from_timestamp_uncheck(1000, 0))
698                    .as_scalar_ref_impl(),
699            ),
700            &timestamp_milli_config,
701        )
702        .unwrap();
703        assert_eq!(ts_value, json!(1000 * 1000));
704
705        let ts_value = datum_to_json_object(
706            &Field {
707                data_type: DataType::Timestamp,
708                ..mock_field.clone()
709            },
710            Some(
711                ScalarImpl::Timestamp(Timestamp::from_timestamp_uncheck(1000, 0))
712                    .as_scalar_ref_impl(),
713            ),
714            &config,
715        )
716        .unwrap();
717        assert_eq!(ts_value, json!("1970-01-01 00:16:40.000000".to_owned()));
718
719        // Represents the number of milliseconds past midnigh, org.apache.kafka.connect.data.Time
720        let time_value = datum_to_json_object(
721            &Field {
722                data_type: DataType::Time,
723                ..mock_field.clone()
724            },
725            Some(
726                ScalarImpl::Time(Time::from_num_seconds_from_midnight_uncheck(1000, 0))
727                    .as_scalar_ref_impl(),
728            ),
729            &config,
730        )
731        .unwrap();
732        assert_eq!(time_value, json!(1000 * 1000));
733
734        let interval_value = datum_to_json_object(
735            &Field {
736                data_type: DataType::Interval,
737                ..mock_field.clone()
738            },
739            Some(
740                ScalarImpl::Interval(Interval::from_month_day_usec(13, 2, 1000000))
741                    .as_scalar_ref_impl(),
742            ),
743            &config,
744        )
745        .unwrap();
746        assert_eq!(interval_value, json!("P1Y1M2DT0H0M1S"));
747
748        let mut decimal_scale = HashMap::default();
749        decimal_scale.insert("aaa".to_owned(), 5_u8);
750        let doris_config = JsonEncoderConfig {
751            time_handling_mode: TimeHandlingMode::String,
752            date_handling_mode: DateHandlingMode::String,
753            timestamp_handling_mode: TimestampHandlingMode::String,
754            timestamptz_handling_mode: TimestamptzHandlingMode::UtcString,
755            custom_json_type: CustomJsonType::Doris(DorisJsonConfig {
756                decimal_scale,
757                variant_columns: HashSet::default(),
758            }),
759            jsonb_handling_mode: JsonbHandlingMode::String,
760        };
761        let decimal = datum_to_json_object(
762            &Field {
763                data_type: DataType::Decimal,
764                name: "aaa".to_owned(),
765            },
766            Some(ScalarImpl::Decimal(Decimal::try_from(1.1111111).unwrap()).as_scalar_ref_impl()),
767            &doris_config,
768        )
769        .unwrap();
770        assert_eq!(decimal, json!("1.11111"));
771
772        let date_value = datum_to_json_object(
773            &Field {
774                data_type: DataType::Date,
775                ..mock_field.clone()
776            },
777            Some(ScalarImpl::Date(Date::from_ymd_uncheck(1970, 1, 1)).as_scalar_ref_impl()),
778            &config,
779        )
780        .unwrap();
781        assert_eq!(date_value, json!(719163));
782
783        let from_epoch_config = JsonEncoderConfig {
784            time_handling_mode: TimeHandlingMode::String,
785            date_handling_mode: DateHandlingMode::FromEpoch,
786            timestamp_handling_mode: TimestampHandlingMode::String,
787            timestamptz_handling_mode: TimestamptzHandlingMode::UtcString,
788            custom_json_type: CustomJsonType::None,
789            jsonb_handling_mode: JsonbHandlingMode::String,
790        };
791        let date_value = datum_to_json_object(
792            &Field {
793                data_type: DataType::Date,
794                ..mock_field.clone()
795            },
796            Some(ScalarImpl::Date(Date::from_ymd_uncheck(1970, 1, 1)).as_scalar_ref_impl()),
797            &from_epoch_config,
798        )
799        .unwrap();
800        assert_eq!(date_value, json!(0));
801
802        let doris_config = JsonEncoderConfig {
803            time_handling_mode: TimeHandlingMode::String,
804            date_handling_mode: DateHandlingMode::String,
805            timestamp_handling_mode: TimestampHandlingMode::String,
806            timestamptz_handling_mode: TimestamptzHandlingMode::UtcString,
807            custom_json_type: CustomJsonType::Doris(DorisJsonConfig {
808                decimal_scale: HashMap::default(),
809                variant_columns: HashSet::default(),
810            }),
811            jsonb_handling_mode: JsonbHandlingMode::String,
812        };
813        let date_value = datum_to_json_object(
814            &Field {
815                data_type: DataType::Date,
816                ..mock_field.clone()
817            },
818            Some(ScalarImpl::Date(Date::from_ymd_uncheck(2010, 10, 10)).as_scalar_ref_impl()),
819            &doris_config,
820        )
821        .unwrap();
822        assert_eq!(date_value, json!("2010-10-10"));
823
824        let value = StructValue::new(vec![
825            Some(3_i32.to_scalar_value()),
826            Some(2_i32.to_scalar_value()),
827            Some(1_i32.to_scalar_value()),
828        ]);
829
830        let interval_value = datum_to_json_object(
831            &Field {
832                data_type: DataType::Struct(StructType::new(vec![
833                    ("v3", DataType::Int32),
834                    ("v2", DataType::Int32),
835                    ("v1", DataType::Int32),
836                ])),
837                ..mock_field.clone()
838            },
839            Some(ScalarRefImpl::Struct(StructRef::ValueRef { val: &value })),
840            &doris_config,
841        )
842        .unwrap();
843        assert_eq!(interval_value, json!("{\"v3\":3,\"v2\":2,\"v1\":1}"));
844
845        let encode_jsonb_obj_config = JsonEncoderConfig {
846            time_handling_mode: TimeHandlingMode::String,
847            date_handling_mode: DateHandlingMode::String,
848            timestamp_handling_mode: TimestampHandlingMode::String,
849            timestamptz_handling_mode: TimestamptzHandlingMode::UtcString,
850            custom_json_type: CustomJsonType::None,
851            jsonb_handling_mode: JsonbHandlingMode::Dynamic,
852        };
853        let json_value = datum_to_json_object(
854            &Field {
855                data_type: DataType::Jsonb,
856                ..mock_field
857            },
858            Some(ScalarImpl::Jsonb(JsonbVal::from(json!([1, 2, 3]))).as_scalar_ref_impl()),
859            &encode_jsonb_obj_config,
860        )
861        .unwrap();
862        assert_eq!(json_value, json!([1, 2, 3]));
863
864        let variant_json_value = datum_to_json_object(
865            &Field {
866                data_type: DataType::Jsonb,
867                name: "variant_col".into(),
868            },
869            Some(
870                ScalarImpl::Jsonb(JsonbVal::from(json!({"nested": [1, 2, 3]})))
871                    .as_scalar_ref_impl(),
872            ),
873            &JsonEncoderConfig {
874                time_handling_mode: TimeHandlingMode::String,
875                date_handling_mode: DateHandlingMode::String,
876                timestamp_handling_mode: TimestampHandlingMode::String,
877                timestamptz_handling_mode: TimestamptzHandlingMode::UtcString,
878                custom_json_type: CustomJsonType::Doris(DorisJsonConfig {
879                    decimal_scale: HashMap::default(),
880                    variant_columns: HashSet::from(["variant_col".to_owned()]),
881                }),
882                jsonb_handling_mode: JsonbHandlingMode::String,
883            },
884        )
885        .unwrap();
886        assert_eq!(variant_json_value, json!({"nested": [1, 2, 3]}));
887    }
888
889    #[test]
890    fn test_generate_json_converter_schema() {
891        let fields = vec![
892            Field {
893                data_type: DataType::Boolean,
894                name: "v1".into(),
895            },
896            Field {
897                data_type: DataType::Int16,
898                name: "v2".into(),
899            },
900            Field {
901                data_type: DataType::Int32,
902                name: "v3".into(),
903            },
904            Field {
905                data_type: DataType::Float32,
906                name: "v4".into(),
907            },
908            Field {
909                data_type: DataType::Decimal,
910                name: "v5".into(),
911            },
912            Field {
913                data_type: DataType::Date,
914                name: "v6".into(),
915            },
916            Field {
917                data_type: DataType::Varchar,
918                name: "v7".into(),
919            },
920            Field {
921                data_type: DataType::Time,
922                name: "v8".into(),
923            },
924            Field {
925                data_type: DataType::Interval,
926                name: "v9".into(),
927            },
928            Field {
929                data_type: DataType::Struct(StructType::new(vec![
930                    ("a", DataType::Timestamp),
931                    ("b", DataType::Timestamptz),
932                    (
933                        "c",
934                        DataType::Struct(StructType::new(vec![
935                            ("aa", DataType::Int64),
936                            ("bb", DataType::Float64),
937                        ])),
938                    ),
939                ])),
940                name: "v10".into(),
941            },
942            Field {
943                data_type: DataType::list(DataType::list(DataType::Struct(StructType::new(vec![
944                    ("aa", DataType::Int64),
945                    ("bb", DataType::Float64),
946                ])))),
947                name: "v11".into(),
948            },
949            Field {
950                data_type: DataType::Jsonb,
951                name: "12".into(),
952            },
953            Field {
954                data_type: DataType::Serial,
955                name: "13".into(),
956            },
957            Field {
958                data_type: DataType::Int256,
959                name: "14".into(),
960            },
961        ];
962        let schema =
963            json_converter_with_schema(json!({}), "test".to_owned(), fields.iter())["schema"]
964                .to_string();
965        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"}"#;
966        assert_eq!(schema, ans);
967    }
968}