From 56e738de06a969e1ec4a5350bb73ca9b15d8f514 Mon Sep 17 00:00:00 2001 From: zhangstar333 Date: Thu, 6 Aug 2026 18:20:20 +0800 Subject: [PATCH 1/2] [opt](paimon) support variant access path in paimon jni reader --- be/src/format/jni/jni_data_bridge.cpp | 2 +- be/src/format/jni/jni_data_bridge.h | 4 +- be/src/format/table/paimon_jni_reader.cpp | 4 +- be/src/format_v2/jni/jni_table_reader.cpp | 4 +- be/src/format_v2/jni/paimon_jni_reader.cpp | 13 + .../format_v2/parquet/native_schema_desc.cpp | 17 +- be/src/format_v2/parquet/native_schema_desc.h | 2 +- .../format_v2/jni/jni_table_reader_test.cpp | 6 +- .../format_v2/jni/paimon_jni_reader_test.cpp | 29 ++- .../format_v2/parquet/parquet_schema_test.cpp | 36 ++- .../doris/paimon/PaimonColumnValue.java | 28 ++- .../apache/doris/paimon/PaimonJniScanner.java | 69 +++++- .../doris/paimon/PaimonVariantProjection.java | 228 ++++++++++++++++++ .../doris/paimon/PaimonJniScannerTest.java | 14 ++ .../paimon/PaimonVariantProjectionTest.java | 147 +++++++++++ .../trees/plans/logical/LogicalFileScan.java | 3 +- .../plans/logical/LogicalFileScanTest.java | 17 ++ .../paimon/test_paimon_catalog_variant.out | 7 + .../paimon/test_paimon_catalog_variant.groovy | 38 +++ 19 files changed, 629 insertions(+), 39 deletions(-) create mode 100644 fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonVariantProjection.java create mode 100644 fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonVariantProjectionTest.java diff --git a/be/src/format/jni/jni_data_bridge.cpp b/be/src/format/jni/jni_data_bridge.cpp index adcf11e196e163..cd4f20b0bbdd63 100644 --- a/be/src/format/jni/jni_data_bridge.cpp +++ b/be/src/format/jni/jni_data_bridge.cpp @@ -521,7 +521,7 @@ std::string JniDataBridge::get_jni_type_with_different_string(const DataTypePtr& } } -std::string JniDataBridge::encode_schema_values(const std::vector& values) { +std::string JniDataBridge::encode_string_list(const std::vector& values) { std::vector encoded_values; encoded_values.reserve(values.size()); for (const auto& value : values) { diff --git a/be/src/format/jni/jni_data_bridge.h b/be/src/format/jni/jni_data_bridge.h index 267d0a1711c06c..339dd600f07214 100644 --- a/be/src/format/jni/jni_data_bridge.h +++ b/be/src/format/jni/jni_data_bridge.h @@ -140,8 +140,8 @@ class JniDataBridge { */ static std::string get_jni_type_with_different_string(const DataTypePtr& data_type); - /** Encodes every list element independently so schema delimiters remain unambiguous. */ - static std::string encode_schema_values(const std::vector& values); + /** Encodes every string independently so delimiters cannot change the list structure. */ + static std::string encode_string_list(const std::vector& values); /** Encodes STRUCT field names inside a JNI type descriptor. */ static std::string get_jni_type_with_encoded_struct_fields(const DataTypePtr& data_type); diff --git a/be/src/format/table/paimon_jni_reader.cpp b/be/src/format/table/paimon_jni_reader.cpp index 04dcc9daeec670..99c24a1ffd6162 100644 --- a/be/src/format/table/paimon_jni_reader.cpp +++ b/be/src/format/table/paimon_jni_reader.cpp @@ -74,8 +74,8 @@ PaimonJniReader::PaimonJniReader(const std::vector& file_slot_d params["columns_types"] = join(column_types, "#"); // V1 and V2 must publish the same safe schema protocol because session routing can select // either producer for the same Paimon scanner. - params["required_fields_base64"] = JniDataBridge::encode_schema_values(column_names); - params["columns_types_base64"] = JniDataBridge::encode_schema_values(encoded_column_types); + params["required_fields_base64"] = JniDataBridge::encode_string_list(column_names); + params["columns_types_base64"] = JniDataBridge::encode_string_list(encoded_column_types); params["time_zone"] = _state->timezone(); if (range_params->__isset.serialized_table) { params["serialized_table"] = range_params->serialized_table; diff --git a/be/src/format_v2/jni/jni_table_reader.cpp b/be/src/format_v2/jni/jni_table_reader.cpp index dfa00c1568ef2c..a5fbe61fd35d36 100644 --- a/be/src/format_v2/jni/jni_table_reader.cpp +++ b/be/src/format_v2/jni/jni_table_reader.cpp @@ -548,9 +548,9 @@ void JniTableReader::_prepare_jni_scanner_schema() { // Only Paimon consumes the paired payload. Keeping it capability-gated avoids recursively // encoding nested types for every split of unrelated V2 JNI connectors. _scanner_params["required_fields_base64"] = - JniDataBridge::encode_schema_values(required_fields); + JniDataBridge::encode_string_list(required_fields); _scanner_params["columns_types_base64"] = - JniDataBridge::encode_schema_values(encoded_column_types); + JniDataBridge::encode_string_list(encoded_column_types); } if (has_replace_type) { _scanner_params["replace_string"] = join(replace_types, ","); diff --git a/be/src/format_v2/jni/paimon_jni_reader.cpp b/be/src/format_v2/jni/paimon_jni_reader.cpp index 01f33c5cdf0396..1ec6b7c9d3b7a9 100644 --- a/be/src/format_v2/jni/paimon_jni_reader.cpp +++ b/be/src/format_v2/jni/paimon_jni_reader.cpp @@ -32,6 +32,7 @@ constexpr std::string_view HADOOP_OPTION_PREFIX = "hadoop."; constexpr std::string_view DORIS_ENABLE_JNI_IO_MANAGER = "jni.enable_jni_io_manager"; constexpr std::string_view DORIS_JNI_IO_MANAGER_TMP_DIR = "jni.io_manager.tmp_dir"; constexpr std::string_view PAIMON_JNI_SCANNER_IO_TMP_DIR = "paimon_jni_scanner_io_tmp"; +constexpr std::string_view VARIANT_ACCESS_PATH_PREFIX = "variant_access_path."; const std::string* get_paimon_predicate(const TFileScanRangeParams* scan_params, const TPaimonFileDesc& paimon_params) { @@ -135,6 +136,18 @@ Status PaimonJniReader::build_scanner_params(std::map* (*params)[std::string(HADOOP_OPTION_PREFIX) + kv.first] = kv.second; } } + // Keep paths aligned with the required_fields order. Each path is encoded independently so + // object keys containing delimiters cannot alter either path or segment cardinality. The Java + // reader validates whether a complete column path set can use Paimon's Variant projection and + // falls back to the full Variant when it cannot. + for (size_t column_idx = 0; column_idx < _projected_columns.size(); ++column_idx) { + const auto& access_paths = _projected_columns[column_idx].variant_access_paths; + for (size_t path_idx = 0; path_idx < access_paths.size(); ++path_idx) { + (*params)[std::string(VARIANT_ACCESS_PATH_PREFIX) + std::to_string(column_idx) + "." + + std::to_string(path_idx)] = + JniDataBridge::encode_string_list(access_paths[path_idx]); + } + } // TODO: Remove legacy split-level paimon_predicate, paimon_options and hadoop_conf from thrift // after the minimum supported FE always sends their scan-level replacements. return Status::OK(); diff --git a/be/src/format_v2/parquet/native_schema_desc.cpp b/be/src/format_v2/parquet/native_schema_desc.cpp index 16f4623bb1bafb..28ac3da4ccf9f9 100644 --- a/be/src/format_v2/parquet/native_schema_desc.cpp +++ b/be/src/format_v2/parquet/native_schema_desc.cpp @@ -286,7 +286,7 @@ class ScopedBoolOverride { Status validate_variant_layout(const NativeFieldSchema& group_field, std::optional specification_version, - bool allow_optional_shredded_metadata) { + bool allow_optional_shredded_fields) { if (specification_version.has_value() && *specification_version != 1) { return Status::NotSupported("Parquet Variant specification version {} is not supported", *specification_version); @@ -332,7 +332,7 @@ Status validate_variant_layout(const NativeFieldSchema& group_field, // unannotated overrides; row materialization still rejects null metadata for a non-null value. const bool valid_metadata_repetition = metadata_repetition == tparquet::FieldRepetitionType::REQUIRED || - (allow_optional_shredded_metadata && typed_value != nullptr && + (allow_optional_shredded_fields && typed_value != nullptr && metadata_repetition == tparquet::FieldRepetitionType::OPTIONAL); if (!metadata->children.empty() || metadata->physical_type != tparquet::Type::BYTE_ARRAY || !valid_metadata_repetition) { @@ -360,8 +360,17 @@ Status validate_variant_layout(const NativeFieldSchema& group_field, std::function validate_typed_value; std::function validate_wrapper; validate_wrapper = [&](const NativeFieldSchema& wrapper, WrapperContext context) -> Status { - if (!wrapper.parquet_schema.__isset.repetition_type || - wrapper.parquet_schema.repetition_type != tparquet::FieldRepetitionType::REQUIRED) { + const bool valid_wrapper_repetition = wrapper.parquet_schema.__isset.repetition_type && + (wrapper.parquet_schema.repetition_type == + tparquet::FieldRepetitionType::REQUIRED || + (allow_optional_shredded_fields && + wrapper.parquet_schema.repetition_type == + tparquet::FieldRepetitionType::OPTIONAL)); + // The Parquet Variant specification requires wrapper groups. Paimon's unannotated + // physical carrier makes them optional, so accept that representation only through the + // table-format override. Materialization still rejects an actually null array element; + // an absent object wrapper represents a missing key. + if (!valid_wrapper_repetition) { return Status::Corruption("Parquet Variant shredded wrapper {} must be required", wrapper.name); } diff --git a/be/src/format_v2/parquet/native_schema_desc.h b/be/src/format_v2/parquet/native_schema_desc.h index be0e739f80b6ff..38cd3bf0269d19 100644 --- a/be/src/format_v2/parquet/native_schema_desc.h +++ b/be/src/format_v2/parquet/native_schema_desc.h @@ -89,7 +89,7 @@ struct NativeFieldSchema { Status validate_variant_layout(const NativeFieldSchema& group_field, std::optional specification_version = std::nullopt, - bool allow_optional_shredded_metadata = false); + bool allow_optional_shredded_fields = false); // V2 owns this schema tree and parser so footer/schema planning never invokes the V1 reader path. class NativeFieldDescriptor { diff --git a/be/test/format_v2/jni/jni_table_reader_test.cpp b/be/test/format_v2/jni/jni_table_reader_test.cpp index a59a4bbdf560dc..b22d25abea69c0 100644 --- a/be/test/format_v2/jni/jni_table_reader_test.cpp +++ b/be/test/format_v2/jni/jni_table_reader_test.cpp @@ -116,10 +116,10 @@ Status init_reader(FakeJniTableReader* reader, const std::shared_ptr #include "core/data_type/data_type_string.h" +#include "core/data_type/data_type_variant_v2.h" #include "format_v2/table_reader.h" #include "gen_cpp/PlanNodes_types.h" #include "runtime/runtime_state.h" @@ -50,9 +51,10 @@ TFileScanRangeParams make_scan_params() { } Status init_reader(PaimonJniReader* reader, TFileScanRangeParams* scan_params, - RuntimeState* runtime_state = nullptr) { + RuntimeState* runtime_state = nullptr, + std::vector projected_columns = {}) { return reader->init({ - .projected_columns = {}, + .projected_columns = std::move(projected_columns), .conjuncts = {}, .format = FileFormat::JNI, .scan_params = scan_params, @@ -68,6 +70,29 @@ Status build_params(PaimonJniReader* reader, const TFileRangeDesc& range, return reader->build_scanner_params(params); } +TEST(PaimonJniReaderTest, PublishesVariantAccessPathsByProjectedColumnPosition) { + auto range = make_paimon_jni_range(); + range.table_format_params.paimon_params.__set_paimon_predicate("serialized-predicate"); + auto scan_params = make_scan_params(); + + ColumnDefinition id; + id.name = "id"; + id.type = std::make_shared(); + ColumnDefinition payload; + payload.name = "payload"; + payload.type = std::make_shared(); + payload.variant_access_paths = {{"name"}, {"profile", "city"}}; + + PaimonJniReader reader; + ASSERT_TRUE(init_reader(&reader, &scan_params, nullptr, {id, payload}).ok()); + + std::map params; + ASSERT_TRUE(build_params(&reader, range, ¶ms).ok()); + EXPECT_FALSE(params.contains("variant_access_path.0.0")); + EXPECT_EQ(params.at("variant_access_path.1.0"), "$bmFtZQ=="); + EXPECT_EQ(params.at("variant_access_path.1.1"), "$cHJvZmlsZQ==,$Y2l0eQ=="); +} + TEST(PaimonJniReaderTest, UsesScanLevelPredicateBeforeLegacySplitPredicate) { auto range = make_paimon_jni_range(); range.table_format_params.paimon_params.__set_paimon_predicate("legacy-predicate"); diff --git a/be/test/format_v2/parquet/parquet_schema_test.cpp b/be/test/format_v2/parquet/parquet_schema_test.cpp index a92a4dee22842e..c475348edc030b 100644 --- a/be/test/format_v2/parquet/parquet_schema_test.cpp +++ b/be/test/format_v2/parquet/parquet_schema_test.cpp @@ -247,19 +247,31 @@ TEST(ParquetSchemaTest, AppliesTableFormatVariantOverrideToUnannotatedGroup) { EXPECT_NE(fields[0]->variant_physical_type, nullptr); } -TEST(ParquetSchemaTest, AppliesPaimonShreddedVariantOverrideWithOptionalMetadata) { - auto schema = shredded_object_variant_schema(); - schema[1].__isset.logicalType = false; - schema[2].__set_repetition_type(tparquet::FieldRepetitionType::OPTIONAL); - NativeFieldDescriptor descriptor; - ASSERT_TRUE(descriptor.parse_from_thrift(schema).ok()); +TEST(ParquetSchemaTest, AppliesPaimonShreddedVariantOverrideWithOptionalFields) { + const auto apply_override = [](std::vector schema) { + NativeFieldDescriptor descriptor; + ASSERT_TRUE(descriptor.parse_from_thrift(schema).ok()); - std::vector> fields; - ASSERT_TRUE(build_parquet_column_schema(descriptor, &fields).ok()); - const std::vector overrides {format::LocalColumnIndex::top_level(format::LocalColumnId(0))}; - const auto status = apply_variant_schema_overrides(descriptor, overrides, &fields); - ASSERT_TRUE(status.ok()) << status; - EXPECT_EQ(fields[0]->kind, ParquetColumnSchemaKind::VARIANT); + std::vector> fields; + ASSERT_TRUE(build_parquet_column_schema(descriptor, &fields).ok()); + const std::vector overrides {format::LocalColumnIndex::top_level(format::LocalColumnId(0))}; + const auto status = apply_variant_schema_overrides(descriptor, overrides, &fields); + ASSERT_TRUE(status.ok()) << status; + ASSERT_EQ(fields.size(), 1); + EXPECT_EQ(fields[0]->kind, ParquetColumnSchemaKind::VARIANT); + }; + + auto object_schema = shredded_object_variant_schema(); + object_schema[1].__isset.logicalType = false; + object_schema[2].__set_repetition_type(tparquet::FieldRepetitionType::OPTIONAL); + object_schema[5].__set_repetition_type(tparquet::FieldRepetitionType::OPTIONAL); + apply_override(std::move(object_schema)); + + auto array_schema = shredded_array_variant_schema(true, true); + array_schema[1].__isset.logicalType = false; + array_schema[2].__set_repetition_type(tparquet::FieldRepetitionType::OPTIONAL); + array_schema[6].__set_repetition_type(tparquet::FieldRepetitionType::OPTIONAL); + apply_override(std::move(array_schema)); } TEST(ParquetSchemaTest, RejectsMalformedUnannotatedVariantOverride) { diff --git a/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonColumnValue.java b/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonColumnValue.java index d92983f5c0c47b..f8efad58e80995 100644 --- a/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonColumnValue.java +++ b/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonColumnValue.java @@ -65,6 +65,10 @@ public class PaimonColumnValue implements ColumnValue { private ColumnType dorisType; private DataType dataType; private ZoneId timeZone; + // have variant sub path project + private PaimonVariantProjection variantProjection; + // rebuild subpath variant to doris + private Variant materializedVariant; // Keep these caches lazy so scalar columns do not pay for complex-type reuse bookkeeping. private List arrayValues; private List mapKeys; @@ -88,13 +92,24 @@ private PaimonColumnValue( } public void setIdx(int idx, ColumnType dorisType, DataType dataType) { + setIdx(idx, dorisType, dataType, null); + } + + public void setIdx( + int idx, + ColumnType dorisType, + DataType dataType, + PaimonVariantProjection variantProjection) { this.idx = idx; this.dorisType = dorisType; this.dataType = dataType; + this.variantProjection = variantProjection; + this.materializedVariant = null; } public void setOffsetRow(InternalRow record) { this.record = record; + this.materializedVariant = null; } public void setTimeZone(String timeZone) { @@ -210,7 +225,16 @@ public byte[] getVariantValue() { } private Variant getVariant() { - return record.getVariant(idx); + // full variant object + if (variantProjection == null) { + return record.getVariant(idx); + } + // sub variant + if (materializedVariant == null) { + materializedVariant = variantProjection.materialize(record, idx); + } + // create another variant with return to doris, only subpath name in payload: element_at(payload, 'name') + return materializedVariant; } @Override @@ -292,6 +316,8 @@ private void reset( this.dorisType = dorisType; this.dataType = dataType; this.timeZone = timeZone; + this.variantProjection = null; + this.materializedVariant = null; } private static ZoneId resolveTimeZone(String timeZone) { diff --git a/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonJniScanner.java b/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonJniScanner.java index c7731dcdc3847f..03a12bc82456c2 100644 --- a/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonJniScanner.java +++ b/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonJniScanner.java @@ -39,8 +39,11 @@ import org.apache.paimon.table.source.Split; import org.apache.paimon.table.source.TableRead; import org.apache.paimon.table.system.SystemTableLoader; +import org.apache.paimon.types.DataField; import org.apache.paimon.types.DataType; +import org.apache.paimon.types.RowType; import org.apache.paimon.types.TimestampType; +import org.apache.paimon.types.VariantType; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -56,6 +59,7 @@ import java.nio.charset.StandardCharsets; import java.nio.file.Files; import java.nio.file.Paths; +import java.util.ArrayList; import java.util.Arrays; import java.util.Base64; import java.util.Collections; @@ -77,6 +81,7 @@ public class PaimonJniScanner extends JniScanner { private static final String ASYNC_READER_THREAD_NAME_PREFIX = "paimon-reader-async-thread"; private static final String FILE_READER_ASYNC_THRESHOLD = "file-reader-async-threshold"; private static final String SERIALIZED_TABLE = "serialized_table"; + private static final String VARIANT_ACCESS_PATH_PREFIX = "variant_access_path."; private static final int MAX_MANIFEST_PARALLELISM = 256; static final String DORIS_MANIFEST_PARALLELISM_CAP = "doris.scan.manifest.parallelism-cap"; @@ -99,6 +104,8 @@ public class PaimonJniScanner extends JniScanner { private final String paimonSplit; private final String paimonPredicate; private final String tableCacheKey; + private final String timeZone; + private final List>> variantAccessPathsByColumn; private Table table; private PaimonTableCache.TableCacheEntry tableCacheEntry; private RecordReader reader; @@ -107,6 +114,7 @@ public class PaimonJniScanner extends JniScanner { private final PaimonColumnValue columnValue = new PaimonColumnValue(); private List paimonAllFieldNames; private List paimonDataTypeList; + private List variantProjections; private RecordReader.RecordIterator recordIterator = null; private final ClassLoader classLoader; private PreExecutionAuthenticator preExecutionAuthenticator; @@ -142,8 +150,9 @@ public PaimonJniScanner(int batchSize, Map params) { tableCacheKey = params.get("serialized_table_cache_key"); Preconditions.checkState(tableCacheKey != null && !tableCacheKey.isEmpty(), "Missing required Paimon scanner parameter: serialized_table_cache_key"); - String timeZone = params.getOrDefault("time_zone", TimeZone.getDefault().getID()); + timeZone = params.getOrDefault("time_zone", TimeZone.getDefault().getID()); columnValue.setTimeZone(timeZone); + this.variantAccessPathsByColumn = variantAccessPathsByColumn(params, requiredFields.length); initTableInfo(columnTypes, requiredFields, batchSize); hadoopOptionParams = params.entrySet().stream() .filter(kv -> kv.getKey().startsWith(HADOOP_OPTION_PREFIX)) @@ -192,11 +201,36 @@ private void initReader() throws IOException { fields.length, paimonAllFieldNames.size())); } int[] projected = getProjected(); - readBuilder.withProjection(projected); + List readFields = new ArrayList<>(projected.length); + variantProjections = new ArrayList<>(projected.length); + boolean hasVariantProjection = false; + for (int outputIndex = 0; outputIndex < projected.length; outputIndex++) { + DataField tableField = table.rowType().getFields().get(projected[outputIndex]); + PaimonVariantProjection projection = tableField.type() instanceof VariantType + ? PaimonVariantProjection.create( + variantAccessPathsByColumn.get(outputIndex), timeZone) + : null; + variantProjections.add(projection); + if (projection == null) { + readFields.add(tableField); + } else { + hasVariantProjection = true; + readFields.add(tableField.newType( + projection.readType().copy(tableField.type().isNullable()))); + } + } + if (hasVariantProjection) { + // Paimon recognizes a metadata-marked RowType as a list of Variant extraction fields. + // It prunes matching shredded Parquet fields per file and reads the raw Variant value + // as a correctness fallback for unshredded or non-matching files. + RowType requestedReadType = new RowType(readFields); + readBuilder.withReadType(requestedReadType); + } else { + readBuilder.withProjection(projected); + } readBuilder.withFilter(getPredicates()); reader = newReadWithOptionalIOManager(readBuilder).executeFilter().createReader(getSplit()); - paimonDataTypeList = - Arrays.stream(projected).mapToObj(i -> table.rowType().getTypeAt(i)).collect(Collectors.toList()); + paimonDataTypeList = readFields.stream().map(DataField::type).collect(Collectors.toList()); } private TableRead newReadWithOptionalIOManager(ReadBuilder readBuilder) throws IOException { @@ -398,7 +432,8 @@ private int readAndProcessNextBatch() throws IOException { rows++; columnValue.setOffsetRow(record); for (int i = 0; i < fields.length; i++) { - columnValue.setIdx(i, types[i], paimonDataTypeList.get(i)); + columnValue.setIdx( + i, types[i], paimonDataTypeList.get(i), variantProjections.get(i)); appendData(i, columnValue); } if (rows >= batchSize) { @@ -536,14 +571,14 @@ static String[] requiredFields(Map params) { } // Each identifier is encoded independently, so delimiters in quoted identifiers cannot // change field cardinality. The legacy parameter remains the rolling-upgrade fallback. - return decodeSchemaValues(encodedFields); + return decodeStringList(encodedFields); } static String[] requiredTypes(Map params) { String encodedTypes = params.get("columns_types_base64"); return encodedTypes == null ? splitParam(params.get("columns_types"), "#") - : decodeSchemaValues(encodedTypes); + : decodeStringList(encodedTypes); } private static boolean usesEncodedSchema(Map params) { @@ -556,7 +591,7 @@ private static boolean usesEncodedSchema(Map params) { return hasFields; } - private static String[] decodeSchemaValues(String encodedValues) { + private static String[] decodeStringList(String encodedValues) { if (encodedValues.isEmpty()) { return new String[0]; } @@ -570,6 +605,24 @@ private static String[] decodeSchemaValues(String encodedValues) { .toArray(String[]::new); } + static List>> variantAccessPathsByColumn( + Map params, int requiredFieldCount) { + List>> result = new ArrayList<>(requiredFieldCount); + for (int columnIndex = 0; columnIndex < requiredFieldCount; columnIndex++) { + List> columnPaths = new ArrayList<>(); + for (int pathIndex = 0; ; pathIndex++) { + String encodedPath = params.get( + VARIANT_ACCESS_PATH_PREFIX + columnIndex + "." + pathIndex); + if (encodedPath == null) { + break; + } + columnPaths.add(Arrays.asList(decodeStringList(encodedPath))); + } + result.add(columnPaths); + } + return result; + } + static int countThreadsByNamePrefix(String threadNamePrefix) { int count = 0; ThreadMXBean threadMXBean = ManagementFactory.getThreadMXBean(); diff --git a/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonVariantProjection.java b/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonVariantProjection.java new file mode 100644 index 00000000000000..172b0f3bb463c4 --- /dev/null +++ b/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonVariantProjection.java @@ -0,0 +1,228 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you 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. + +package org.apache.doris.paimon; + +import org.apache.paimon.data.DataGetters; +import org.apache.paimon.data.InternalRow; +import org.apache.paimon.data.variant.GenericVariant; +import org.apache.paimon.data.variant.GenericVariantBuilder; +import org.apache.paimon.data.variant.Variant; +import org.apache.paimon.data.variant.VariantMetadataUtils; +import org.apache.paimon.types.DataField; +import org.apache.paimon.types.DataTypes; +import org.apache.paimon.types.RowType; + +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; + +/** + * Bridges Doris Variant access paths to Paimon's metadata-marked Variant extraction RowType. + * + *

Paimon returns one Variant value for every requested path. Doris expressions still consume a + * single Variant slot, so this class rebuilds a partial object containing only those paths. Missing + * paths are omitted, while a JSON null remains a present Variant null value. + */ +final class PaimonVariantProjection { + private static final String FIELD_NAME_PREFIX = "__doris_variant_field_"; + private static final int VARIANT_VALUE_INDEX = 0; + private static final int VARIANT_METADATA_INDEX = 1; + private static final int VARIANT_FIELD_COUNT = 2; + + private final RowType readType; + private final PathNode root; + + private PaimonVariantProjection(RowType readType, PathNode root) { + this.readType = readType; + this.root = root; + } + + /** + * Creates the metadata-marked RowType understood by Paimon's Variant reader. + * + *

Each Doris access path becomes one Variant field in {@link #readType}. Returning null + * means that the complete Variant column must be read instead. This all-or-nothing fallback is + * important because Doris still evaluates every original element_at expression after the scan. + */ + static PaimonVariantProjection create(List> paths, String timeZone) { + if (paths == null || paths.isEmpty()) { + return null; + } + + List fields = new ArrayList<>(paths.size()); + PathNode root = new PathNode(); + for (int fieldIndex = 0; fieldIndex < paths.size(); fieldIndex++) { + List path = paths.get(fieldIndex); + if (!supportsObjectPath(path) || !root.add(path, fieldIndex)) { + // Doris access paths currently do not retain whether a numeric segment came from + // an array index or an object key. Falling back avoids changing either meaning. + return null; + } + fields.add(new DataField( + fieldIndex, + FIELD_NAME_PREFIX + fieldIndex, + DataTypes.VARIANT(), + VariantMetadataUtils.buildVariantMetadata(toPaimonPath(path), false, timeZone))); + } + return new PaimonVariantProjection(new RowType(fields), root); + } + + /** Returns the logical type passed to Paimon's ReadBuilder.withReadType. */ + RowType readType() { + return readType; + } + + /** + * Rebuilds the extracted path values as one partial Variant object for Doris. + * + *

For example, Paimon returns separate fields for $.name and $.profile.city. This method + * turns them into {"name": ..., "profile": {"city": ...}}, so the unchanged Doris + * element_at expressions can continue to read a normal Variant slot. + */ + Variant materialize(DataGetters record, int fieldIndex) { + InternalRow extracted = record.getRow(fieldIndex, readType.getFieldCount()); + GenericVariantBuilder builder = new GenericVariantBuilder(false); + appendObject(builder, root, extracted); + return builder.result(); + } + + /** + * Checks whether a path can be represented unambiguously by Paimon's current Variant metadata. + * Array indexes and delimiter-bearing keys fall back to reading the complete Variant. + */ + private static boolean supportsObjectPath(List path) { + if (path == null || path.isEmpty()) { + return false; + } + for (String segment : path) { + if (segment == null || segment.isEmpty() || segment.indexOf('.') >= 0 + || segment.indexOf('[') >= 0 || segment.indexOf(';') >= 0 + || isIntegerSegment(segment)) { + return false; + } + } + return true; + } + + private static boolean isIntegerSegment(String segment) { + int offset = segment.startsWith("-") ? 1 : 0; + if (offset == segment.length()) { + return false; + } + for (int i = offset; i < segment.length(); i++) { + if (!Character.isDigit(segment.charAt(i))) { + return false; + } + } + return true; + } + + /** Converts Doris path segments such as [profile, city] to Paimon's $.profile.city syntax. */ + private static String toPaimonPath(List path) { + return "$." + String.join(".", path); + } + + /** Returns whether this path node or any descendant was present in the source Variant. */ + private static boolean hasValue(PathNode node, InternalRow extracted) { + if (node.fieldIndex >= 0) { + return hasExtractedVariant(extracted, node.fieldIndex); + } + for (PathNode child : node.children.values()) { + if (hasValue(child, extracted)) { + return true; + } + } + return false; + } + + private static boolean hasExtractedVariant(InternalRow extracted, int fieldIndex) { + // Paimon 1.4.2's RowToColumnConverter writes a Variant's value and metadata children but + // does not advance the enclosing HeapRowVector. Its null bitmap can therefore be shifted + // when a batch mixes present and missing paths. The two binary children remain aligned, + // so use them as the source of truth instead of extracted.isNullAt(fieldIndex). + InternalRow variant = extracted.getRow(fieldIndex, VARIANT_FIELD_COUNT); + boolean valueIsNull = variant.isNullAt(VARIANT_VALUE_INDEX); + boolean metadataIsNull = variant.isNullAt(VARIANT_METADATA_INDEX); + if (valueIsNull != metadataIsNull) { + throw new IllegalStateException( + "Paimon projected Variant must contain both value and metadata"); + } + return !valueIsNull; + } + + /** Reads the aligned value and metadata children as one Paimon Variant. */ + private static Variant getExtractedVariant(InternalRow extracted, int fieldIndex) { + InternalRow variant = extracted.getRow(fieldIndex, VARIANT_FIELD_COUNT); + return new GenericVariant( + variant.getBinary(VARIANT_VALUE_INDEX), + variant.getBinary(VARIANT_METADATA_INDEX)); + } + + /** + * Writes one object node recursively, omitting missing paths while preserving present JSON + * null values. Child insertion order follows the requested access-path order. + */ + private static void appendObject( + GenericVariantBuilder builder, PathNode node, InternalRow extracted) { + int start = builder.getWritePos(); + ArrayList fields = new ArrayList<>(); + for (Map.Entry entry : node.children.entrySet()) { + PathNode child = entry.getValue(); + if (!hasValue(child, extracted)) { + continue; + } + String key = entry.getKey(); + int dictionaryId = builder.addKey(key); + fields.add(new GenericVariantBuilder.FieldEntry( + key, dictionaryId, builder.getWritePos() - start)); + if (child.fieldIndex >= 0) { + Variant value = getExtractedVariant(extracted, child.fieldIndex); + builder.appendVariant(new GenericVariant(value.value(), value.metadata())); + } else { + appendObject(builder, child, extracted); + } + } + builder.finishWritingObject(start, fields); + } + + private static final class PathNode { + private final Map children = new LinkedHashMap<>(); + private int fieldIndex = -1; + + /** + * Adds one leaf path and records its position in Paimon's extracted Row. + * Parent/child overlaps and duplicate paths are rejected because one node cannot safely be + * materialized as both a leaf Variant and an object containing descendants. + */ + private boolean add(List path, int index) { + PathNode node = this; + for (String segment : path) { + if (node.fieldIndex >= 0) { + return false; + } + node = node.children.computeIfAbsent(segment, ignored -> new PathNode()); + } + if (!node.children.isEmpty() || node.fieldIndex >= 0) { + return false; + } + node.fieldIndex = index; + return true; + } + } +} diff --git a/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonJniScannerTest.java b/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonJniScannerTest.java index 8c1209a138753d..e9bd2e1a5a0931 100644 --- a/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonJniScannerTest.java +++ b/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonJniScannerTest.java @@ -64,6 +64,7 @@ import java.util.Base64; import java.util.Collections; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.CountDownLatch; @@ -496,6 +497,19 @@ public void testConstructorDecodesDelimiterSafeNestedFieldNamesAndTypes() { structType.getChildNames()); } + @Test + public void testVariantAccessPathsStayAlignedWithRequiredFields() { + Map params = createBaseParams(); + params.put("variant_access_path.1.0", encodeFields("name")); + params.put("variant_access_path.1.1", encodeFields("profile", "city")); + + List>> paths = PaimonJniScanner.variantAccessPathsByColumn(params, 3); + Assert.assertTrue(paths.get(0).isEmpty()); + Assert.assertEquals(Collections.singletonList("name"), paths.get(1).get(0)); + Assert.assertEquals(Arrays.asList("profile", "city"), paths.get(1).get(1)); + Assert.assertTrue(paths.get(2).isEmpty()); + } + @Test public void testConstructorDistinguishesEmptyIdentifierFromEmptyProjection() { Map oneEmptyField = createBaseParams(); diff --git a/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonVariantProjectionTest.java b/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonVariantProjectionTest.java new file mode 100644 index 00000000000000..7262068d06e043 --- /dev/null +++ b/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonVariantProjectionTest.java @@ -0,0 +1,147 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you 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. + +package org.apache.doris.paimon; + +import org.apache.paimon.data.GenericRow; +import org.apache.paimon.data.columnar.ColumnVector; +import org.apache.paimon.data.columnar.heap.HeapBytesVector; +import org.apache.paimon.data.columnar.heap.HeapRowVector; +import org.apache.paimon.data.variant.GenericVariant; +import org.apache.paimon.data.variant.Variant; +import org.apache.paimon.data.variant.VariantMetadataUtils; +import org.junit.Assert; +import org.junit.Test; + +import java.util.Arrays; +import java.util.Collections; + +public class PaimonVariantProjectionTest { + @Test + public void testBuildsMetadataMarkedReadTypeAndPartialVariant() { + PaimonVariantProjection projection = PaimonVariantProjection.create( + Arrays.asList( + Collections.singletonList("name"), + Arrays.asList("profile", "city"), + Collections.singletonList("missing")), + "Asia/Shanghai"); + + Assert.assertNotNull(projection); + Assert.assertEquals("$.name", VariantMetadataUtils.path( + projection.readType().getFields().get(0).description())); + Assert.assertEquals("$.profile.city", VariantMetadataUtils.path( + projection.readType().getFields().get(1).description())); + Assert.assertFalse(VariantMetadataUtils.failOnError( + projection.readType().getFields().get(1).description())); + + GenericRow extracted = projectedRecord( + GenericVariant.fromJson("\"alice\""), + GenericVariant.fromJson("\"beijing\""), + null); + Variant result = projection.materialize(extracted, 0); + Assert.assertEquals( + "{\"name\":\"alice\",\"profile\":{\"city\":\"beijing\"}}", + result.toJson()); + } + + @Test + public void testPreservesJsonNullButOmitsMissingPath() { + PaimonVariantProjection projection = PaimonVariantProjection.create( + Arrays.asList( + Collections.singletonList("present"), + Collections.singletonList("missing")), + "UTC"); + + Variant result = projection.materialize( + projectedRecord(GenericVariant.fromJson("null"), null), 0); + Assert.assertEquals("{\"present\":null}", result.toJson()); + } + + @Test + public void testIgnoresMisalignedPaimonVariantNullBitmap() { + PaimonVariantProjection projection = PaimonVariantProjection.create( + Collections.singletonList(Collections.singletonList("name")), "UTC"); + + GenericVariant alice = GenericVariant.fromJson("\"alice\""); + GenericVariant bob = GenericVariant.fromJson("\"bob\""); + HeapBytesVector values = new HeapBytesVector(3); + HeapBytesVector metadata = new HeapBytesVector(3); + appendVariant(values, metadata, alice); + appendVariant(values, metadata, bob); + + HeapRowVector variants = new HeapRowVector(3, values, metadata); + // Paimon 1.4.2 does not advance this row vector for non-null Variants. Appending the + // missing third value consequently marks row 0 null even though its binary children hold + // Alice; the binary children themselves still have the correct row alignment. + variants.appendNull(); + HeapRowVector extractedRows = new HeapRowVector(3, variants); + extractedRows.appendRow(); + extractedRows.appendRow(); + extractedRows.appendRow(); + + Assert.assertEquals( + "{\"name\":\"alice\"}", + projection.materialize(GenericRow.of(extractedRows.getRow(0)), 0).toJson()); + Assert.assertEquals( + "{\"name\":\"bob\"}", + projection.materialize(GenericRow.of(extractedRows.getRow(1)), 0).toJson()); + Assert.assertEquals( + "{}", projection.materialize(GenericRow.of(extractedRows.getRow(2)), 0).toJson()); + } + + private static void appendVariant( + HeapBytesVector values, HeapBytesVector metadata, GenericVariant variant) { + values.appendByteArray(variant.value(), 0, variant.value().length); + metadata.appendByteArray(variant.metadata(), 0, variant.metadata().length); + } + + private static GenericRow projectedRecord(Variant... extractedValues) { + ColumnVector[] fields = new ColumnVector[extractedValues.length]; + for (int i = 0; i < extractedValues.length; i++) { + HeapBytesVector values = new HeapBytesVector(1); + HeapBytesVector metadata = new HeapBytesVector(1); + HeapRowVector variant = new HeapRowVector(1, values, metadata); + if (extractedValues[i] == null) { + variant.appendNull(); + } else { + GenericVariant value = new GenericVariant( + extractedValues[i].value(), extractedValues[i].metadata()); + appendVariant(values, metadata, value); + variant.appendRow(); + } + fields[i] = variant; + } + HeapRowVector extracted = new HeapRowVector(1, fields); + extracted.appendRow(); + return GenericRow.of(extracted.getRow(0)); + } + + @Test + public void testFallsBackForAmbiguousOrUnsupportedPaths() { + Assert.assertNull(PaimonVariantProjection.create( + Collections.singletonList(Collections.singletonList("1")), "UTC")); + Assert.assertNull(PaimonVariantProjection.create( + Collections.singletonList(Collections.singletonList("a.b")), "UTC")); + Assert.assertNull(PaimonVariantProjection.create( + Collections.singletonList(Collections.singletonList("a;b")), "UTC")); + Assert.assertNull(PaimonVariantProjection.create( + Arrays.asList( + Collections.singletonList("profile"), + Arrays.asList("profile", "city")), + "UTC")); + } +} diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/logical/LogicalFileScan.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/logical/LogicalFileScan.java index 957c4b142e3639..53480d193ba54a 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/logical/LogicalFileScan.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/logical/LogicalFileScan.java @@ -329,7 +329,8 @@ public List computeAsteriskOutput() { @Override public boolean supportPruneNestedColumn() { ExternalTable table = getTable(); - if (table instanceof IcebergExternalTable || table instanceof IcebergSysExternalTable) { + if (table instanceof IcebergExternalTable || table instanceof IcebergSysExternalTable + || table instanceof PaimonExternalTable || table instanceof PaimonSysExternalTable) { return true; } else if (table instanceof HMSExternalTable) { HMSExternalTable hmsTable = (HMSExternalTable) table; diff --git a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/logical/LogicalFileScanTest.java b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/logical/LogicalFileScanTest.java index e3e3c56bc85c4f..eaaa24ebaf0f03 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/logical/LogicalFileScanTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/logical/LogicalFileScanTest.java @@ -98,6 +98,23 @@ public void testPaimonOptionsBindRelationScopedSnapshotSchema() { Mockito.verify(table, Mockito.never()).initSelectedPartitions(Mockito.any()); } + @Test + public void testPaimonSupportsNestedColumnPruning() { + PaimonExternalTable table = Mockito.mock(PaimonExternalTable.class); + Mockito.when(table.getName()).thenReturn("paimon_tbl"); + TableScanParams scanParams = new TableScanParams( + TableScanParams.OPTIONS, + Collections.singletonMap("scan.snapshot-id", "1"), + Collections.emptyList()); + Mockito.when(table.getFullSchema(scanParams)).thenReturn(Collections.emptyList()); + + LogicalFileScan scan = new LogicalFileScan(new RelationId(2), table, + Collections.singletonList("db"), Collections.emptyList(), + Optional.empty(), Optional.empty(), Optional.of(scanParams), Optional.empty()); + + Assertions.assertTrue(scan.supportPruneNestedColumn()); + } + @Test public void testCapturingRelationSchemaDoesNotAllocateOutputExprIds() throws Exception { StatementScopeIdGenerator.clear(); diff --git a/regression-test/data/external_table_p0/paimon/test_paimon_catalog_variant.out b/regression-test/data/external_table_p0/paimon/test_paimon_catalog_variant.out index b33530f97e034e..338936b6d20589 100644 --- a/regression-test/data/external_table_p0/paimon/test_paimon_catalog_variant.out +++ b/regression-test/data/external_table_p0/paimon/test_paimon_catalog_variant.out @@ -40,6 +40,13 @@ payload variant= 20 + order by id + """ + + // A table can contain files written before and after shredding was enabled. Paimon must + // apply physical projection per file while returning one consistent partial Variant. + order_qt_jni_mixed_us_projection """ + select id, + cast(payload['name'] as string), + cast(payload['layout'] as string) + from variant_mixed_us + order by id + """ + // Native reader cases. Reset every relevant switch explicitly so the JNI cases above do // not leak their session state into this block. sql """set enable_variant_v2 = true""" @@ -130,6 +157,17 @@ suite("test_paimon_catalog_variant", "p0,external,doris,external_docker,external } } + explain { + sql "select id, cast(payload['name'] as string) from variant_shredded order by id" + contains "all access paths: [payload.name]" + check { explainString -> + def nativeSplits = explainString =~ /paimonNativeReadSplits=(\d+)\/(\d+)/ + return nativeSplits.find() + && nativeSplits.group(1).toInteger() > 0 + && nativeSplits.group(1) == nativeSplits.group(2) + } + } + order_qt_native_full_variant """ select id, payload from variant_smoke From 64634e985e03c32c2a3630c49e169142f8ab54ba Mon Sep 17 00:00:00 2001 From: zhangstar333 Date: Fri, 7 Aug 2026 11:40:05 +0800 Subject: [PATCH 2/2] update prune type and test case --- .../paimon/run13.sql | 6 ++--- .../rules/rewrite/SlotTypeReplacer.java | 8 ++++-- .../trees/plans/logical/LogicalFileScan.java | 12 +++++++++ .../logical/SupportPruneNestedColumn.java | 7 ++++++ .../paimon/source/PaimonScanNodeTest.java | 21 ++++++++++++++-- .../plans/logical/LogicalFileScanTest.java | 14 ++++++++++- .../paimon/test_paimon_catalog_variant.out | 8 ++++++ .../paimon/test_paimon_catalog_variant.groovy | 25 +++++++++++++++++++ 8 files changed, 93 insertions(+), 8 deletions(-) diff --git a/docker/thirdparties/docker-compose/iceberg/scripts/create_preinstalled_scripts/paimon/run13.sql b/docker/thirdparties/docker-compose/iceberg/scripts/create_preinstalled_scripts/paimon/run13.sql index fc4506d0a1ebcc..dae545d3c04722 100644 --- a/docker/thirdparties/docker-compose/iceberg/scripts/create_preinstalled_scripts/paimon/run13.sql +++ b/docker/thirdparties/docker-compose/iceberg/scripts/create_preinstalled_scripts/paimon/run13.sql @@ -25,12 +25,12 @@ create table variant_shredded ( partitioned by (event_date) tblproperties ( 'file.format' = 'parquet', - 'parquet.variant.shreddingSchema' = '{"type":"ROW","fields":[{"name":"payload","type":{"type":"ROW","fields":[{"name":"name","type":"STRING"},{"name":"age","type":"INT"}]}}]}' + 'parquet.variant.shreddingSchema' = '{"type":"ROW","fields":[{"name":"payload","type":{"type":"ROW","fields":[{"name":"name","type":"STRING"},{"name":"age","type":"INT"},{"name":"profile","type":{"type":"ROW","fields":[{"name":"city","type":"STRING"}]}}]}}]}' ); insert into variant_shredded values - (1, date '2026-06-01', parse_json('{"name":"alice","age":18,"extra":"shredded"}')), - (2, date '2026-06-01', parse_json('{"name":"bob","age":30}')); + (1, date '2026-06-01', parse_json('{"name":"alice","age":18,"profile":{"city":"beijing"},"tags":["flink","paimon"],"extra":"shredded"}')), + (2, date '2026-06-01', parse_json('{"name":"bob","age":30,"profile":{"city":"shanghai"},"tags":["doris"]}')); drop table if exists variant_mixed_us; create table variant_mixed_us ( diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/SlotTypeReplacer.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/SlotTypeReplacer.java index c90d85d55e667d..de3de9adef620a 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/SlotTypeReplacer.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/SlotTypeReplacer.java @@ -723,14 +723,18 @@ private void replaceIcebergAccessPathToId(List originPath, int index, Da } private void tryRecordReplaceSlots(Plan plan, Object checkObj, Set shouldReplaceSlots) { - if (checkObj instanceof SupportPruneNestedColumn - && ((SupportPruneNestedColumn) checkObj).supportPruneNestedColumn()) { + if (checkObj instanceof SupportPruneNestedColumn) { + SupportPruneNestedColumn supportPruneNestedColumn = (SupportPruneNestedColumn) checkObj; + if (!supportPruneNestedColumn.supportPruneNestedColumn()) { + return; + } List output = plan.getOutput(); boolean shouldPrune = false; for (Slot slot : output) { int slotId = slot.getExprId().asInt(); if ((slot.getDataType() instanceof NestedColumnPrunable || slot.getDataType().isVariantType()) + && supportPruneNestedColumn.supportPruneNestedColumn(slot.getDataType()) && replacedDataTypes.containsKey(slotId)) { shouldReplaceSlots.add(slotId); shouldPrune = true; diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/logical/LogicalFileScan.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/logical/LogicalFileScan.java index 53480d193ba54a..c6727bf93c1eed 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/logical/LogicalFileScan.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/logical/LogicalFileScan.java @@ -43,6 +43,7 @@ import org.apache.doris.nereids.trees.plans.PlanType; import org.apache.doris.nereids.trees.plans.RelationId; import org.apache.doris.nereids.trees.plans.visitor.PlanVisitor; +import org.apache.doris.nereids.types.DataType; import org.apache.doris.nereids.util.Utils; import org.apache.doris.qe.ConnectContext; import org.apache.doris.qe.SessionVariable; @@ -357,6 +358,17 @@ public boolean supportPruneNestedColumn() { return false; } + @Override + public boolean supportPruneNestedColumn(DataType dataType) { + ExternalTable table = getTable(); + if (table instanceof PaimonExternalTable || table instanceof PaimonSysExternalTable) { + // Paimon JNI currently supports nested projection only for Variant. Its static complex + // types still return the full value, which would misalign pruned ROW/ARRAY/MAP slots. + return dataType.isVariantType(); + } + return supportPruneNestedColumn(); + } + private boolean hasSameSnapshot(Optional left, Optional right) { if (!left.isPresent() || !right.isPresent()) { return left.isPresent() == right.isPresent(); diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/logical/SupportPruneNestedColumn.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/logical/SupportPruneNestedColumn.java index 44f1b733dd109d..9776bb920a5a26 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/logical/SupportPruneNestedColumn.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/logical/SupportPruneNestedColumn.java @@ -17,8 +17,15 @@ package org.apache.doris.nereids.trees.plans.logical; +import org.apache.doris.nereids.types.DataType; + /** SupportPruneNestedColumn */ public interface SupportPruneNestedColumn { // return false will not prune the nested column boolean supportPruneNestedColumn(); + + // Allows a scan implementation to restrict pruning to selected root types. + default boolean supportPruneNestedColumn(DataType dataType) { + return supportPruneNestedColumn(); + } } diff --git a/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/source/PaimonScanNodeTest.java b/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/source/PaimonScanNodeTest.java index 50c37928ac93a8..c933c627de66c2 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/source/PaimonScanNodeTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/source/PaimonScanNodeTest.java @@ -22,7 +22,12 @@ import org.apache.doris.analysis.TableScanParams; import org.apache.doris.analysis.TupleDescriptor; import org.apache.doris.analysis.TupleId; +import org.apache.doris.catalog.ArrayType; import org.apache.doris.catalog.Column; +import org.apache.doris.catalog.MapType; +import org.apache.doris.catalog.StructField; +import org.apache.doris.catalog.StructType; +import org.apache.doris.catalog.Type; import org.apache.doris.catalog.VariantType; import org.apache.doris.common.ExceptionChecker; import org.apache.doris.common.UserException; @@ -108,10 +113,22 @@ public class PaimonScanNodeTest { private PaimonFileExternalCatalog paimonFileExternalCatalog; @Test - public void testVariantProjectionRequiresVariantV2() throws UserException { + public void testVariantProjectionRequiresVariantV2Recursively() throws UserException { + List variantTypes = Arrays.asList( + VariantType.COMPUTE_V2_INSTANCE, + new ArrayType(VariantType.COMPUTE_V2_INSTANCE), + new MapType(Type.STRING, VariantType.COMPUTE_V2_INSTANCE), + new StructType(new StructField("payload", VariantType.COMPUTE_V2_INSTANCE))); + + for (Type variantType : variantTypes) { + assertVariantProjectionRequiresVariantV2(variantType); + } + } + + private void assertVariantProjectionRequiresVariantV2(Type variantType) throws UserException { TupleDescriptor desc = new TupleDescriptor(new TupleId(0)); SlotDescriptor slot = new SlotDescriptor(new SlotId(0), desc); - slot.setColumn(new Column("payload", VariantType.COMPUTE_V2_INSTANCE)); + slot.setColumn(new Column("payload", variantType)); desc.addSlot(slot); ExceptionChecker.expectThrowsWithMsg(UserException.class, diff --git a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/logical/LogicalFileScanTest.java b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/logical/LogicalFileScanTest.java index eaaa24ebaf0f03..5fc6d422e09eb2 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/logical/LogicalFileScanTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/logical/LogicalFileScanTest.java @@ -34,6 +34,12 @@ import org.apache.doris.nereids.trees.expressions.StatementScopeIdGenerator; import org.apache.doris.nereids.trees.plans.RelationId; import org.apache.doris.nereids.trees.plans.logical.LogicalFileScan.SelectedPartitions; +import org.apache.doris.nereids.types.ArrayType; +import org.apache.doris.nereids.types.IntegerType; +import org.apache.doris.nereids.types.MapType; +import org.apache.doris.nereids.types.StructField; +import org.apache.doris.nereids.types.StructType; +import org.apache.doris.nereids.types.VariantType; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; @@ -99,7 +105,7 @@ public void testPaimonOptionsBindRelationScopedSnapshotSchema() { } @Test - public void testPaimonSupportsNestedColumnPruning() { + public void testPaimonSupportsOnlyVariantNestedColumnPruning() { PaimonExternalTable table = Mockito.mock(PaimonExternalTable.class); Mockito.when(table.getName()).thenReturn("paimon_tbl"); TableScanParams scanParams = new TableScanParams( @@ -113,6 +119,12 @@ public void testPaimonSupportsNestedColumnPruning() { Optional.empty(), Optional.empty(), Optional.of(scanParams), Optional.empty()); Assertions.assertTrue(scan.supportPruneNestedColumn()); + Assertions.assertTrue(scan.supportPruneNestedColumn(VariantType.INSTANCE)); + Assertions.assertFalse(scan.supportPruneNestedColumn(ArrayType.of(IntegerType.INSTANCE))); + Assertions.assertFalse(scan.supportPruneNestedColumn( + MapType.of(IntegerType.INSTANCE, IntegerType.INSTANCE))); + Assertions.assertFalse(scan.supportPruneNestedColumn(new StructType(Collections.singletonList( + new StructField("field", IntegerType.INSTANCE, true, ""))))); } @Test diff --git a/regression-test/data/external_table_p0/paimon/test_paimon_catalog_variant.out b/regression-test/data/external_table_p0/paimon/test_paimon_catalog_variant.out index 338936b6d20589..31c7fe287ae5fa 100644 --- a/regression-test/data/external_table_p0/paimon/test_paimon_catalog_variant.out +++ b/regression-test/data/external_table_p0/paimon/test_paimon_catalog_variant.out @@ -43,6 +43,14 @@ payload variant BE -> JNI -> + // Paimon readType path for a nested object projection, beyond the Java projection UT. + explain { + sql "select id, cast(payload['profile']['city'] as string) from variant_shredded" + contains "paimonNativeReadSplits=0/1" + contains "all access paths: [payload.profile.city]" + } + + order_qt_jni_nested_shredded_path """ + select id, cast(payload['profile']['city'] as string) + from variant_shredded + order by id + """ + + // Doris Variant array indexes are one-based. The numeric path segment is intentionally + // unsupported by Paimon's metadata projection, so it must make the whole Variant column + // fall back even though payload.name alone is projectable. + order_qt_jni_unsupported_array_path_fallback """ + select id, + cast(payload['name'] as string), + cast(payload['tags'][1] as string) + from variant_shredded + order by id + """ + // A table can contain files written before and after shredding was enabled. Paimon must // apply physical projection per file while returning one consistent partial Variant. order_qt_jni_mixed_us_projection """