1use std::collections::{HashMap, HashSet};
16use std::sync::Arc;
17
18use arrow_schema::extension::ExtensionType;
19use datatypes::arrow::datatypes::{DataType as ArrowDataType, Field, Schema, SchemaRef};
20use datatypes::arrow::record_batch::RecordBatch;
21use datatypes::extension::json::{JSON2_REMAINDER_FIELD_NAME, Json2ExtensionType, JsonMetadata};
22use datatypes::json::{JSON2_DEFAULT_MAX_AUTO_EXPANDED_PATHS, JsonSettings, JsonTypeHint};
23use datatypes::prelude::ConcreteDataType;
24use datatypes::types::json_type::JsonNativeType;
25use datatypes::vectors::json::array::JsonArray;
26use datatypes::vectors::json::json2_physical_data_type;
27use parquet::arrow::parquet_to_arrow_schema;
28use parquet::file::metadata::ParquetMetaData;
29use snafu::{OptionExt, ResultExt, ensure};
30use store_api::metadata::RegionMetadataRef;
31
32use crate::error::{
33 ConvertValueSnafu, DataTypeMismatchSnafu, InvalidRecordBatchSnafu, NewRecordBatchSnafu, Result,
34};
35
36pub(crate) struct Json2RewritePlan {
43 logical_settings: JsonSettings,
45 pub(super) target_layout: JsonSettings,
47}
48
49pub(crate) type Json2RewritePlans = HashMap<String, Json2RewritePlan>;
51
52#[derive(Clone)]
53struct Json2LeafPathStats {
54 rows: u64,
55 data_type: JsonNativeType,
56 is_type_conflicted: bool,
57}
58
59pub(super) fn collect_json2_rewrite_plans_from_parquet(
70 metadata: &RegionMetadataRef,
71 parquet_metadata: &[Arc<ParquetMetaData>],
72) -> Result<Json2RewritePlans> {
73 let schemas = parquet_metadata
74 .iter()
75 .map(|metadata| {
76 let file = metadata.file_metadata();
77 let schema = parquet_to_arrow_schema(file.schema_descr(), file.key_value_metadata())
78 .map_err(|error| {
79 InvalidRecordBatchSnafu {
80 reason: format!("Failed to read compaction input Arrow schema: {error}"),
81 }
82 .build()
83 })?;
84 let rows = metadata
85 .row_groups()
86 .iter()
87 .map(|x| x.num_rows())
88 .sum::<i64>() as u64;
89 Ok((Arc::new(schema), rows))
90 })
91 .collect::<Result<Vec<_>>>()?;
92
93 collect_json2_rewrite_plans(metadata, &schemas)
94}
95
96pub(crate) fn collect_json2_rewrite_plans(
101 metadata: &RegionMetadataRef,
102 schemas: &[(SchemaRef, u64)],
103) -> Result<Json2RewritePlans> {
104 let json2_columns = metadata
105 .column_metadatas
106 .iter()
107 .filter_map(|x| {
108 x.column_schema
109 .data_type
110 .is_json2()
111 .then_some(&x.column_schema)
112 })
113 .collect::<Vec<_>>();
114
115 let mut plans = HashMap::with_capacity(json2_columns.len());
116 for column in json2_columns {
117 let extension = column
118 .extension_type::<Json2ExtensionType>()
119 .context(DataTypeMismatchSnafu)?
120 .with_context(|| InvalidRecordBatchSnafu {
121 reason: format!("JSON2 column '{}' has no extension metadata", column.name),
122 })?;
123 ensure!(
127 extension.metadata().is_version_2(),
128 InvalidRecordBatchSnafu {
129 reason: format!("JSON2 column '{}' is not layout v2", column.name),
130 }
131 );
132
133 let settings = extension.metadata().json_settings();
134 let hint_paths = settings
135 .type_hints()
136 .iter()
137 .map(|hint| hint.path.iter().map(String::as_str).collect::<Vec<_>>())
138 .collect::<HashSet<_>>();
139 let mut stats = HashMap::new();
140 for (schema, rows) in schemas {
141 let Some((_, field)) = schema.fields().find(&column.name) else {
142 continue;
143 };
144 collect_json2_path_stats(field, *rows, &hint_paths, &mut stats)?;
145 }
146
147 let mut hints = settings.type_hints().to_vec();
148 hints.extend(select_dynamic_hints(settings, &hint_paths, &stats));
149 let target_layout = JsonSettings::try_new(hints, Some(0)).context(DataTypeMismatchSnafu)?;
150 plans.insert(
151 column.name.clone(),
152 Json2RewritePlan {
153 logical_settings: settings.clone(),
154 target_layout,
155 },
156 );
157 }
158 Ok(plans)
159}
160
161fn collect_json2_path_stats<'a>(
162 field: &'a Field,
163 rows: u64,
164 hint_paths: &HashSet<Vec<&str>>,
165 stats: &mut HashMap<Vec<&'a str>, Json2LeafPathStats>,
166) -> Result<()> {
167 let ArrowDataType::Struct(fields) = field.data_type() else {
168 return InvalidRecordBatchSnafu {
169 reason: format!("JSON2 column '{}' is not a struct", field.name()),
170 }
171 .fail();
172 };
173 let mut paths = Vec::new();
174 for field in fields {
175 if field.name() == JSON2_REMAINDER_FIELD_NAME {
176 continue;
177 }
178 collect_leaf_path_types(field, &mut Vec::new(), &mut paths)?;
179 }
180
181 for (path, data_type) in paths {
182 if hint_paths.contains(path.as_slice()) {
183 continue;
184 }
185 let Some(stat) = stats.get_mut(&path) else {
186 stats.insert(
187 path,
188 Json2LeafPathStats {
189 rows,
190 data_type,
191 is_type_conflicted: false,
192 },
193 );
194 continue;
195 };
196 if stat.data_type != data_type {
197 stat.is_type_conflicted = true;
198 } else {
199 stat.rows += rows;
200 }
201 }
202 Ok(())
203}
204
205fn collect_leaf_path_types<'a>(
206 field: &'a Field,
207 path: &mut Vec<&'a str>,
208 paths: &mut Vec<(Vec<&'a str>, JsonNativeType)>,
209) -> Result<()> {
210 path.push(field.name());
211 if let ArrowDataType::Struct(fields) = field.data_type()
212 && !fields.is_empty()
213 {
214 for field in fields {
215 collect_leaf_path_types(field, path, paths)?;
216 }
217 } else {
218 let json_type =
219 JsonNativeType::try_from(field.data_type()).context(DataTypeMismatchSnafu)?;
220 paths.push((path.clone(), json_type));
221 }
222 path.pop();
223 Ok(())
224}
225
226fn select_dynamic_hints(
227 settings: &JsonSettings,
228 hint_paths: &HashSet<Vec<&str>>,
229 stats: &HashMap<Vec<&str>, Json2LeafPathStats>,
230) -> Vec<JsonTypeHint> {
231 let all_paths = stats
232 .keys()
233 .map(Vec::as_slice)
234 .chain(hint_paths.iter().map(Vec::as_slice))
235 .collect::<HashSet<_>>();
236 let has_ancestor_path =
237 |path: &[&str]| (1..path.len()).any(|len| all_paths.contains(&path[..len]));
238
239 let prefixes = all_paths
240 .iter()
241 .copied()
242 .flat_map(|path| (1..path.len()).map(|len| &path[..len]))
243 .collect::<HashSet<_>>();
244 let has_descendant_path = |path: &[&str]| prefixes.contains(path);
245
246 let mut candidates = stats
247 .iter()
248 .filter(|(path, stat)| {
249 !stat.is_type_conflicted
250 && stat.data_type.is_primitive()
254 && !has_ancestor_path(path)
255 && !has_descendant_path(path)
256 })
257 .collect::<Vec<_>>();
258 candidates.sort_unstable_by(|(x_path, x), (y_path, y)| {
259 y.rows.cmp(&x.rows).then_with(|| x_path.cmp(y_path))
260 });
261 candidates
262 .into_iter()
263 .take(
264 settings
265 .max_auto_expanded_paths()
266 .unwrap_or(JSON2_DEFAULT_MAX_AUTO_EXPANDED_PATHS) as usize,
267 )
268 .map(|(path, stat)| JsonTypeHint {
269 path: path.iter().map(|x| (*x).to_owned()).collect(),
270 data_type: ConcreteDataType::from_arrow_type(&stat.data_type.as_arrow_type()),
271 nullable: true,
272 default_constraint: None,
273 inverted_index: false,
274 })
275 .collect()
276}
277
278pub(crate) fn rewrite_json2_schema(schema: &SchemaRef, plans: &Json2RewritePlans) -> SchemaRef {
280 if plans.is_empty() {
281 return schema.clone();
282 }
283 let fields = schema
284 .fields()
285 .iter()
286 .map(|field| {
287 let Some(plan) = plans.get(field.name()) else {
288 return field.clone();
289 };
290 let mut field = Field::clone(field);
291 field.set_data_type(json2_physical_data_type(&plan.target_layout));
292 field = field.with_extension_type(Json2ExtensionType::new(Arc::new(
293 JsonMetadata::new(plan.logical_settings.clone()),
294 )));
295 Arc::new(field)
296 })
297 .collect::<Vec<_>>();
298 Arc::new(Schema::new_with_metadata(fields, schema.metadata().clone()))
299}
300
301pub(crate) fn rewrite_json2_batch(
303 batch: RecordBatch,
304 plans: &Json2RewritePlans,
305) -> Result<RecordBatch> {
306 if plans.is_empty() {
307 return Ok(batch);
308 }
309 let mut fields = Vec::with_capacity(batch.num_columns());
310 let mut columns = Vec::with_capacity(batch.num_columns());
311
312 for (field, array) in batch.schema_ref().fields().iter().zip(batch.columns()) {
313 let Some(plan) = plans.get(field.name()) else {
314 fields.push(field.clone());
315 columns.push(array.clone());
316 continue;
317 };
318
319 let array = JsonArray::from(array)
320 .rewrite_to_v2(field, &plan.logical_settings, &plan.target_layout)
321 .context(ConvertValueSnafu)?;
322 debug_assert_eq!(
323 &json2_physical_data_type(&plan.target_layout),
324 array.data_type()
325 );
326
327 let mut field = Field::clone(field);
328 field.set_data_type(array.data_type().clone());
329 field = field.with_extension_type(Json2ExtensionType::new(Arc::new(JsonMetadata::new(
330 plan.logical_settings.clone(),
331 ))));
332 fields.push(Arc::new(field));
333 columns.push(array);
334 }
335
336 let schema = Arc::new(Schema::new_with_metadata(
337 fields,
338 batch.schema_ref().metadata().clone(),
339 ));
340 RecordBatch::try_new(schema, columns).context(NewRecordBatchSnafu)
341}
342
343#[cfg(test)]
344mod tests {
345 use datatypes::extension::json::{
346 JSON2_REMAINDER_FIELD_NAME, Json2PhysicalLayout, JsonMetadata,
347 };
348 use datatypes::json::JsonTypeHint;
349 use datatypes::prelude::{ConcreteDataType, DataType};
350 use datatypes::schema::ColumnSchema;
351 use datatypes::types::json_type::{JsonNativeType, JsonObjectType};
352 use serde_json::json;
353
354 use super::*;
355
356 #[test]
357 fn test_select_dynamic_hints_rejects_type_and_prefix_conflicts()
358 -> Result<(), Box<dyn std::error::Error>> {
359 let settings = JsonSettings::try_new(
360 vec![JsonTypeHint {
361 path: vec!["hint".to_string()],
362 data_type: ConcreteDataType::string_datatype(),
363 nullable: true,
364 default_constraint: None,
365 inverted_index: false,
366 }],
367 Some(2),
368 )?;
369 let stat = |rows, data_type, is_type_conflicted| Json2LeafPathStats {
370 rows,
371 data_type,
372 is_type_conflicted,
373 };
374 let stats = HashMap::from([
375 (
376 vec!["hint", "nested"],
377 stat(10, JsonNativeType::String, false),
378 ),
379 (vec!["popular"], stat(9, JsonNativeType::String, false)),
380 (
381 vec!["popular", "nested"],
382 stat(8, JsonNativeType::u64(), false),
383 ),
384 (
385 vec!["type_conflicted"],
386 stat(7, JsonNativeType::String, true),
387 ),
388 (
389 vec!["array"],
390 stat(
391 6,
392 JsonNativeType::Array(Box::new(JsonNativeType::String)),
393 false,
394 ),
395 ),
396 (vec!["variant"], stat(5, JsonNativeType::Variant, false)),
397 (vec!["tie_a"], stat(2, JsonNativeType::String, false)),
398 (vec!["tie_b"], stat(2, JsonNativeType::Bool, false)),
399 ]);
400
401 let hint_paths = settings
402 .type_hints()
403 .iter()
404 .map(|hint| hint.path.iter().map(String::as_str).collect::<Vec<_>>())
405 .collect::<HashSet<_>>();
406 let hints = select_dynamic_hints(&settings, &hint_paths, &stats);
407 assert_eq!(
408 vec![vec!["tie_a".to_string()], vec!["tie_b".to_string()]],
409 hints.into_iter().map(|x| x.path).collect::<Vec<_>>()
410 );
411 Ok(())
412 }
413
414 #[test]
415 fn test_rewrite_json2_v1_batch_to_target_layout() -> Result<(), Box<dyn std::error::Error>> {
416 let settings = JsonSettings::try_new(
417 vec![JsonTypeHint {
418 path: vec!["kind".to_string()],
419 data_type: ConcreteDataType::string_datatype(),
420 nullable: true,
421 default_constraint: None,
422 inverted_index: false,
423 }],
424 Some(0),
425 )?;
426 let target = json2_physical_data_type(&settings);
427 let plans = HashMap::from([(
428 "j".to_string(),
429 Json2RewritePlan {
430 logical_settings: settings.clone(),
431 target_layout: settings,
432 },
433 )]);
434 let values = [
435 json!({"kind": "a", "extra": {"x": 1}}),
436 json!({"kind": "b", "extra": {"x": 2}}),
437 ];
438 let source_settings = JsonSettings::default();
439 let source_extension =
440 Json2ExtensionType::new(Arc::new(JsonMetadata::new_v1(source_settings.clone())));
441 let mut source_column = ColumnSchema::new(
442 "j",
443 ConcreteDataType::json2(JsonNativeType::Object(JsonObjectType::new())),
444 true,
445 );
446 source_column.with_extension_type(&source_extension);
447 let mut source_builder = source_column.data_type.create_mutable_vector(values.len());
448 for value in &values {
449 let value = source_settings.encode(value.clone())?;
450 source_builder.try_push_value_ref(&value.as_value_ref())?;
451 }
452 let source = source_builder.to_vector().to_arrow_array();
453 let field =
454 Field::new("j", source.data_type().clone(), true).with_extension_type(source_extension);
455 let batch = RecordBatch::try_new(Arc::new(Schema::new(vec![field])), vec![source])?;
456
457 let batch = rewrite_json2_batch(batch, &plans)?;
458 let field = batch.schema_ref().field(0);
459 assert!(Json2PhysicalLayout::try_from_root(field)?.is_version_2());
460 assert_eq!(&target, field.data_type());
461 let projected =
462 JsonArray::from(batch.column(0)).project_to_v2(field, &ArrowDataType::Binary)?;
463 let projected = JsonArray::from(&projected);
464 for (i, expected) in values.into_iter().enumerate() {
465 assert_eq!(expected, projected.try_get_value(i)?);
466 }
467 Ok(())
468 }
469
470 #[test]
471 fn test_rewrite_json2_v2_source_to_narrower_target() -> Result<(), Box<dyn std::error::Error>> {
472 let target_settings = JsonSettings::try_new(
473 vec![JsonTypeHint {
474 path: vec!["kind".to_string()],
475 data_type: ConcreteDataType::string_datatype(),
476 nullable: true,
477 default_constraint: None,
478 inverted_index: false,
479 }],
480 Some(0),
481 )?;
482 let target_extension =
483 Json2ExtensionType::new(Arc::new(JsonMetadata::new(target_settings.clone())));
484 let target_type = json2_physical_data_type(&target_settings);
485 let plans = HashMap::from([(
486 "j".to_string(),
487 Json2RewritePlan {
488 logical_settings: target_settings.clone(),
489 target_layout: target_settings,
490 },
491 )]);
492
493 let source_settings = JsonSettings::try_new(
494 vec![
495 JsonTypeHint {
496 path: vec!["kind".to_string()],
497 data_type: ConcreteDataType::string_datatype(),
498 nullable: true,
499 default_constraint: None,
500 inverted_index: false,
501 },
502 JsonTypeHint {
503 path: vec!["source_only".to_string()],
504 data_type: ConcreteDataType::int64_datatype(),
505 nullable: true,
506 default_constraint: None,
507 inverted_index: false,
508 },
509 ],
510 None,
511 )?;
512 let source_extension =
513 Json2ExtensionType::new(Arc::new(JsonMetadata::new(source_settings.clone())));
514 let mut source_column = ColumnSchema::new(
515 "j",
516 ConcreteDataType::json2(JsonNativeType::Object(JsonObjectType::new())),
517 true,
518 );
519 source_column.with_extension_type(&source_extension);
520 let expected = json!({
521 "kind": "a",
522 "source_only": 7,
523 "dynamic": {"nested": true}
524 });
525 let mut builder = source_column.create_mutable_vector(1);
526 let value = source_settings.encode(expected.clone())?;
527 builder.try_push_value_ref(&value.as_value_ref())?;
528 let source = builder.to_vector().to_arrow_array();
529 let source_field =
530 Field::new("j", source.data_type().clone(), true).with_extension_type(source_extension);
531 let source =
532 JsonArray::from(&source).project_to_v2(&source_field, &ArrowDataType::Binary)?;
533 let field = Field::new("j", ArrowDataType::Binary, true)
534 .with_extension_type(target_extension.clone());
535 let batch = RecordBatch::try_new(Arc::new(Schema::new(vec![field])), vec![source])?;
536
537 let batch = rewrite_json2_batch(batch, &plans)?;
538 let field = batch.schema_ref().field(0);
539 assert!(Json2PhysicalLayout::try_from_root(field)?.is_version_2());
540 assert_eq!(&target_type, field.data_type());
541 let ArrowDataType::Struct(fields) = field.data_type() else {
542 unreachable!()
543 };
544 assert_eq!(
545 vec![JSON2_REMAINDER_FIELD_NAME, "kind"],
546 fields.iter().map(|x| x.name().as_str()).collect::<Vec<_>>()
547 );
548
549 let projected =
550 JsonArray::from(batch.column(0)).project_to_v2(field, &ArrowDataType::Binary)?;
551 assert_eq!(expected, JsonArray::from(&projected).try_get_value(0)?);
552
553 let first_schema = batch.schema();
554 let field =
555 Field::new("j", ArrowDataType::Binary, true).with_extension_type(target_extension);
556 let batch = RecordBatch::try_new(Arc::new(Schema::new(vec![field])), vec![projected])?;
557 let batch = rewrite_json2_batch(batch, &plans)?;
558 assert_eq!(first_schema, batch.schema());
559
560 let field = batch.schema_ref().field(0);
561 let projected =
562 JsonArray::from(batch.column(0)).project_to_v2(field, &ArrowDataType::Binary)?;
563 assert_eq!(expected, JsonArray::from(&projected).try_get_value(0)?);
564 Ok(())
565 }
566
567 #[test]
568 fn test_rewrite_json2_v2_source_to_wider_target() -> Result<(), Box<dyn std::error::Error>> {
569 let logical_settings = JsonSettings::try_new(
570 vec![JsonTypeHint {
571 path: vec!["kind".to_string()],
572 data_type: ConcreteDataType::string_datatype(),
573 nullable: true,
574 default_constraint: None,
575 inverted_index: false,
576 }],
577 Some(0),
578 )?;
579 let target_layout = JsonSettings::try_new(
580 vec![
581 JsonTypeHint {
582 path: vec!["kind".to_string()],
583 data_type: ConcreteDataType::string_datatype(),
584 nullable: true,
585 default_constraint: None,
586 inverted_index: false,
587 },
588 JsonTypeHint {
589 path: vec!["promoted".to_string()],
590 data_type: ConcreteDataType::int64_datatype(),
591 nullable: true,
592 default_constraint: None,
593 inverted_index: false,
594 },
595 ],
596 Some(0),
597 )?;
598 let target_type = json2_physical_data_type(&target_layout);
599 let plans = HashMap::from([(
600 "j".to_string(),
601 Json2RewritePlan {
602 logical_settings: logical_settings.clone(),
603 target_layout,
604 },
605 )]);
606
607 let extension =
608 Json2ExtensionType::new(Arc::new(JsonMetadata::new(logical_settings.clone())));
609 let mut column = ColumnSchema::new(
610 "j",
611 ConcreteDataType::json2(JsonNativeType::Object(JsonObjectType::new())),
612 true,
613 );
614 column.with_extension_type(&extension);
615 let expected = json!({
616 "kind": "a",
617 "promoted": 7,
618 "dynamic": {"nested": true}
619 });
620 let mut builder = column.create_mutable_vector(1);
621 let value = logical_settings.encode(expected.clone())?;
622 builder.try_push_value_ref(&value.as_value_ref())?;
623 let source = builder.to_vector().to_arrow_array();
624 assert_eq!(
625 &json2_physical_data_type(&logical_settings),
626 source.data_type()
627 );
628 let field =
629 Field::new("j", source.data_type().clone(), true).with_extension_type(extension);
630 let batch = RecordBatch::try_new(Arc::new(Schema::new(vec![field])), vec![source])?;
631
632 let batch = rewrite_json2_batch(batch, &plans)?;
633 let field = batch.schema_ref().field(0);
634 assert_eq!(&target_type, field.data_type());
635 let ArrowDataType::Struct(fields) = field.data_type() else {
636 unreachable!()
637 };
638 assert_eq!(
639 vec![JSON2_REMAINDER_FIELD_NAME, "kind", "promoted"],
640 fields.iter().map(|x| x.name().as_str()).collect::<Vec<_>>()
641 );
642
643 let array = batch
644 .column(0)
645 .as_any()
646 .downcast_ref::<datatypes::arrow::array::StructArray>()
647 .unwrap();
648 let promoted = array
649 .column_by_name("promoted")
650 .unwrap()
651 .as_any()
652 .downcast_ref::<datatypes::arrow::array::Int64Array>()
653 .unwrap();
654 assert_eq!(7, promoted.value(0));
655
656 let projected =
657 JsonArray::from(batch.column(0)).project_to_v2(field, &ArrowDataType::Binary)?;
658 assert_eq!(expected, JsonArray::from(&projected).try_get_value(0)?);
659 Ok(())
660 }
661}