1#[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 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 (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 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 (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 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 (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
461pub(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); 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 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 ×tamp_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 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}