From 1c01b541ec617e580d4c8a563bccf56886e39957 Mon Sep 17 00:00:00 2001 From: fys <40801205+fengys1996@users.noreply.github.com> Date: Wed, 16 Sep 2026 04:44:32 +0000 Subject: [PATCH] refactor(json2): JSON2 parquet projection and schema alignment (#9137) * refactor(mito2): make JSON schema alignment targets explicit Signed-off-by: fys * fix: cr * fix: cargo fmt --------- Signed-off-by: fys --- src/datatypes/src/vectors/json/array.rs | 451 +------------- .../{json_align/stream.rs => json_align.rs} | 505 ++++++++++++---- src/mito2/src/sst/parquet/json_align/mod.rs | 24 - src/mito2/src/sst/parquet/read_columns.rs | 568 +++++++----------- src/mito2/src/sst/parquet/reader.rs | 39 +- 5 files changed, 639 insertions(+), 948 deletions(-) rename src/mito2/src/sst/parquet/{json_align/stream.rs => json_align.rs} (53%) delete mode 100644 src/mito2/src/sst/parquet/json_align/mod.rs diff --git a/src/datatypes/src/vectors/json/array.rs b/src/datatypes/src/vectors/json/array.rs index 153d6b490ef..e8b75dee259 100644 --- a/src/datatypes/src/vectors/json/array.rs +++ b/src/datatypes/src/vectors/json/array.rs @@ -12,7 +12,6 @@ // See the License for the specific language governing permissions and // limitations under the License. -use std::cmp::Ordering; use std::sync::Arc; use arrow::compute::{can_cast_types, cast}; @@ -21,13 +20,12 @@ use arrow_array::types::{ Float32Type, Float64Type, Int8Type, Int16Type, Int32Type, Int64Type, UInt8Type, UInt16Type, UInt32Type, UInt64Type, }; -use arrow_array::{Array, ArrayRef, GenericListArray, ListArray, StructArray, new_null_array}; -use arrow_schema::{DataType, Field, FieldRef}; -use common_telemetry::trace; +use arrow_array::{Array, ArrayRef, GenericListArray, StructArray, new_null_array}; +use arrow_schema::{DataType, Field}; use serde_json::Value; use snafu::{OptionExt, ResultExt}; -use crate::arrow_array::{MutableBinaryArray, binary_array_value, string_array_value}; +use crate::arrow_array::{binary_array_value, string_array_value}; use crate::data_type::ConcreteDataType; use crate::error::{ AlignJsonArraySnafu, ArrowComputeSnafu, InvalidJsonSnafu, InvalidJsonbSnafu, Result, @@ -200,171 +198,6 @@ impl JsonArray<'_> { Ok(values) } - /// Normalizes a JSON2 array to the wider `expect` data type without losing - /// information. - /// - /// This is mainly used for write/flush-time JSON2 schema alignment: - /// - fields missing from the source are filled with typed null arrays; - /// - fields present in the source must also exist in `expect`; - /// - fields present in both are widened recursively when their types differ. - /// - /// Narrowing conversions and any other conversions that may lose information - /// are rejected. - pub fn widen_to(&self, expect: &DataType) -> Result { - let data_type = self.inner.data_type(); - - if data_type == expect { - return Ok(self.inner.clone()); - } - - trace!( - "Try aligning JSON array {} to data type {}", - data_type, expect - ); - - let struct_array = self.inner.as_struct_opt().context(AlignJsonArraySnafu { - reason: "expect struct array", - })?; - let array_fields = struct_array.fields(); - let array_columns = struct_array.columns(); - let DataType::Struct(expect_fields) = expect else { - return AlignJsonArraySnafu { - reason: "expect struct datatype", - } - .fail(); - }; - let mut aligned = Vec::with_capacity(expect_fields.len()); - - // Compare the fields in the JSON array and the to-be-aligned schema, amending with null - // arrays on the way. It's very important to note that fields in the JSON array and those - // in the JSON type are both **SORTED**, which can be guaranteed because the fields in the - // JSON type implementation are sorted. - debug_assert!(expect_fields.iter().map(|f| f.name()).is_sorted()); - debug_assert!(array_fields.iter().map(|f| f.name()).is_sorted()); - - let mut i = 0; // point to the expect fields - let mut j = 0; // point to the array fields - while i < expect_fields.len() && j < array_fields.len() { - let expect_field = &expect_fields[i]; - let array_field = &array_fields[j]; - match expect_field.name().cmp(array_field.name()) { - Ordering::Equal => { - if expect_field.data_type() == array_field.data_type() { - aligned.push(array_columns[j].clone()); - } else { - let expect_type = expect_field.data_type(); - let array_type = array_field.data_type(); - let array = match (expect_type, array_type) { - (DataType::Struct(_), DataType::Struct(_)) => { - JsonArray::from(&array_columns[j]).widen_to(expect_type)? - } - (DataType::List(expect_item), DataType::List(array_item)) => { - let list_array = array_columns[j].as_list::(); - widen_list(list_array, array_item, expect_item)? - } - _ => JsonArray::from(&array_columns[j]).widen_scalar_to(expect_type)?, - }; - aligned.push(array); - } - i += 1; - j += 1; - } - Ordering::Less => { - aligned.push(new_null_array(expect_field.data_type(), struct_array.len())); - i += 1; - } - Ordering::Greater => { - return AlignJsonArraySnafu { - reason: format!( - "source field {} does not exist in target schema", - array_field.name() - ), - } - .fail(); - } - } - } - if j < array_fields.len() { - return AlignJsonArraySnafu { - reason: format!( - "source field {} does not exist in target schema", - array_fields[j].name() - ), - } - .fail(); - } - if i < expect_fields.len() { - for field in &expect_fields[i..] { - aligned.push(new_null_array(field.data_type(), struct_array.len())); - } - } - - let json_array = StructArray::try_new_with_length( - expect_fields.clone(), - aligned, - struct_array.nulls().cloned(), - struct_array.len(), - ) - .map_err(|e| { - AlignJsonArraySnafu { - reason: e.to_string(), - } - .build() - })?; - Ok(Arc::new(json_array)) - } - - /// Widens an array to the merged JSON2 physical type without losing information. - /// - /// Supported conversions: - /// - identical types are returned unchanged; - /// - null arrays become typed null arrays; - /// - concrete JSON values are encoded as JSONB when the target type is binary. - /// - /// All other conversions are rejected. - fn widen_scalar_to(&self, to_type: &DataType) -> Result { - let from_type = self.inner.data_type(); - if from_type == to_type { - return Ok(self.inner.clone()); - } - - if from_type == &DataType::Null { - return Ok(new_null_array(to_type, self.inner.len())); - } - - if !from_type.is_binary() && to_type.is_binary() { - return self.encode_variant(); - } - - AlignJsonArraySnafu { - reason: format!("unable to widen {from_type} to {to_type}"), - } - .fail() - } - - fn encode_variant(&self) -> Result { - let len = self.inner.len(); - let mut encoded = Vec::with_capacity(len); - let mut total_bytes = 0; - - for i in 0..len { - let value = self.try_get_value(i)?; - if value.is_null() { - encoded.push(None); - } else { - let bytes = encode_serde_json_as_jsonb(value); - total_bytes += bytes.len(); - encoded.push(Some(bytes)); - } - } - - let mut builder = MutableBinaryArray::with_capacity(len, total_bytes); - for value in encoded { - builder.append_option(value); - } - Ok(Arc::new(builder.finish())) - } - /// Projects this JSON array to `target` for query evaluation. /// /// Unlike [`Self::widen_to`], projection tolerates lossy conversions: @@ -567,28 +400,6 @@ fn project_json_value_to_type(value: Value, to_type: &ConcreteDataType) -> Resul Ok(to_type.try_cast(value).unwrap_or(GreptimeValue::Null)) } -fn widen_list(list_array: &ListArray, actual: &FieldRef, expected: &FieldRef) -> Result { - let item_aligned = match (actual.data_type(), expected.data_type()) { - (DataType::Struct(_), DataType::Struct(_)) => { - JsonArray::from(list_array.values()).widen_to(expected.data_type())? - } - (DataType::List(actual), DataType::List(expected)) => { - let list_array = list_array.values().as_list::(); - widen_list(list_array, actual, expected)? - } - _ => JsonArray::from(list_array.values()).widen_scalar_to(expected.data_type())?, - }; - Ok(Arc::new( - GenericListArray::::try_new( - expected.clone(), - list_array.offsets().clone(), - item_aligned, - list_array.nulls().cloned(), - ) - .context(ArrowComputeSnafu)?, - )) -} - impl<'a> From<&'a ArrayRef> for JsonArray<'a> { fn from(inner: &'a ArrayRef) -> Self { Self { inner } @@ -740,262 +551,6 @@ mod test { Ok(()) } - #[test] - fn test_widen_null_to_any_type() -> Result<()> { - let nulls = new_null_array(&DataType::Null, 2); - let target_types = [ - DataType::Boolean, - DataType::UInt64, - DataType::Utf8View, - DataType::Binary, - DataType::List(Arc::new(Field::new_list_field(DataType::Int64, true))), - DataType::Struct(Fields::from(vec![Field::new( - "value", - DataType::Int64, - true, - )])), - ]; - - for target_type in target_types { - let widened = JsonArray::from(&nulls).widen_scalar_to(&target_type)?; - assert_eq!(&target_type, widened.data_type()); - assert_eq!(2, widened.len()); - assert_eq!(2, widened.null_count()); - } - - Ok(()) - } - - #[test] - fn test_widen_non_null_to_utf8_view_fails() { - let bools: ArrayRef = Arc::new(BooleanArray::from(vec![true])); - let err = JsonArray::from(&bools) - .widen_scalar_to(&DataType::Utf8View) - .unwrap_err(); - - assert_eq!( - "Failed to align JSON array, reason: unable to widen Boolean to Utf8View", - err.to_string() - ); - } - - #[test] - fn test_widen_variant_to_non_binary_fails() { - let value = jsonb::parse_value(b"true").unwrap().to_vec(); - let variants: ArrayRef = Arc::new(BinaryArray::from(vec![value.as_slice()])); - let err = JsonArray::from(&variants) - .widen_scalar_to(&DataType::Boolean) - .unwrap_err(); - - assert_eq!( - "Failed to align JSON array, reason: unable to widen Binary to Boolean", - err.to_string() - ); - } - - #[test] - fn test_widen_between_number_types_fails() { - let values: ArrayRef = Arc::new(UInt64Array::from(vec![1])); - let err = JsonArray::from(&values) - .widen_scalar_to(&DataType::Int64) - .unwrap_err(); - - assert_eq!( - "Failed to align JSON array, reason: unable to widen UInt64 to Int64", - err.to_string() - ); - } - - #[test] - fn test_widen_numbers_to_variant_preserves_values() -> Result<()> { - let cases: [(ArrayRef, Value); 3] = [ - (Arc::new(UInt64Array::from(vec![u64::MAX])), json!(u64::MAX)), - (Arc::new(Int64Array::from(vec![i64::MIN])), json!(i64::MIN)), - (Arc::new(Float64Array::from(vec![1.25])), json!(1.25)), - ]; - - for (values, expected) in cases { - let widened = JsonArray::from(&values).widen_scalar_to(&DataType::Binary)?; - assert_eq!(&DataType::Binary, widened.data_type()); - assert_eq!(expected, JsonArray::from(&widened).try_get_value(0)?); - } - - Ok(()) - } - - #[test] - fn test_align_json_array() -> Result<()> { - struct TestCase { - json_array: ArrayRef, - schema_type: DataType, - expected: std::result::Result, - } - - impl TestCase { - fn new( - json_array: StructArray, - schema_type: Fields, - expected: std::result::Result, String>, - ) -> Self { - Self { - json_array: Arc::new(json_array), - schema_type: DataType::Struct(schema_type.clone()), - expected: expected - .map(|x| Arc::new(StructArray::new(schema_type, x, None)) as ArrayRef), - } - } - - fn test(self) -> Result<()> { - let result = JsonArray::from(&self.json_array).widen_to(&self.schema_type); - match (result, self.expected) { - (Ok(json_array), Ok(expected)) => assert_eq!(&json_array, &expected), - (Ok(json_array), Err(e)) => { - panic!("expecting error {e} but actually get: {json_array:?}") - } - (Err(e), Err(expected)) => assert_eq!(e.to_string(), expected), - (Err(e), Ok(_)) => return Err(e), - } - Ok(()) - } - } - - // Test empty json array can be aligned with a complex json type. - TestCase::new( - StructArray::new_empty_fields(2, None), - Fields::from(vec![ - Field::new("int", DataType::Int64, true), - Field::new_struct( - "nested", - vec![Field::new("bool", DataType::Boolean, true)], - true, - ), - Field::new("string", DataType::Utf8, true), - ]), - Ok(vec![ - Arc::new(Int64Array::new_null(2)) as ArrayRef, - Arc::new(StructArray::new_null( - Fields::from(vec![Arc::new(Field::new("bool", DataType::Boolean, true))]), - 2, - )), - Arc::new(StringArray::new_null(2)), - ]), - ) - .test()?; - - // Test simple json array alignment. - TestCase::new( - StructArray::from(vec![( - Arc::new(Field::new("float", DataType::Float64, true)), - Arc::new(Float64Array::from(vec![1.0, 2.0, 3.0])) as ArrayRef, - )]), - Fields::from(vec![ - Field::new("float", DataType::Float64, true), - Field::new("string", DataType::Utf8, true), - ]), - Ok(vec![ - Arc::new(Float64Array::from(vec![1.0, 2.0, 3.0])) as ArrayRef, - Arc::new(StringArray::new_null(3)), - ]), - ) - .test()?; - - // Test complex json array alignment. - TestCase::new( - StructArray::from(vec![ - ( - Arc::new(Field::new_list( - "list", - Field::new_list_field(DataType::Int64, true), - true, - )), - Arc::new(ListArray::from_iter_primitive::(vec![ - Some(vec![Some(1)]), - None, - Some(vec![Some(2), Some(3)]), - ])) as ArrayRef, - ), - ( - Arc::new(Field::new_struct( - "nested", - vec![Field::new("int", DataType::Int64, true)], - true, - )), - Arc::new(StructArray::from(vec![( - Arc::new(Field::new("int", DataType::Int64, true)), - Arc::new(Int64Array::from(vec![-1, -2, -3])) as ArrayRef, - )])), - ), - ( - Arc::new(Field::new("string", DataType::Utf8, true)), - Arc::new(StringArray::from(vec!["a", "b", "c"])), - ), - ]), - Fields::from(vec![ - Field::new("bool", DataType::Boolean, true), - Field::new_list("list", Field::new_list_field(DataType::Int64, true), true), - Field::new_struct( - "nested", - vec![ - Field::new("float", DataType::Float64, true), - Field::new("int", DataType::Int64, true), - ], - true, - ), - Field::new("string", DataType::Utf8, true), - ]), - Ok(vec![ - Arc::new(BooleanArray::new_null(3)) as ArrayRef, - Arc::new(ListArray::from_iter_primitive::(vec![ - Some(vec![Some(1)]), - None, - Some(vec![Some(2), Some(3)]), - ])), - Arc::new(StructArray::from(vec![ - ( - Arc::new(Field::new("float", DataType::Float64, true)), - Arc::new(Float64Array::new_null(3)) as ArrayRef, - ), - ( - Arc::new(Field::new("int", DataType::Int64, true)), - Arc::new(Int64Array::from(vec![-1, -2, -3])), - ), - ])), - Arc::new(StringArray::from(vec!["a", "b", "c"])), - ]), - ) - .test()?; - - // Source fields that do not exist in the target schema must not be discarded. - TestCase::new( - StructArray::from(vec![( - Arc::new(Field::new("a", DataType::Boolean, true)), - Arc::new(BooleanArray::from(vec![true])) as ArrayRef, - )]), - Fields::from(vec![Field::new("b", DataType::Boolean, true)]), - Err( - "Failed to align JSON array, reason: source field a does not exist in target schema" - .to_string(), - ), - ) - .test()?; - - // Trailing source fields must also be rejected after all target fields are processed. - TestCase::new( - StructArray::from(vec![( - Arc::new(Field::new("b", DataType::Boolean, true)), - Arc::new(BooleanArray::from(vec![true])) as ArrayRef, - )]), - Fields::from(vec![Field::new("a", DataType::Boolean, true)]), - Err( - "Failed to align JSON array, reason: source field b does not exist in target schema" - .to_string(), - ), - ) - .test()?; - - Ok(()) - } - #[test] fn test_align_variant_to_struct() -> Result<()> { let encode = |json: &[u8]| jsonb::parse_value(json).unwrap().to_vec(); diff --git a/src/mito2/src/sst/parquet/json_align/stream.rs b/src/mito2/src/sst/parquet/json_align.rs similarity index 53% rename from src/mito2/src/sst/parquet/json_align/stream.rs rename to src/mito2/src/sst/parquet/json_align.rs index a474fd781d0..8af88f6ac0f 100644 --- a/src/mito2/src/sst/parquet/json_align/stream.rs +++ b/src/mito2/src/sst/parquet/json_align.rs @@ -19,12 +19,15 @@ use std::task::{Context, Poll}; use datafusion_common::cast_column; use datafusion_common::format::DEFAULT_CAST_OPTIONS; use datatypes::arrow::array::{ArrayRef, new_null_array}; -use datatypes::arrow::datatypes::{DataType, Field, FieldRef, SchemaRef}; +use datatypes::arrow::datatypes::{DataType, Field, FieldRef, Schema, SchemaRef}; use datatypes::arrow::record_batch::RecordBatch; use datatypes::extension::json::{JsonMetadata, is_json2_extension_type}; use datatypes::json::JsonSettings; use datatypes::vectors::json::array::JsonArray; +use datatypes::vectors::json::json2_physical_data_type; use futures::Stream; +use futures::stream::BoxStream; +use serde_json::from_str; use snafu::{ResultExt, ensure}; use crate::error::{ @@ -32,32 +35,44 @@ use crate::error::{ }; use crate::sst::parquet::Json2TargetLayout; +pub(crate) type ProjectedRecordBatchStream = BoxStream<'static, Result>; + +/// Specifies how JSON columns in a record batch are aligned. #[derive(Debug)] -struct Json2RewriteSettings { - logical_settings: JsonSettings, - target_layout: JsonSettings, +pub(crate) enum AlignMode { + /// Aligns JSON columns to the logical fields in the output schema. + AlignToSchema, + /// Rewrites JSON columns to physical layouts, typically for compaction. + Rewrite { + /// Target layouts keyed by root column name, not nested field path. + /// + /// Only listed columns are rewritten. Other existing arrays are reused + /// unchanged and must already match their output field types. + /// An empty map therefore only fills missing roots. + columns: HashMap, + }, } -/// Aligns projected batches to the expected output schema for nested projections. +/// Alignment mode with parsed rewrite metadata and validated target layouts. +#[derive(Debug)] +enum ResolvedAlignMode { + AlignToSchema, + Rewrite { + columns: HashMap, + }, +} + +/// Adapts Parquet record batches to the output schema expected by the reader. /// -/// Background -/// ---------- -/// Nested projection may ask parquet to read leaves under a root column. If none -/// of the requested leaves exists in the current parquet file, parquet decoding -/// omits the whole root from the physical [`RecordBatch`]. +/// Nested projection can return only part of a JSON2 column, or omit its root +/// entirely when no requested leaves are read. This stream restores missing +/// roots with null arrays of the expected types. /// -/// In addition, after nested-path filtering, returned struct arrays may contain -/// only a subset of fields. The current output schema is not pruned by nested -/// paths, so physical struct fields can be a subset of the expected struct -/// fields, and their nested schema can differ from the expected output schema. -/// -/// To keep projected batches schema-consistent before entering upper readers: -/// - Root-column presence alignment restores missing projected root columns by -/// inserting root-level null arrays. -/// - Nested struct alignment aligns struct arrays to the expected nested field -/// layout. +/// Existing JSON2 columns are aligned to the logical schema inferred from +/// type hints or rewritten to the specified JSON2 physical layout, as selected by +/// [`AlignMode`]. #[derive(derive_more::Debug)] -pub struct NestedSchemaAligner { +pub struct JsonSchemaAligner { #[debug(skip)] inner: S, /// Output schema expected by the upper reader. @@ -70,79 +85,55 @@ pub struct NestedSchemaAligner { /// Whether all projected roots are present and the stream can pass batches /// through. all_roots_present: bool, - /// JSON2 columns that require semantic source-to-target layout rewriting. - json2_rewrite_targets: HashMap, + /// Alignment mode with parsed and validated rewrite settings. + mode: ResolvedAlignMode, /// The cache for whether incoming batches already match output schema. is_schema_matched: Option, } -impl NestedSchemaAligner +impl JsonSchemaAligner where S: Stream>, { - pub fn new( + /// Creates an aligner with a shared output schema and an explicit operation. + /// Parses rewrite metadata once and validates layouts against output field types. + pub(crate) fn new( inner: S, projected_root_presence: Vec, output_schema: SchemaRef, - ) -> Result> { + mode: AlignMode, + ) -> Result> { ensure!( projected_root_presence.len() == output_schema.fields().len(), UnexpectedSnafu { reason: format!( - "NestedSchemaAligner projected root presence len {} does not match output schema columns {}", + "JsonSchemaAligner projected root presence len {} does not match output schema columns {}", projected_root_presence.len(), output_schema.fields().len() ), } ); + let mode = resolve_align_mode(mode, output_schema.as_ref())?; + let expected_input_col_num = projected_root_presence .iter() .filter(|matched| **matched) .count(); let all_roots_present = projected_root_presence.iter().all(|&m| m); - Ok(NestedSchemaAligner { + Ok(JsonSchemaAligner { inner, output_schema, projected_root_presence, expected_input_col_num, all_roots_present, - json2_rewrite_targets: HashMap::new(), + mode, is_schema_matched: None, }) } - - /// Sets JSON2 columns that must be rewritten into the output field layout. - pub(crate) fn with_json2_rewrite_targets( - mut self, - targets: &HashMap, - ) -> Result { - self.json2_rewrite_targets = targets - .iter() - .map(|(name, layout)| { - let metadata = serde_json::from_str::(&layout.extension_metadata) - .map_err(|e| { - UnexpectedSnafu { - reason: format!( - "invalid JSON2 extension metadata for column '{name}': {e}" - ), - } - .build() - })?; - Ok(( - name.clone(), - Json2RewriteSettings { - logical_settings: metadata.into_json_settings(), - target_layout: layout.target_layout.clone(), - }, - )) - }) - .collect::>()?; - Ok(self) - } } -impl Stream for NestedSchemaAligner +impl Stream for JsonSchemaAligner where S: Stream> + Unpin, { @@ -153,7 +144,8 @@ where match Pin::new(&mut this.inner).poll_next(cx) { Poll::Ready(Some(Ok(rb))) => { - let is_schema_matched = this.all_roots_present + let is_schema_matched = matches!(this.mode, ResolvedAlignMode::AlignToSchema) + && this.all_roots_present && *this .is_schema_matched .get_or_insert_with(|| rb.schema() == this.output_schema); @@ -166,7 +158,7 @@ where &this.output_schema, &this.projected_root_presence, this.expected_input_col_num, - &this.json2_rewrite_targets, + &this.mode, ))) } } @@ -182,13 +174,13 @@ fn align_projected_batch( output_schema: &SchemaRef, projected_root_presence: &[bool], expected_input_col_num: usize, - json2_rewrite_targets: &HashMap, + mode: &ResolvedAlignMode, ) -> Result { ensure!( rb.columns().len() == expected_input_col_num, UnexpectedSnafu { reason: format!( - "NestedSchemaAligner expected {} input columns but got {}", + "JsonSchemaAligner expected {} input columns but got {}", expected_input_col_num, rb.columns().len() ), @@ -205,12 +197,16 @@ fn align_projected_batch( continue; } - cols.push(align_array( - rb.column(idx), - input_schema.field(idx), - field, - json2_rewrite_targets.get(field.name()), - )?); + let array = match mode { + ResolvedAlignMode::AlignToSchema => { + align_array(rb.column(idx), input_schema.field(idx), field)? + } + ResolvedAlignMode::Rewrite { columns } => match columns.get(field.name()) { + Some(settings) => rewrite_array(rb.column(idx), input_schema.field(idx), settings)?, + None => rb.column(idx).clone(), + }, + }; + cols.push(array); idx += 1; } @@ -218,36 +214,111 @@ fn align_projected_batch( } fn align_array( - array: &ArrayRef, - source: &Field, - field: &FieldRef, - rewrite_settings: Option<&Json2RewriteSettings>, + source_array: &ArrayRef, + source_field: &Field, + target_field: &FieldRef, ) -> Result { - if let Some(settings) = rewrite_settings { - return JsonArray::from(array) - .rewrite_to_v2(source, &settings.logical_settings, &settings.target_layout) - .context(DataTypeMismatchSnafu); - } - if array.data_type() == field.data_type() { - return Ok(array.clone()); + if source_array.data_type() == target_field.data_type() { + return Ok(source_array.clone()); } - if is_json2_extension_type(field) { - if is_json2_extension_type(source) { - return JsonArray::from(array) - .project_to_v2(source, field.data_type()) + if is_json2_extension_type(target_field) { + if is_json2_extension_type(source_field) { + return JsonArray::from(source_array) + .project_to_v2(source_field, target_field.data_type()) .context(DataTypeMismatchSnafu); } - return JsonArray::from(array) - .project_to(field.data_type()) + return JsonArray::from(source_array) + .project_to(target_field.data_type()) .context(DataTypeMismatchSnafu); } - if !matches!(field.data_type(), DataType::Struct(_)) { - return Ok(array.clone()); + if !matches!(target_field.data_type(), DataType::Struct(_)) { + return Ok(source_array.clone()); } - cast_column(array, field.data_type(), &DEFAULT_CAST_OPTIONS).context(CastColumnSnafu) + cast_column( + source_array, + target_field.data_type(), + &DEFAULT_CAST_OPTIONS, + ) + .context(CastColumnSnafu) +} + +fn rewrite_array( + source_array: &ArrayRef, + source_field: &Field, + settings: &RewriteSettings, +) -> Result { + JsonArray::from(source_array) + .rewrite_to_v2( + source_field, + &settings.logical_settings, + &settings.target_layout, + ) + .context(DataTypeMismatchSnafu) +} + +/// Resolved settings for rewriting one JSON2 column to a target physical layout. +/// +/// Created from [`Json2TargetLayout`] when resolving the alignment mode, so +/// extension metadata is parsed once and reused across batches. +#[derive(Debug)] +struct RewriteSettings { + /// Logical settings parsed from extension metadata and applied to JSON values + /// before encoding them into the target layout. + logical_settings: JsonSettings, + /// Settings defining the physical Arrow layout of the rewritten column. + target_layout: JsonSettings, +} + +impl TryFrom<&Json2TargetLayout> for RewriteSettings { + type Error = crate::error::Error; + + fn try_from(layout: &Json2TargetLayout) -> Result { + let metadata = from_str::(&layout.extension_metadata).map_err(|e| { + UnexpectedSnafu { + reason: format!("invalid JSON2 extension metadata: {e}"), + } + .build() + })?; + Ok(Self { + logical_settings: metadata.into_json_settings(), + target_layout: layout.target_layout.clone(), + }) + } +} + +/// Parses rewrite metadata and validates target layouts against the output schema. +fn resolve_align_mode(mode: AlignMode, output_schema: &Schema) -> Result { + let AlignMode::Rewrite { columns } = mode else { + return Ok(ResolvedAlignMode::AlignToSchema); + }; + + let mut rewrite_columns = HashMap::with_capacity(columns.len()); + for (name, layout) in columns { + let settings = RewriteSettings::try_from(&layout)?; + let field = output_schema.field_with_name(&name).map_err(|_| { + UnexpectedSnafu { + reason: format!("JSON2 rewrite column '{name}' is missing from output schema"), + } + .build() + })?; + ensure!( + is_json2_extension_type(field) + && field.data_type() == &json2_physical_data_type(&settings.target_layout), + UnexpectedSnafu { + reason: format!( + "JSON2 rewrite layout for column '{name}' does not match output field" + ), + } + ); + rewrite_columns.insert(name, settings); + } + + Ok(ResolvedAlignMode::Rewrite { + columns: rewrite_columns, + }) } #[cfg(test)] @@ -279,14 +350,21 @@ mod tests { target_layout: target_layout.clone(), }, )]); - let aligner = NestedSchemaAligner::new( + let aligner = JsonSchemaAligner::new( stream::empty::>(), - vec![], - schema(Vec::::new()), - )? - .with_json2_rewrite_targets(&rewrite_targets)?; - - let settings = &aligner.json2_rewrite_targets["j"]; + vec![false], + schema([ + Field::new("j", json2_physical_data_type(&target_layout), true) + .with_extension_type(Json2ExtensionType::default()), + ]), + AlignMode::Rewrite { + columns: rewrite_targets, + }, + )?; + let ResolvedAlignMode::Rewrite { columns } = &aligner.mode else { + panic!("expected rewrite mode"); + }; + let settings = &columns["j"]; assert_eq!(logical_settings, settings.logical_settings); assert_eq!(target_layout, settings.target_layout); Ok(()) @@ -305,8 +383,13 @@ mod tests { .unwrap(); let stream = stream::iter([Ok(input.clone())]); - let mut aligner = - NestedSchemaAligner::new(stream, vec![true, true], output_schema.clone()).unwrap(); + let mut aligner = JsonSchemaAligner::new( + stream, + vec![true, true], + output_schema.clone(), + AlignMode::AlignToSchema, + ) + .unwrap(); let output = aligner.next().await.unwrap().unwrap(); assert_eq!(input, output); @@ -324,9 +407,13 @@ mod tests { let input = RecordBatch::try_new(input_schema, vec![int_array([10, 20])]).unwrap(); let stream = stream::iter([Ok(input)]); - let mut aligner = - NestedSchemaAligner::new(stream, vec![true, false, false], output_schema.clone()) - .unwrap(); + let mut aligner = JsonSchemaAligner::new( + stream, + vec![true, false, false], + output_schema.clone(), + AlignMode::AlignToSchema, + ) + .unwrap(); let output = aligner.next().await.unwrap().unwrap(); assert_eq!(output_schema, output.schema()); @@ -362,8 +449,13 @@ mod tests { let input = RecordBatch::try_new(input_schema, vec![int_array([10, 20])]).unwrap(); let stream = stream::iter([Ok(input)]); - let mut aligner = - NestedSchemaAligner::new(stream, vec![true, false], output_schema.clone()).unwrap(); + let mut aligner = JsonSchemaAligner::new( + stream, + vec![true, false], + output_schema.clone(), + AlignMode::AlignToSchema, + ) + .unwrap(); let output = aligner.next().await.unwrap().unwrap(); assert_eq!(output_schema, output.schema()); @@ -377,8 +469,13 @@ mod tests { let output_schema = schema([Field::new("a", DataType::Int64, true)]); let stream = stream::iter([]); - let err = match NestedSchemaAligner::new(stream, vec![true, false], output_schema) { - Ok(_) => panic!("NestedSchemaAligner should reject projection length mismatch"), + let err = match JsonSchemaAligner::new( + stream, + vec![true, false], + output_schema, + AlignMode::AlignToSchema, + ) { + Ok(_) => panic!("JsonSchemaAligner should reject projection length mismatch"), Err(err) => err, }; @@ -399,8 +496,13 @@ mod tests { let input = RecordBatch::try_new(input_schema, vec![int_array([1, 2])]).unwrap(); let stream = stream::iter([Ok(input)]); - let mut aligner = - NestedSchemaAligner::new(stream, vec![true, true, false], output_schema).unwrap(); + let mut aligner = JsonSchemaAligner::new( + stream, + vec![true, true, false], + output_schema, + AlignMode::AlignToSchema, + ) + .unwrap(); let err = aligner.next().await.unwrap().unwrap_err(); assert!( @@ -410,7 +512,7 @@ mod tests { } #[tokio::test] - async fn test_nested_schema_aligner_aligns_struct_field() { + async fn test_json_schema_aligner_aligns_struct_field() { let output_schema = schema([Field::new( "nested", DataType::Struct(Fields::from(vec![ @@ -432,9 +534,13 @@ mod tests { ) .unwrap(); - let mut aligner = - NestedSchemaAligner::new(stream::iter([Ok(input)]), vec![true], output_schema.clone()) - .unwrap(); + let mut aligner = JsonSchemaAligner::new( + stream::iter([Ok(input)]), + vec![true], + output_schema.clone(), + AlignMode::AlignToSchema, + ) + .unwrap(); let output = aligner.next().await.unwrap().unwrap(); assert_eq!(output_schema, output.schema()); @@ -448,7 +554,7 @@ mod tests { } #[tokio::test] - async fn test_nested_schema_aligner_decodes_variant_to_struct() { + async fn test_json_schema_aligner_decodes_variant_to_struct() { let source_values = [ Some(parse_string_to_jsonb("1").unwrap()), Some(parse_string_to_jsonb(r#"{"b":2}"#).unwrap()), @@ -481,9 +587,13 @@ mod tests { true, ) .with_extension_type(Json2ExtensionType::default())]); - let mut aligner = - NestedSchemaAligner::new(stream::iter([Ok(input)]), vec![true], output_schema.clone()) - .unwrap(); + let mut aligner = JsonSchemaAligner::new( + stream::iter([Ok(input)]), + vec![true], + output_schema.clone(), + AlignMode::AlignToSchema, + ) + .unwrap(); let output = aligner.next().await.unwrap().unwrap(); assert_eq!(output_schema, output.schema()); @@ -516,7 +626,7 @@ mod tests { } #[tokio::test] - async fn test_nested_schema_aligner_preserves_struct_siblings() { + async fn test_json_schema_aligner_preserves_struct_siblings() { let source_values = [ Some(parse_string_to_jsonb(r#"{"x":1}"#).unwrap()), Some(parse_string_to_jsonb(r#"{"x":2}"#).unwrap()), @@ -570,9 +680,13 @@ mod tests { true, ) .with_extension_type(Json2ExtensionType::default())]); - let mut aligner = - NestedSchemaAligner::new(stream::iter([Ok(input)]), vec![true], output_schema.clone()) - .unwrap(); + let mut aligner = JsonSchemaAligner::new( + stream::iter([Ok(input)]), + vec![true], + output_schema.clone(), + AlignMode::AlignToSchema, + ) + .unwrap(); let output = aligner.next().await.unwrap().unwrap(); assert_eq!(output_schema, output.schema()); @@ -605,6 +719,159 @@ mod tests { ); } + #[tokio::test] + async fn test_rewrite_multiple_columns_and_fill_missing_roots() { + let logical_settings = JsonSettings::default(); + let target_layout = JsonSettings::try_new(vec![], Some(0)).unwrap(); + let target_type = json2_physical_data_type(&target_layout); + let output_schema = schema([ + Field::new("j", target_type.clone(), true) + .with_extension_type(Json2ExtensionType::default()), + Field::new("missing", target_type.clone(), true) + .with_extension_type(Json2ExtensionType::default()), + Field::new("k", target_type.clone(), true) + .with_extension_type(Json2ExtensionType::default()), + Field::new("a", DataType::Int64, true), + ]); + let values = [Some(parse_string_to_jsonb(r#"{"x":1}"#).unwrap()), None]; + let source = Arc::new(BinaryArray::from_iter( + values.iter().map(|value| value.as_deref()), + )) as ArrayRef; + let source_field = Field::new("j", DataType::Binary, true) + .with_extension_type(Json2ExtensionType::default()); + let expected = JsonArray::from(&source) + .rewrite_to_v2(&source_field, &logical_settings, &target_layout) + .unwrap(); + let input = RecordBatch::try_new( + schema([ + source_field, + Field::new("k", DataType::Binary, true) + .with_extension_type(Json2ExtensionType::default()), + Field::new("a", DataType::Int64, true), + ]), + vec![source.clone(), source, int_array([10, 20])], + ) + .unwrap(); + let columns = ["j", "missing", "k"] + .into_iter() + .map(|name| { + ( + name.to_string(), + Json2TargetLayout { + extension_metadata: serde_json::to_string(&JsonMetadata::new( + logical_settings.clone(), + )) + .unwrap(), + target_layout: target_layout.clone(), + }, + ) + }) + .collect(); + let mut aligner = JsonSchemaAligner::new( + stream::iter([Ok(input)]), + vec![true, false, true, true], + output_schema.clone(), + AlignMode::Rewrite { columns }, + ) + .unwrap(); + let output = aligner.next().await.unwrap().unwrap(); + assert_eq!(output_schema, output.schema()); + assert_eq!(expected.as_ref(), output.column(0).as_ref()); + assert_eq!(expected.as_ref(), output.column(2).as_ref()); + assert_eq!(&target_type, output.column(1).data_type()); + assert_eq!(2, output.column(1).null_count()); + assert_eq!(int_array([10, 20]).as_ref(), output.column(3).as_ref()); + } + + #[test] + fn test_rewrite_rejects_mismatched_output_layout() { + let columns = HashMap::from([( + "j".to_string(), + Json2TargetLayout { + extension_metadata: serde_json::to_string(&JsonMetadata::new( + JsonSettings::default(), + )) + .unwrap(), + target_layout: JsonSettings::try_new(vec![], Some(0)).unwrap(), + }, + )]); + let result = JsonSchemaAligner::new( + stream::empty::>(), + vec![false], + schema([Field::new("j", DataType::Binary, true) + .with_extension_type(Json2ExtensionType::default())]), + AlignMode::Rewrite { columns }, + ); + assert!( + result + .unwrap_err() + .to_string() + .contains("does not match output field") + ); + } + + #[tokio::test] + async fn test_empty_rewrite_only_fills_missing_roots() { + let source = int_array([10, 20]); + let input = RecordBatch::try_new( + schema([Field::new("a", DataType::Int64, true)]), + vec![source.clone()], + ) + .unwrap(); + let output_schema = schema([ + Field::new("missing", DataType::Utf8, true), + Field::new("a", DataType::Int64, true), + ]); + let mut aligner = JsonSchemaAligner::new( + stream::iter([Ok(input)]), + vec![false, true], + output_schema.clone(), + AlignMode::Rewrite { + columns: HashMap::new(), + }, + ) + .unwrap(); + let output = aligner.next().await.unwrap().unwrap(); + assert_eq!(output_schema, output.schema()); + assert_eq!(2, output.num_rows()); + assert_eq!(2, output.column(0).null_count()); + assert!(Arc::ptr_eq(&source, output.column(1))); + } + + #[tokio::test] + async fn test_empty_rewrite_does_not_align_existing_struct() { + let source = Arc::new(StructArray::from(vec![( + Arc::new(Field::new("x", DataType::Int64, true)), + int_array([1, 2]), + )])) as ArrayRef; + let input = RecordBatch::try_new( + schema([Field::new("j", source.data_type().clone(), true)]), + vec![source], + ) + .unwrap(); + let output_schema = schema([Field::new( + "j", + DataType::Struct(Fields::from(vec![ + Field::new("x", DataType::Int64, true), + Field::new("y", DataType::Utf8, true), + ])), + true, + )]); + let mut aligner = JsonSchemaAligner::new( + stream::iter([Ok(input)]), + vec![true], + output_schema, + AlignMode::Rewrite { + columns: HashMap::new(), + }, + ) + .unwrap(); + assert!(matches!( + aligner.next().await.unwrap(), + Err(crate::error::Error::NewRecordBatch { .. }) + )); + } + fn schema(fields: impl IntoIterator) -> SchemaRef { Arc::new(Schema::new(fields.into_iter().collect::>())) } diff --git a/src/mito2/src/sst/parquet/json_align/mod.rs b/src/mito2/src/sst/parquet/json_align/mod.rs deleted file mode 100644 index 54c97bdbb89..00000000000 --- a/src/mito2/src/sst/parquet/json_align/mod.rs +++ /dev/null @@ -1,24 +0,0 @@ -// Copyright 2023 Greptime Team -// -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -// See the License for the specific language governing permissions and -// limitations under the License. - -use datatypes::arrow::record_batch::RecordBatch; -use futures::stream::BoxStream; - -mod stream; - -pub(crate) use stream::NestedSchemaAligner; - -use crate::error::Result; - -pub(crate) type ProjectedRecordBatchStream = BoxStream<'static, Result>; diff --git a/src/mito2/src/sst/parquet/read_columns.rs b/src/mito2/src/sst/parquet/read_columns.rs index 58c7d0d2f8d..07614375dcf 100644 --- a/src/mito2/src/sst/parquet/read_columns.rs +++ b/src/mito2/src/sst/parquet/read_columns.rs @@ -12,13 +12,17 @@ // See the License for the specific language governing permissions and // limitations under the License. -use std::collections::{HashMap, HashSet}; +use std::collections::HashSet; +use std::ops::Range; -use datatypes::extension::json::JSON2_REMAINDER_FIELD_NAME; +use datatypes::arrow::datatypes::Schema as ArrowSchema; +use datatypes::extension::json::{JSON2_REMAINDER_FIELD_NAME, is_json2_extension_type}; use parquet::arrow::ProjectionMask; use parquet::basic::{ConvertedType, Type as PhysicalType}; use parquet::schema::types::{ColumnDescriptor, SchemaDescriptor}; +use crate::error::Result as MitoResult; + /// A nested field access path inside one parquet root column. pub type ParquetNestedPath = Vec; @@ -147,6 +151,32 @@ impl ParquetReadColumn { } } +/// Nested leaf selection semantics. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum NestedSelectionPolicy { + /// Also read the nearest JSONB Variant ancestor, or the remainder when it may + /// contain an unresolved path or children of a materialized object. + /// + /// For example, if `j.cold` has no matching field or Variant ancestor, read the + /// remainder because it may contain `j.cold`. When requesting a materialized + /// object such as `j.commit`, also read the remainder because it may contain + /// unmaterialized children of `j.commit`. + Json2, +} + +impl NestedSelectionPolicy { + /// Selects all leaves needed for one JSON2 root, including fallback data. + fn select_leaves( + self, + schema: &SchemaDescriptor, + leaf_range: Range, + col: &ParquetReadColumn, + selected: &mut HashSet, + ) { + select_json2_leaves(schema, leaf_range, col, selected); + } +} + /// Projection plan built for a parquet file. #[derive(Clone)] pub struct ProjectionMaskPlan { @@ -173,6 +203,9 @@ pub struct ProjectionMaskPlan { /// file. It is used to resolve requested nested paths to actual leaf /// column indices. /// +/// `source_schema` is the Arrow schema of the current parquet file. It is used +/// only for nested projections to identify JSON2 root fields. +/// /// See [`ProjectionMaskPlan`] for the returned value. /// /// For example, if the query requests `j.a` and `k`, but the current @@ -183,18 +216,20 @@ pub struct ProjectionMaskPlan { pub(crate) fn build_projection_plan( parquet_read_cols: &ParquetReadColumns, parquet_schema_desc: &SchemaDescriptor, -) -> ProjectionMaskPlan { + source_schema: &ArrowSchema, +) -> MitoResult { if !parquet_read_cols.has_nested() { let mask = ProjectionMask::roots(parquet_schema_desc, parquet_read_cols.root_indices_iter()); - return ProjectionMaskPlan { + + return Ok(ProjectionMaskPlan { mask, projected_root_presence: vec![true; parquet_read_cols.columns().len()], - }; + }); } let (matched_leaves, matched_roots) = - build_parquet_leaves_indices(parquet_schema_desc, parquet_read_cols); + build_parquet_leaves_indices(parquet_schema_desc, parquet_read_cols, source_schema)?; let projected_root_presence = parquet_read_cols .columns() @@ -203,10 +238,11 @@ pub(crate) fn build_projection_plan( .collect(); let mask = ProjectionMask::leaves(parquet_schema_desc, matched_leaves); - ProjectionMaskPlan { + + Ok(ProjectionMaskPlan { mask, projected_root_presence, - } + }) } /// Builds parquet leaf-column indices for reading a parquet file. @@ -217,89 +253,114 @@ pub(crate) fn build_projection_plan( fn build_parquet_leaves_indices( parquet_schema_desc: &SchemaDescriptor, projection: &ParquetReadColumns, -) -> (Vec, HashSet) { - let mut map = HashMap::with_capacity(projection.cols.len()); - for col in &projection.cols { - map.insert(col.root_index, col); - } + source_schema: &ArrowSchema, +) -> MitoResult<(Vec, HashSet)> { + let root_leaf_ranges = group_requested_leaf_ranges(parquet_schema_desc, projection); let mut matched_leaves = HashSet::new(); - let mut matched_roots = HashSet::with_capacity(projection.cols.len()); + let mut matched_roots = HashSet::with_capacity(projection.columns().len()); - // Tracks whether each requested nested path matched by prefix without fallback. - let mut prefix_matched = HashMap::>::new(); - for col in &projection.cols { - prefix_matched.insert(col.root_index, vec![false; col.nested_paths.len()]); - } + for col in projection.columns() { + let before = matched_leaves.len(); - // First select root/nested prefix matches. - for (leaf_idx, leaf_col) in parquet_schema_desc.columns().iter().enumerate() { - let root_idx = parquet_schema_desc.get_column_root_idx(leaf_idx); - let Some(col) = map.get(&root_idx) else { - continue; - }; - if col.nested_paths.is_empty() { - matched_leaves.insert(leaf_idx); - matched_roots.insert(root_idx); - continue; + let leaf_range = root_leaf_ranges[col.root_index()].clone(); + + if col.nested_paths().is_empty() { + matched_leaves.extend(leaf_range); + } else if is_json2_extension_type(&source_schema.fields()[col.root_index()]) { + NestedSelectionPolicy::Json2.select_leaves( + parquet_schema_desc, + leaf_range, + col, + &mut matched_leaves, + ); + } else { + select_prefix_leaves(parquet_schema_desc, leaf_range, col, &mut matched_leaves); } - let leaf_path = leaf_col.path().parts(); - let mut matched = false; - for (path_idx, _) in col - .nested_paths - .iter() - .enumerate() - .filter(|(_, nested_path)| leaf_path.starts_with(nested_path)) - { - prefix_matched.get_mut(&root_idx).unwrap()[path_idx] = true; - matched = true; - } - - if matched { - matched_leaves.insert(leaf_idx); - matched_roots.insert(root_idx); - } - } - - // Then include v2 remainder leaves or fallback prefix misses to their nearest variant parent. - // TODO(fys): Gate fallback planning on the root being JSON2. A raw Binary - // leaf is a JSONB variant only under a JSON2 root; plain struct Binary - // children should not enter this fallback path. - for col in &projection.cols { - let path_matches = &prefix_matched[&col.root_index]; - let mut needs_remainder = false; - for (matched, nested_path) in path_matches.iter().zip(&col.nested_paths) { - if *matched { - if !needs_remainder { - needs_remainder = - path_points_to_struct(parquet_schema_desc, col.root_index, nested_path); - } - continue; - } - - if let Some(leaf_idx) = - find_nearest_variant_parent(parquet_schema_desc, col.root_index, nested_path) - { - matched_leaves.insert(leaf_idx); - matched_roots.insert(col.root_index); - } else { - needs_remainder = true; - } - } - - if needs_remainder { - let remainder_leaves = find_remainder_leaves(parquet_schema_desc, col.root_index); - if !remainder_leaves.is_empty() { - matched_leaves.extend(remainder_leaves); - matched_roots.insert(col.root_index); - } + // Requested roots are unique, and leaves belong to exactly one root. + if matched_leaves.len() > before { + matched_roots.insert(col.root_index()); } } let mut matched_leaves = matched_leaves.into_iter().collect::>(); matched_leaves.sort_unstable(); - (matched_leaves, matched_roots) + + Ok((matched_leaves, matched_roots)) +} + +/// Groups file-level leaf ranges for requested roots. +fn group_requested_leaf_ranges( + schema: &SchemaDescriptor, + projection: &ParquetReadColumns, +) -> Vec> { + let root_count = schema.root_schema().get_fields().len(); + let mut requested = vec![false; root_count]; + for col in projection.columns() { + requested[col.root_index()] = true; + } + + let mut ranges = vec![0..0; root_count]; + let mut leaf_idx = 0; + while leaf_idx < schema.num_columns() { + let root_idx = schema.get_column_root_idx(leaf_idx); + let start = leaf_idx; + while leaf_idx < schema.num_columns() && schema.get_column_root_idx(leaf_idx) == root_idx { + leaf_idx += 1; + } + if requested[root_idx] { + ranges[root_idx] = start..leaf_idx; + } + } + ranges +} + +/// V2 can additionally store missing paths and object children in the remainder. +fn select_json2_leaves( + schema: &SchemaDescriptor, + leaf_range: Range, + col: &ParquetReadColumn, + selected: &mut HashSet, +) { + let prefix_matched = select_prefix_leaves(schema, leaf_range, col, selected); + let mut needs_remainder = false; + for (matched, path) in prefix_matched.iter().zip(&col.nested_paths) { + if *matched { + needs_remainder |= path_points_to_struct(schema, col.root_index, path); + } else if let Some(idx) = find_nearest_variant_parent(schema, col.root_index, path) { + selected.insert(idx); + } else { + needs_remainder = true; + } + } + if needs_remainder { + selected.extend(find_remainder_leaves(schema, col.root_index)); + } +} + +/// Selects prefix matches and returns a match flag for each requested nested path. +fn select_prefix_leaves( + schema: &SchemaDescriptor, + leaf_range: Range, + col: &ParquetReadColumn, + selected: &mut HashSet, +) -> Vec { + let mut prefix_matched = vec![false; col.nested_paths().len()]; + for leaf_idx in leaf_range { + let leaf_path = schema.columns()[leaf_idx].path().parts(); + let mut matched_leaf = false; + for (path, matched) in col.nested_paths().iter().zip(&mut prefix_matched) { + if leaf_path.starts_with(path) { + *matched = true; + matched_leaf = true; + } + } + if matched_leaf { + selected.insert(leaf_idx); + } + } + prefix_matched } /// Returns whether a nested path points to an explicitly materialized object. @@ -386,7 +447,10 @@ fn is_variant_leaf(leaf_col: &ColumnDescriptor) -> bool { mod tests { use std::sync::Arc; - use parquet::basic::{ConvertedType, LogicalType, Repetition, VariantType}; + use datatypes::arrow::datatypes::{DataType, Field, Fields, Schema as ArrowSchema}; + use datatypes::extension::json::{Json2ExtensionType, JsonMetadata}; + use datatypes::json::JsonSettings; + use parquet::basic::{LogicalType, Repetition, VariantType}; use parquet::errors::ParquetError; use parquet::schema::types::Type; @@ -397,7 +461,7 @@ mod tests { let parquet_schema_desc = build_test_nested_parquet_schema(); let projection = ParquetReadColumns::from_deduped_root_indices([0, 1]); - let plan = build_projection_plan(&projection, &parquet_schema_desc); + let plan = build_projection_plan(&projection, &parquet_schema_desc, &[None; 2]); assert_eq!(vec![true, true], plan.projected_root_presence); assert_eq!( @@ -412,8 +476,11 @@ mod tests { let projection = ParquetReadColumns::from_deduped(vec![ParquetReadColumn::new(0)]); - let (matched_leaves, matched_roots) = - build_parquet_leaves_indices(&parquet_schema_desc, &projection); + let (matched_leaves, matched_roots) = build_parquet_leaves_indices_with_policies( + &parquet_schema_desc, + &projection, + &[None; 2], + ); assert_eq!(vec![0, 1, 2], matched_leaves); assert_eq!(HashSet::from([0]), matched_roots); } @@ -428,8 +495,11 @@ mod tests { ParquetReadColumn::new(1), ]); - let (matched_leaves, matched_roots) = - build_parquet_leaves_indices(&parquet_schema_desc, &projection); + let (matched_leaves, matched_roots) = build_parquet_leaves_indices_with_policies( + &parquet_schema_desc, + &projection, + &[None; 2], + ); assert_eq!(vec![1, 2, 3], matched_leaves); assert_eq!(HashSet::from([0, 1]), matched_roots); } @@ -443,8 +513,11 @@ mod tests { .with_nested_paths(vec![vec!["j".to_string(), "b".to_string()]]), ]); - let (matched_leaves, matched_roots) = - build_parquet_leaves_indices(&parquet_schema_desc, &projection); + let (matched_leaves, matched_roots) = build_parquet_leaves_indices_with_policies( + &parquet_schema_desc, + &projection, + &[None; 2], + ); assert_eq!(vec![1, 2], matched_leaves); assert_eq!(HashSet::from([0]), matched_roots); } @@ -460,8 +533,11 @@ mod tests { let read_column = ParquetReadColumn::new(0).with_nested_paths(nested_paths); let projection = ParquetReadColumns::from_deduped(vec![read_column]); - let (matched_leaves, matched_roots) = - build_parquet_leaves_indices(&parquet_schema_desc, &projection); + let (matched_leaves, matched_roots) = build_parquet_leaves_indices_with_policies( + &parquet_schema_desc, + &projection, + &[None; 2], + ); assert_eq!(vec![1, 2], matched_leaves); assert_eq!(HashSet::from([0]), matched_roots); } @@ -475,8 +551,11 @@ mod tests { vec![vec!["j".to_string(), "b".to_string(), "c".to_string()]], )]); - let (matched_leaves, matched_roots) = - build_parquet_leaves_indices(&parquet_schema_desc, &projection); + let (matched_leaves, matched_roots) = build_parquet_leaves_indices_with_policies( + &parquet_schema_desc, + &projection, + &[None; 2], + ); assert_eq!(vec![1], matched_leaves); assert_eq!(HashSet::from([0]), matched_roots); } @@ -491,7 +570,7 @@ mod tests { ParquetReadColumn::new(1), ]); - let plan = build_projection_plan(&projection, &parquet_schema_desc); + let plan = build_projection_plan(&projection, &parquet_schema_desc, &[None; 2]); assert_eq!(vec![false, true], plan.projected_root_presence); assert_eq!( @@ -511,7 +590,8 @@ mod tests { ], )]); - let plan = build_projection_plan(&projection, &parquet); + let plan = + build_projection_plan(&projection, &parquet, &[Some(NestedSelectionPolicy::Json2)]); assert_eq!(vec![true], plan.projected_root_presence); assert_eq!(ProjectionMask::leaves(&parquet, [0, 1]), plan.mask); @@ -526,7 +606,8 @@ mod tests { .with_nested_paths(vec![vec!["j".to_string(), "hot".to_string()]]), ]); - let plan = build_projection_plan(&projection, &parquet); + let plan = + build_projection_plan(&projection, &parquet, &[Some(NestedSelectionPolicy::Json2)]); assert_eq!(vec![true], plan.projected_root_presence); assert_eq!(ProjectionMask::leaves(&parquet, [3]), plan.mask); @@ -541,7 +622,8 @@ mod tests { .with_nested_paths(vec![vec!["j".to_string(), "commit".to_string()]]), ]); - let plan = build_projection_plan(&projection, &parquet); + let plan = + build_projection_plan(&projection, &parquet, &[Some(NestedSelectionPolicy::Json2)]); assert_eq!(vec![true], plan.projected_root_presence); assert_eq!(ProjectionMask::leaves(&parquet, [0, 1, 2]), plan.mask); @@ -563,7 +645,8 @@ mod tests { ]], )]); - let plan = build_projection_plan(&projection, &parquet); + let plan = + build_projection_plan(&projection, &parquet, &[Some(NestedSelectionPolicy::Json2)]); assert_eq!(vec![true], plan.projected_root_presence); assert_eq!(ProjectionMask::leaves(&parquet, [4]), plan.mask); @@ -582,8 +665,11 @@ mod tests { ], )]); - let (matched_leaves, matched_roots) = - build_parquet_leaves_indices(&parquet_schema_desc, &projection); + let (matched_leaves, matched_roots) = build_parquet_leaves_indices_with_policies( + &parquet_schema_desc, + &projection, + &[None; 2], + ); assert_eq!(vec![0, 2], matched_leaves); assert_eq!(HashSet::from([0]), matched_roots); } @@ -614,137 +700,6 @@ mod tests { assert!(col.nested_paths().is_empty()); } - #[test] - fn test_fallback_to_nearest_variant_parent() { - let parquet_schema_desc = build_test_variant_parent_schema(); - let projection = - ParquetReadColumns::from_deduped(vec![ParquetReadColumn::new(0).with_nested_paths( - vec![vec!["j".to_string(), "a".to_string(), "x".to_string()]], - )]); - - let plan = build_projection_plan(&projection, &parquet_schema_desc); - - assert_eq!(vec![true], plan.projected_root_presence); - assert_eq!( - ProjectionMask::leaves(&parquet_schema_desc, vec![0]), - plan.mask - ); - } - - #[test] - fn test_prefix_match_prevents_variant_parent_fallback() { - let parquet_schema_desc = build_test_variant_parent_schema(); - let projection = - ParquetReadColumns::from_deduped(vec![ParquetReadColumn::new(0).with_nested_paths( - vec![vec!["j".to_string(), "b".to_string(), "x".to_string()]], - )]); - - let plan = build_projection_plan(&projection, &parquet_schema_desc); - - assert_eq!(vec![true], plan.projected_root_presence); - assert_eq!( - ProjectionMask::leaves(&parquet_schema_desc, vec![1]), - plan.mask - ); - } - - #[test] - fn test_mixed_prefix_and_fallback_paths() { - let parquet_schema_desc = build_test_variant_parent_schema(); - let nested_paths = vec![ - vec!["j".to_string(), "a".to_string(), "x".to_string()], - vec!["j".to_string(), "b".to_string(), "x".to_string()], - ]; - let projection = ParquetReadColumns::from_deduped(vec![ - ParquetReadColumn::new(0).with_nested_paths(nested_paths), - ]); - - let plan = build_projection_plan(&projection, &parquet_schema_desc); - - assert_eq!(vec![true], plan.projected_root_presence); - assert_eq!( - ProjectionMask::leaves(&parquet_schema_desc, vec![0, 1]), - plan.mask - ); - } - - #[test] - fn test_fallback_selects_multiple_variant_parents() { - let parquet_schema_desc = build_test_two_variant_parents_schema(); - let nested_paths = vec![ - vec!["j".to_string(), "a".to_string(), "y".to_string()], - vec!["j".to_string(), "b".to_string(), "d".to_string()], - vec!["j".to_string(), "a".to_string(), "x".to_string()], - ]; - let projection = ParquetReadColumns::from_deduped(vec![ - ParquetReadColumn::new(0).with_nested_paths(nested_paths), - ]); - - let plan = build_projection_plan(&projection, &parquet_schema_desc); - - assert_eq!(vec![true], plan.projected_root_presence); - assert_eq!( - ProjectionMask::leaves(&parquet_schema_desc, vec![0, 1]), - plan.mask - ); - } - - #[test] - fn test_nested_paths_fallback_to_variant_parent_by_default() { - let parquet_schema_desc = build_test_variant_parent_schema(); - let projection = - ParquetReadColumns::from_deduped(vec![ParquetReadColumn::new(0).with_nested_paths( - vec![vec!["j".to_string(), "a".to_string(), "x".to_string()]], - )]); - - let plan = build_projection_plan(&projection, &parquet_schema_desc); - - assert_eq!(vec![true], plan.projected_root_presence); - assert_eq!( - ProjectionMask::leaves(&parquet_schema_desc, vec![0]), - plan.mask - ); - } - - #[test] - fn test_non_variant_parent_does_not_fallback() { - let parquet_schema_desc = build_test_nested_parquet_schema(); - let projection = - ParquetReadColumns::from_deduped(vec![ParquetReadColumn::new(0).with_nested_paths( - vec![vec!["j".to_string(), "a".to_string(), "x".to_string()]], - )]); - - let plan = build_projection_plan(&projection, &parquet_schema_desc); - - assert_eq!(vec![false], plan.projected_root_presence); - } - - #[test] - fn test_utf8_parent_does_not_fallback() { - let parquet_schema_desc = build_test_utf8_parent_schema(); - let projection = - ParquetReadColumns::from_deduped(vec![ParquetReadColumn::new(0).with_nested_paths( - vec![vec!["j".to_string(), "a".to_string(), "x".to_string()]], - )]); - - let plan = build_projection_plan(&projection, &parquet_schema_desc); - - assert_eq!(vec![false], plan.projected_root_presence); - } - - #[test] - fn test_root_variant_does_not_fallback() { - let parquet_schema_desc = build_test_root_variant_schema(); - let projection = ParquetReadColumns::from_deduped(vec![ - ParquetReadColumn::new(0) - .with_nested_paths(vec![vec!["j".to_string(), "a".to_string()]]), - ]); - - let plan = build_projection_plan(&projection, &parquet_schema_desc); - - assert_eq!(vec![false], plan.projected_root_presence); - } - // Test schema: // schema // |- j @@ -859,127 +814,44 @@ mod tests { ))) } - // Test schema: - // schema - // `- j - // |- a: BYTE_ARRAY - // `- b - // `- x: INT64 - fn build_test_variant_parent_schema() -> SchemaDescriptor { - let leaf_a = Arc::new( - Type::primitive_type_builder("a", parquet::basic::Type::BYTE_ARRAY) - .with_repetition(Repetition::REQUIRED) - .build() - .unwrap(), - ); - let leaf_x = Arc::new( - Type::primitive_type_builder("x", parquet::basic::Type::INT64) - .with_repetition(Repetition::REQUIRED) - .build() - .unwrap(), - ); - let group_b = Arc::new( - Type::group_type_builder("b") - .with_repetition(Repetition::REQUIRED) - .with_fields(vec![leaf_x]) - .build() - .unwrap(), - ); - let root_j = Arc::new( - Type::group_type_builder("j") - .with_repetition(Repetition::REQUIRED) - .with_fields(vec![leaf_a, group_b]) - .build() - .unwrap(), - ); - let schema = Arc::new( - Type::group_type_builder("schema") - .with_fields(vec![root_j]) - .build() - .unwrap(), - ); - - SchemaDescriptor::new(schema) + fn source_schema(policies: &[Option]) -> ArrowSchema { + ArrowSchema::new( + policies + .iter() + .enumerate() + .map(|(index, policy)| { + let field = Field::new( + format!("root_{index}"), + DataType::Struct(Fields::empty()), + true, + ); + match policy { + None => field, + Some(NestedSelectionPolicy::Json2) => { + field.with_extension_type(Json2ExtensionType::new(Arc::new( + JsonMetadata::new(JsonSettings::default()), + ))) + } + } + }) + .collect::>(), + ) } - // Test schema: - // schema - // `- j - // |- a: BYTE_ARRAY - // `- b: BYTE_ARRAY - fn build_test_two_variant_parents_schema() -> SchemaDescriptor { - let leaf_a = Arc::new( - Type::primitive_type_builder("a", parquet::basic::Type::BYTE_ARRAY) - .with_repetition(Repetition::REQUIRED) - .build() - .unwrap(), - ); - let leaf_b = Arc::new( - Type::primitive_type_builder("b", parquet::basic::Type::BYTE_ARRAY) - .with_repetition(Repetition::REQUIRED) - .build() - .unwrap(), - ); - let root_j = Arc::new( - Type::group_type_builder("j") - .with_repetition(Repetition::REQUIRED) - .with_fields(vec![leaf_a, leaf_b]) - .build() - .unwrap(), - ); - let schema = Arc::new( - Type::group_type_builder("schema") - .with_fields(vec![root_j]) - .build() - .unwrap(), - ); - - SchemaDescriptor::new(schema) + fn build_projection_plan( + cols: &ParquetReadColumns, + parquet_schema: &SchemaDescriptor, + policies: &[Option], + ) -> ProjectionMaskPlan { + super::build_projection_plan(cols, parquet_schema, &source_schema(policies)).unwrap() } - fn build_test_utf8_parent_schema() -> SchemaDescriptor { - let leaf_a = Arc::new( - Type::primitive_type_builder("a", parquet::basic::Type::BYTE_ARRAY) - .with_repetition(Repetition::REQUIRED) - .with_logical_type(Some(LogicalType::String)) - .with_converted_type(ConvertedType::UTF8) - .build() - .unwrap(), - ); - let root_j = Arc::new( - Type::group_type_builder("j") - .with_repetition(Repetition::REQUIRED) - .with_fields(vec![leaf_a]) - .build() - .unwrap(), - ); - let schema = Arc::new( - Type::group_type_builder("schema") - .with_fields(vec![root_j]) - .build() - .unwrap(), - ); - - SchemaDescriptor::new(schema) - } - - // Test schema: - // schema - // `- j: BYTE_ARRAY - fn build_test_root_variant_schema() -> SchemaDescriptor { - let root_j = Arc::new( - Type::primitive_type_builder("j", parquet::basic::Type::BYTE_ARRAY) - .with_repetition(Repetition::REQUIRED) - .build() - .unwrap(), - ); - let schema = Arc::new( - Type::group_type_builder("schema") - .with_fields(vec![root_j]) - .build() - .unwrap(), - ); - - SchemaDescriptor::new(schema) + fn build_parquet_leaves_indices_with_policies( + parquet_schema: &SchemaDescriptor, + projection: &ParquetReadColumns, + policies: &[Option], + ) -> (Vec, HashSet) { + super::build_parquet_leaves_indices(parquet_schema, projection, &source_schema(policies)) + .unwrap() } } diff --git a/src/mito2/src/sst/parquet/reader.rs b/src/mito2/src/sst/parquet/reader.rs index 59ef12d9e1e..e3232c189d6 100644 --- a/src/mito2/src/sst/parquet/reader.rs +++ b/src/mito2/src/sst/parquet/reader.rs @@ -85,7 +85,7 @@ use crate::sst::parquet::file_range::{ }; use crate::sst::parquet::flat_format::{FlatReadFormat, primary_key_column_index}; use crate::sst::parquet::format::{INTERNAL_COLUMN_NUM, need_override_sequence}; -use crate::sst::parquet::json_align::{NestedSchemaAligner, ProjectedRecordBatchStream}; +use crate::sst::parquet::json_align::{AlignMode, JsonSchemaAligner, ProjectedRecordBatchStream}; use crate::sst::parquet::metadata::MetadataLoader; use crate::sst::parquet::prefilter::{ PrefilterContextBuilder, build_reader_filter_plan, execute_prefilter, @@ -522,13 +522,14 @@ impl ParquetReaderBuilder { let file_metadata = parquet_meta.file_metadata(); let parquet_schema_desc = file_metadata.schema_descr(); - let file_schema = + let file_schema = Arc::new( parquet_to_arrow_schema(parquet_schema_desc, file_metadata.key_value_metadata()) - .context(ParquetToArrowSchemaSnafu { file: &file_path })?; + .context(ParquetToArrowSchemaSnafu { file: &file_path })?, + ); let mut read_format = FlatReadFormat::new( region_meta.clone(), read_cols, - Some(Arc::new(file_schema)), + Some(file_schema.clone()), &file_path, skip_auto_convert, )?; @@ -565,7 +566,8 @@ impl ParquetReaderBuilder { // Computes the projection mask. let parquet_read_cols = read_format.parquet_read_columns(); - let projection_plan = build_projection_plan(parquet_read_cols, parquet_schema_desc); + let projection_plan = + build_projection_plan(parquet_read_cols, parquet_schema_desc, &file_schema)?; let has_nested_projection = parquet_read_cols.has_nested(); let selection = self .row_groups_to_read(&read_format, &parquet_meta, &mut metrics.filter_metrics) @@ -2070,12 +2072,20 @@ impl RowGroupReaderBuilder { return Ok(stream); } - Ok(NestedSchemaAligner::new( + let mode = if self.json2_rewrite_targets.is_empty() { + AlignMode::AlignToSchema + } else { + AlignMode::Rewrite { + columns: self.json2_rewrite_targets.clone(), + } + }; + + Ok(JsonSchemaAligner::new( stream, self.projection.projected_root_presence.clone(), self.output_schema.clone(), + mode, )? - .with_json2_rewrite_targets(&self.json2_rewrite_targets)? .boxed()) } @@ -2693,7 +2703,17 @@ mod tests { let output_schema = read_format.arrow_schema().clone(); let parquet_schema = parquet_meta.file_metadata().schema_descr(); - let projection = build_projection_plan(read_format.parquet_read_columns(), parquet_schema); + let source_schema = parquet_to_arrow_schema( + parquet_schema, + parquet_meta.file_metadata().key_value_metadata(), + ) + .unwrap(); + let projection = build_projection_plan( + read_format.parquet_read_columns(), + parquet_schema, + &source_schema, + ) + .unwrap(); let arrow_metadata = ArrowReaderMetadata::try_new(parquet_meta.clone(), ArrowReaderOptions::new()).unwrap(); ( @@ -2989,7 +3009,8 @@ mod tests { ParquetReadColumns::from_deduped(vec![ParquetReadColumn::new(0).with_nested_paths( vec![vec!["j".to_string(), "a".to_string(), "x".to_string()]], )]); - let projection_plan = build_projection_plan(&projection, parquet_schema); + let projection_plan = + build_projection_plan(&projection, parquet_schema, batch.schema_ref()).unwrap(); assert_eq!(vec![true], projection_plan.projected_root_presence); assert_eq!( projection_plan.mask,