Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion be/src/format/jni/jni_data_bridge.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
4 changes: 2 additions & 2 deletions be/src/format/jni/jni_data_bridge.h
Original file line number Diff line number Diff line change
Expand Up @@ -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<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);
Expand Down
4 changes: 2 additions & 2 deletions be/src/format/table/paimon_jni_reader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
4 changes: 2 additions & 2 deletions be/src/format_v2/jni/jni_table_reader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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, ",");
Expand Down
13 changes: 13 additions & 0 deletions be/src/format_v2/jni/paimon_jni_reader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -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();
Expand Down
17 changes: 13 additions & 4 deletions be/src/format_v2/parquet/native_schema_desc.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -286,7 +286,7 @@ class ScopedBoolOverride {

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);
Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -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);
}
Expand Down
2 changes: 1 addition & 1 deletion be/src/format_v2/parquet/native_schema_desc.h
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
6 changes: 3 additions & 3 deletions be/test/format_v2/jni/jni_table_reader_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
29 changes: 27 additions & 2 deletions be/test/format_v2/jni/paimon_jni_reader_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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,
Expand All @@ -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, &params).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");
Expand Down
36 changes: 24 additions & 12 deletions be/test/format_v2/parquet/parquet_schema_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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) {
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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) {
Expand Down
Loading
Loading