This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new 6a4700877d6 [opt](paimon) support variant access path in paimon jni
reader (#66547)
6a4700877d6 is described below
commit 6a4700877d60795ff04030da6265967c1bb7b7cd
Author: zhangstar333 <[email protected]>
AuthorDate: Fri Aug 7 14:34:30 2026 +0800
[opt](paimon) support variant access path in paimon jni reader (#66547)
### What problem does this PR solve?
Problem Summary:
support read variant sub path in jni reader using access path
```
| 0:VPAIMON_SCAN_NODE(87)
|
| table: test_variant.variant_db.variant_shredded
|
| inputSplitNum=1, totalFileSize=0, scanRanges=1
|
| partition=1/0
|
| cardinality=2, numNodes=1
|
| nested columns:
|
| payload:
|
| origin type: variant
|
| all access paths: [payload.name]
|
| pushdown agg=NONE
|
| paimonNativeReadSplits=0/1
|
| predicatesFromPaimon: NONE
|
| final projections: id[#0], CAST(element_at(payload[#1], 'name') AS
text) |
| final project output tuple id: 1
```
### Release note
None
### Check List (For Author)
- Test <!-- At least one of them must be included. -->
- [ ] Regression test
- [x] Unit Test
- [ ] Manual test (add detailed scripts or steps below)
- [ ] No need to test or manual test. Explain why:
- [ ] This is a refactor/code format and no logic has been changed.
- [ ] Previous test can cover this change.
- [ ] No code files have been changed.
- [ ] Other reason <!-- Add your reason? -->
- Behavior changed:
- [ ] No.
- [ ] Yes. <!-- Explain the behavior change -->
- Does this need documentation?
- [ ] No.
- [ ] Yes. <!-- Add document PR link here. eg:
https://github.com/apache/doris-website/pull/1214 -->
### Check List (For Reviewer who merge this PR)
- [ ] Confirm the release note
- [ ] Confirm test cases
- [ ] Confirm document
- [ ] Add branch pick label <!-- Add branch pick label that this PR
should merge into -->
---
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 ++
be/src/format_v2/parquet/native_schema_desc.cpp | 17 +-
be/src/format_v2/parquet/native_schema_desc.h | 2 +-
be/test/format_v2/jni/jni_table_reader_test.cpp | 6 +-
be/test/format_v2/jni/paimon_jni_reader_test.cpp | 29 ++-
be/test/format_v2/parquet/parquet_schema_test.cpp | 36 ++--
.../create_preinstalled_scripts/paimon/run13.sql | 6 +-
.../org/apache/doris/paimon/PaimonColumnValue.java | 28 ++-
.../org/apache/doris/paimon/PaimonJniScanner.java | 69 ++++++-
.../doris/paimon/PaimonVariantProjection.java | 228 +++++++++++++++++++++
.../apache/doris/paimon/PaimonJniScannerTest.java | 14 ++
.../doris/paimon/PaimonVariantProjectionTest.java | 147 +++++++++++++
.../nereids/rules/rewrite/SlotTypeReplacer.java | 8 +-
.../trees/plans/logical/LogicalFileScan.java | 15 +-
.../plans/logical/SupportPruneNestedColumn.java | 7 +
.../paimon/source/PaimonScanNodeTest.java | 21 +-
.../trees/plans/logical/LogicalFileScanTest.java | 29 +++
.../paimon/test_paimon_catalog_variant.out | 15 ++
.../paimon/test_paimon_catalog_variant.groovy | 63 ++++++
23 files changed, 721 insertions(+), 46 deletions(-)
diff --git a/be/src/format/jni/jni_data_bridge.cpp
b/be/src/format/jni/jni_data_bridge.cpp
index adcf11e196e..cd4f20b0bbd 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<std::string>& values) {
+std::string JniDataBridge::encode_string_list(const std::vector<std::string>&
values) {
std::vector<std::string> 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 267d0a1711c..339dd600f07 100644
--- a/be/src/format/jni/jni_data_bridge.h
+++ b/be/src/format/jni/jni_data_bridge.h
@@ -140,8 +140,8 @@ public:
*/
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<std::string>&
values);
+ /** Encodes every string independently so delimiters cannot change the
list structure. */
+ static std::string encode_string_list(const std::vector<std::string>&
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 04dcc9daeec..99c24a1ffd6 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<SlotDescriptor*>& 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 dfa00c1568e..a5fbe61fd35 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 01f33c5cdf0..1ec6b7c9d3b 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<std::string, std::string>*
(*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 16f4623bb1b..28ac3da4ccf 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 @@ private:
Status validate_variant_layout(const NativeFieldSchema& group_field,
std::optional<int8_t> 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<Status(const NativeFieldSchema&)> validate_typed_value;
std::function<Status(const NativeFieldSchema&, WrapperContext)>
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 be0e739f80b..38cd3bf0269 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<int8_t> 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 a59a4bbdf56..b22d25abea6 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<io::IOConte
}
TEST(JniTableReaderTest, RequiredFieldEncodingPreservesQuotedIdentifiers) {
- EXPECT_EQ(JniDataBridge::encode_schema_values({"region,code", "hash#name",
"地区 名"}),
+ EXPECT_EQ(JniDataBridge::encode_string_list({"region,code", "hash#name",
"地区 名"}),
"$cmVnaW9uLGNvZGU=,$aGFzaCNuYW1l,$5Zyw5Yy6IOWQjQ==");
- EXPECT_EQ(JniDataBridge::encode_schema_values({}), "");
- EXPECT_EQ(JniDataBridge::encode_schema_values({""}), "$");
+ EXPECT_EQ(JniDataBridge::encode_string_list({}), "");
+ EXPECT_EQ(JniDataBridge::encode_string_list({""}), "$");
}
TEST(JniTableReaderTest,
EncodedTypeDescriptorsPreserveNestedQuotedIdentifiers) {
diff --git a/be/test/format_v2/jni/paimon_jni_reader_test.cpp
b/be/test/format_v2/jni/paimon_jni_reader_test.cpp
index b97e6e6126b..8ce9e025a13 100644
--- a/be/test/format_v2/jni/paimon_jni_reader_test.cpp
+++ b/be/test/format_v2/jni/paimon_jni_reader_test.cpp
@@ -24,6 +24,7 @@
#include <utility>
#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<ColumnDefinition> 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<DataTypeString>();
+ ColumnDefinition payload;
+ payload.name = "payload";
+ payload.type = std::make_shared<DataTypeVariantV2>();
+ payload.variant_access_paths = {{"name"}, {"profile", "city"}};
+
+ PaimonJniReader reader;
+ ASSERT_TRUE(init_reader(&reader, &scan_params, nullptr, {id,
payload}).ok());
+
+ std::map<std::string, std::string> 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 a92a4dee228..c475348edc0 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<tparquet::SchemaElement>
schema) {
+ NativeFieldDescriptor descriptor;
+ ASSERT_TRUE(descriptor.parse_from_thrift(schema).ok());
- std::vector<std::unique_ptr<ParquetColumnSchema>> 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<std::unique_ptr<ParquetColumnSchema>> 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/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 fc4506d0a1e..dae545d3c04 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/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 d92983f5c0c..f8efad58e80 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<PaimonColumnValue> arrayValues;
private List<PaimonColumnValue> mapKeys;
@@ -88,13 +92,24 @@ public class PaimonColumnValue implements ColumnValue {
}
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 class PaimonColumnValue implements ColumnValue {
}
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 @@ public class PaimonColumnValue implements ColumnValue {
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 c7731dcdc38..03a12bc8245 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.ReadBuilder;
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.lang.reflect.InvocationTargetException;
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<List<List<String>>> variantAccessPathsByColumn;
private Table table;
private PaimonTableCache.TableCacheEntry tableCacheEntry;
private RecordReader<InternalRow> reader;
@@ -107,6 +114,7 @@ public class PaimonJniScanner extends JniScanner {
private final PaimonColumnValue columnValue = new PaimonColumnValue();
private List<String> paimonAllFieldNames;
private List<DataType> paimonDataTypeList;
+ private List<PaimonVariantProjection> variantProjections;
private RecordReader.RecordIterator<InternalRow> recordIterator = null;
private final ClassLoader classLoader;
private PreExecutionAuthenticator preExecutionAuthenticator;
@@ -142,8 +150,9 @@ public class PaimonJniScanner extends JniScanner {
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 @@ public class PaimonJniScanner extends JniScanner {
fields.length, paimonAllFieldNames.size()));
}
int[] projected = getProjected();
- readBuilder.withProjection(projected);
+ List<DataField> 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 @@ public class PaimonJniScanner extends JniScanner {
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 @@ public class PaimonJniScanner extends JniScanner {
}
// 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<String, String> 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<String, String> params) {
@@ -556,7 +591,7 @@ public class PaimonJniScanner extends JniScanner {
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 @@ public class PaimonJniScanner extends JniScanner {
.toArray(String[]::new);
}
+ static List<List<List<String>>> variantAccessPathsByColumn(
+ Map<String, String> params, int requiredFieldCount) {
+ List<List<List<String>>> result = new ArrayList<>(requiredFieldCount);
+ for (int columnIndex = 0; columnIndex < requiredFieldCount;
columnIndex++) {
+ List<List<String>> 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 00000000000..172b0f3bb46
--- /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.
+ *
+ * <p>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.
+ *
+ * <p>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<List<String>> paths, String
timeZone) {
+ if (paths == null || paths.isEmpty()) {
+ return null;
+ }
+
+ List<DataField> fields = new ArrayList<>(paths.size());
+ PathNode root = new PathNode();
+ for (int fieldIndex = 0; fieldIndex < paths.size(); fieldIndex++) {
+ List<String> 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.
+ *
+ * <p>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<String> 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<String> 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<GenericVariantBuilder.FieldEntry> fields = new ArrayList<>();
+ for (Map.Entry<String, PathNode> 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<String, PathNode> 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<String> 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 8c1209a1387..e9bd2e1a5a0 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.Arrays;
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 class PaimonJniScannerTest {
structType.getChildNames());
}
+ @Test
+ public void testVariantAccessPathsStayAlignedWithRequiredFields() {
+ Map<String, String> params = createBaseParams();
+ params.put("variant_access_path.1.0", encodeFields("name"));
+ params.put("variant_access_path.1.1", encodeFields("profile", "city"));
+
+ List<List<List<String>>> 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<String, String> 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 00000000000..7262068d06e
--- /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/rules/rewrite/SlotTypeReplacer.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/SlotTypeReplacer.java
index c90d85d55e6..de3de9adef6 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 @@ public class SlotTypeReplacer extends
DefaultPlanRewriter<Void> {
}
private void tryRecordReplaceSlots(Plan plan, Object checkObj,
Set<Integer> shouldReplaceSlots) {
- if (checkObj instanceof SupportPruneNestedColumn
- && ((SupportPruneNestedColumn)
checkObj).supportPruneNestedColumn()) {
+ if (checkObj instanceof SupportPruneNestedColumn) {
+ SupportPruneNestedColumn supportPruneNestedColumn =
(SupportPruneNestedColumn) checkObj;
+ if (!supportPruneNestedColumn.supportPruneNestedColumn()) {
+ return;
+ }
List<Slot> 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 957c4b142e3..c6727bf93c1 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.Plan;
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;
@@ -329,7 +330,8 @@ public class LogicalFileScan extends LogicalCatalogRelation
implements SupportPr
@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;
@@ -356,6 +358,17 @@ public class LogicalFileScan extends
LogicalCatalogRelation implements SupportPr
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<TableSnapshot> left,
Optional<TableSnapshot> 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 44f1b733dd1..9776bb920a5 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 50c37928ac9..c933c627de6 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.SlotId;
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<Type> 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 e3e3c56bc85..5fc6d422e09 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.Slot;
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;
@@ -98,6 +104,29 @@ public class LogicalFileScanTest {
Mockito.verify(table,
Mockito.never()).initSelectedPartitions(Mockito.any());
}
+ @Test
+ public void testPaimonSupportsOnlyVariantNestedColumnPruning() {
+ 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());
+
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
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 b33530f97e0..31c7fe287ae 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,21 @@ payload variant<PROPERTIES
("variant_max_subcolumns_count" = "0","variant_enable
1 alice-updated
2 bob
+-- !jni_shredded_projection --
+2 bob 30
+
+-- !jni_nested_shredded_path --
+1 beijing
+2 shanghai
+
+-- !jni_unsupported_array_path_fallback --
+1 alice flink
+2 bob doris
+
+-- !jni_mixed_us_projection --
+1 alice unshredded
+2 bob shredded
+
-- !native_full_variant --
1
{"active":true,"age":18,"missing":null,"name":"alice","profile":{"city":"beijing","zip":100000},"score":98.5,"tags":["flink","paimon"]}
2
{"active":false,"age":30,"extra":{"levels":[1,2,3]},"name":"bob","profile":{"city":"shanghai"},"tags":["doris"]}
diff --git
a/regression-test/suites/external_table_p0/paimon/test_paimon_catalog_variant.groovy
b/regression-test/suites/external_table_p0/paimon/test_paimon_catalog_variant.groovy
index fdba0df3f8a..dcc24bbac18 100644
---
a/regression-test/suites/external_table_p0/paimon/test_paimon_catalog_variant.groovy
+++
b/regression-test/suites/external_table_p0/paimon/test_paimon_catalog_variant.groovy
@@ -44,6 +44,12 @@ suite("test_paimon_catalog_variant",
"p0,external,doris,external_docker,external
contains "paimonNativeReadSplits=0/1"
}
+ explain {
+ sql "select id, cast(payload['name'] as string) from
variant_shredded order by id"
+ contains "paimonNativeReadSplits=0/1"
+ contains "all access paths: [payload.name]"
+ }
+
order_qt_desc """desc variant_smoke"""
order_qt_full_variant """
@@ -113,6 +119,52 @@ suite("test_paimon_catalog_variant",
"p0,external,doris,external_docker,external
order by id
"""
+ // Exercise Paimon's metadata-marked read type against a physically
shredded file. The
+ // unshredded and primary-key cases above cover Paimon's raw-value and
merge fallbacks.
+ order_qt_jni_shredded_projection """
+ select id,
+ cast(payload['name'] as string),
+ cast(payload['age'] as int)
+ from variant_shredded
+ where cast(payload['age'] as int) >= 20
+ order by id
+ """
+
+ // profile.city is physically shredded. This covers the complete FE ->
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 """
+ 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 +182,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
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]