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