mirror of
https://github.com/GreptimeTeam/greptimedb.git
synced 2026-09-30 17:15:38 +00:00
refactor(json2): JSON2 parquet projection and schema alignment (#9137)
* refactor(mito2): make JSON schema alignment targets explicit Signed-off-by: fys <fengys1996@gmail.com> * fix: cr * fix: cargo fmt --------- Signed-off-by: fys <fengys1996@gmail.com>
This commit is contained in:
@@ -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<ArrayRef> {
|
||||
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::<i32>();
|
||||
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<ArrayRef> {
|
||||
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<ArrayRef> {
|
||||
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<ArrayRef> {
|
||||
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::<i32>();
|
||||
widen_list(list_array, actual, expected)?
|
||||
}
|
||||
_ => JsonArray::from(list_array.values()).widen_scalar_to(expected.data_type())?,
|
||||
};
|
||||
Ok(Arc::new(
|
||||
GenericListArray::<i32>::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<ArrayRef, String>,
|
||||
}
|
||||
|
||||
impl TestCase {
|
||||
fn new(
|
||||
json_array: StructArray,
|
||||
schema_type: Fields,
|
||||
expected: std::result::Result<Vec<ArrayRef>, 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::<Int64Type, _, _>(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::<Int64Type, _, _>(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();
|
||||
|
||||
+386
-119
@@ -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<RecordBatch>>;
|
||||
|
||||
/// 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<String, Json2TargetLayout>,
|
||||
},
|
||||
}
|
||||
|
||||
/// 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<String, RewriteSettings>,
|
||||
},
|
||||
}
|
||||
|
||||
/// 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<S> {
|
||||
pub struct JsonSchemaAligner<S> {
|
||||
#[debug(skip)]
|
||||
inner: S,
|
||||
/// Output schema expected by the upper reader.
|
||||
@@ -70,79 +85,55 @@ pub struct NestedSchemaAligner<S> {
|
||||
/// 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<String, Json2RewriteSettings>,
|
||||
/// Alignment mode with parsed and validated rewrite settings.
|
||||
mode: ResolvedAlignMode,
|
||||
/// The cache for whether incoming batches already match output schema.
|
||||
is_schema_matched: Option<bool>,
|
||||
}
|
||||
|
||||
impl<S> NestedSchemaAligner<S>
|
||||
impl<S> JsonSchemaAligner<S>
|
||||
where
|
||||
S: Stream<Item = Result<RecordBatch>>,
|
||||
{
|
||||
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<bool>,
|
||||
output_schema: SchemaRef,
|
||||
) -> Result<NestedSchemaAligner<S>> {
|
||||
mode: AlignMode,
|
||||
) -> Result<JsonSchemaAligner<S>> {
|
||||
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<String, Json2TargetLayout>,
|
||||
) -> Result<Self> {
|
||||
self.json2_rewrite_targets = targets
|
||||
.iter()
|
||||
.map(|(name, layout)| {
|
||||
let metadata = serde_json::from_str::<JsonMetadata>(&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::<Result<_>>()?;
|
||||
Ok(self)
|
||||
}
|
||||
}
|
||||
|
||||
impl<S> Stream for NestedSchemaAligner<S>
|
||||
impl<S> Stream for JsonSchemaAligner<S>
|
||||
where
|
||||
S: Stream<Item = Result<RecordBatch>> + 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<String, Json2RewriteSettings>,
|
||||
mode: &ResolvedAlignMode,
|
||||
) -> Result<RecordBatch> {
|
||||
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<ArrayRef> {
|
||||
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<ArrayRef> {
|
||||
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<Self> {
|
||||
let metadata = from_str::<JsonMetadata>(&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<ResolvedAlignMode> {
|
||||
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::<Result<RecordBatch>>(),
|
||||
vec![],
|
||||
schema(Vec::<Field>::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::<Result<RecordBatch>>(),
|
||||
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<Item = Field>) -> SchemaRef {
|
||||
Arc::new(Schema::new(fields.into_iter().collect::<Vec<_>>()))
|
||||
}
|
||||
@@ -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<RecordBatch>>;
|
||||
@@ -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<String>;
|
||||
|
||||
@@ -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<usize>,
|
||||
col: &ParquetReadColumn,
|
||||
selected: &mut HashSet<usize>,
|
||||
) {
|
||||
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<ProjectionMaskPlan> {
|
||||
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<usize>, HashSet<usize>) {
|
||||
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<usize>, HashSet<usize>)> {
|
||||
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::<usize, Vec<bool>>::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::<Vec<_>>();
|
||||
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<Range<usize>> {
|
||||
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<usize>,
|
||||
col: &ParquetReadColumn,
|
||||
selected: &mut HashSet<usize>,
|
||||
) {
|
||||
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<usize>,
|
||||
col: &ParquetReadColumn,
|
||||
selected: &mut HashSet<usize>,
|
||||
) -> Vec<bool> {
|
||||
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<NestedSelectionPolicy>]) -> 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::<Vec<_>>(),
|
||||
)
|
||||
}
|
||||
|
||||
// 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<NestedSelectionPolicy>],
|
||||
) -> 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<NestedSelectionPolicy>],
|
||||
) -> (Vec<usize>, HashSet<usize>) {
|
||||
super::build_parquet_leaves_indices(parquet_schema, projection, &source_schema(policies))
|
||||
.unwrap()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
Reference in New Issue
Block a user