blob: 16516670a065c289d209044e8331c4af19e3e43b [file]
// 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.
#include "format/table/iceberg_reader.h"
#include <gen_cpp/Descriptors_types.h>
#include <gen_cpp/Metrics_types.h>
#include <gen_cpp/PlanNodes_types.h>
#include <gen_cpp/parquet_types.h>
#include <glog/logging.h>
#include <parallel_hashmap/phmap.h>
#include <rapidjson/document.h>
#include <algorithm>
#include <cstring>
#include <functional>
#include <memory>
#include "common/compiler_util.h" // IWYU pragma: keep
#include "common/consts.h"
#include "common/status.h"
#include "core/assert_cast.h"
#include "core/block/block.h"
#include "core/block/column_with_type_and_name.h"
#include "core/column/column.h"
#include "core/column/column_nullable.h"
#include "core/column/column_string.h"
#include "core/column/column_vector.h"
#include "core/data_type/data_type_factory.hpp"
#include "core/data_type/data_type_nullable.h"
#include "core/data_type/define_primitive_type.h"
#include "core/data_type/primitive_type.h"
#include "core/string_ref.h"
#include "exprs/aggregate/aggregate_function.h"
#include "format/format_common.h"
#include "format/generic_reader.h"
#include "format/orc/vorc_reader.h"
#include "format/parquet/schema_desc.h"
#include "format/parquet/vparquet_column_chunk_reader.h"
#include "format/table/deletion_vector_reader.h"
#include "format/table/iceberg/iceberg_orc_nested_column_utils.h"
#include "format/table/iceberg/iceberg_parquet_nested_column_utils.h"
#include "format/table/iceberg_scan_semantics.h"
#include "format/table/nested_column_access_helper.h"
#include "format/table/table_schema_change_helper.h"
#include "runtime/runtime_state.h"
#include "util/coding.h"
#include "util/string_util.h"
namespace cctz {
class time_zone;
} // namespace cctz
namespace doris {
class RowDescriptor;
class SlotDescriptor;
class TupleDescriptor;
namespace io {
struct IOContext;
} // namespace io
class VExprContext;
} // namespace doris
namespace doris {
namespace {
constexpr auto kIcebergOrcAttribute = "iceberg.id";
bool orc_subtree_has_iceberg_id(const orc::Type* type, const std::string& attribute) {
if (type->hasAttributeKey(attribute)) {
return true;
}
for (uint64_t idx = 0; idx < type->getSubtypeCount(); ++idx) {
if (orc_subtree_has_iceberg_id(type->getSubtype(idx), attribute)) {
return true;
}
}
return false;
}
bool parquet_subtree_has_iceberg_id(const FieldSchema& field) {
if (field.field_id >= 0) {
return true;
}
return std::ranges::any_of(field.children, parquet_subtree_has_iceberg_id);
}
struct ParquetEqualityFieldPath {
std::vector<const FieldSchema*> fields;
std::vector<size_t> child_indexes;
};
bool find_parquet_equality_field_path_by_id(const FieldDescriptor* descriptor, int32_t field_id,
ParquetEqualityFieldPath* result) {
DORIS_CHECK(descriptor != nullptr);
DORIS_CHECK(result != nullptr);
const auto find = [field_id](const auto& self, const FieldSchema* field,
ParquetEqualityFieldPath* path) -> bool {
DORIS_CHECK(field != nullptr);
path->fields.push_back(field);
if (field->field_id == field_id) {
return true;
}
for (size_t index = 0; index < field->children.size(); ++index) {
path->child_indexes.push_back(index);
if (self(self, &field->children[index], path)) {
return true;
}
path->child_indexes.pop_back();
}
path->fields.pop_back();
return false;
};
for (int index = 0; index < descriptor->size(); ++index) {
if (find(find, descriptor->get_column(index), result)) {
return true;
}
}
return false;
}
bool find_parquet_equality_field_prefix_by_id_path(
const FieldDescriptor* descriptor,
const std::vector<const schema::external::TField*>& table_path,
ParquetEqualityFieldPath* result) {
DORIS_CHECK(descriptor != nullptr);
DORIS_CHECK(result != nullptr);
DORIS_CHECK(!table_path.empty());
const std::vector<FieldSchema>* candidates = nullptr;
for (size_t path_index = 0; path_index < table_path.size(); ++path_index) {
const auto* table_field = table_path[path_index];
DORIS_CHECK(table_field != nullptr);
DORIS_CHECK(table_field->__isset.id);
const FieldSchema* match = nullptr;
size_t match_index = 0;
const size_t candidate_count =
candidates == nullptr ? cast_set<size_t>(descriptor->size()) : candidates->size();
for (size_t candidate_index = 0; candidate_index < candidate_count; ++candidate_index) {
const auto* candidate = candidates == nullptr
? descriptor->get_column(cast_set<int>(candidate_index))
: &(*candidates)[candidate_index];
if (candidate != nullptr && candidate->field_id == table_field->id) {
match = candidate;
match_index = candidate_index;
break;
}
}
if (match == nullptr) {
const auto wrapper =
candidates == nullptr
? TableSchemaChangeHelper::BuildTableInfoUtil::
find_unique_idless_parquet_wrapper_index(
*table_field, descriptor->get_fields_schema())
: TableSchemaChangeHelper::BuildTableInfoUtil::
find_unique_idless_parquet_wrapper_index(*table_field,
*candidates);
if (wrapper.has_value()) {
match_index = *wrapper;
match = candidates == nullptr ? descriptor->get_column(cast_set<int>(match_index))
: &(*candidates)[match_index];
}
}
if (match == nullptr) {
return false;
}
if (!result->fields.empty()) {
result->child_indexes.push_back(match_index);
}
result->fields.push_back(match);
candidates = &match->children;
}
return true;
}
std::vector<std::string> equality_field_name_candidates(const schema::external::TField& table_field,
const std::string* leaf_fallback) {
std::vector<std::string> candidates;
if (table_field.__isset.name_mapping) {
candidates.insert(candidates.end(), table_field.name_mapping.begin(),
table_field.name_mapping.end());
if (table_field.__isset.name_mapping_is_authoritative &&
table_field.name_mapping_is_authoritative) {
return candidates;
}
}
if (table_field.__isset.name) {
candidates.push_back(table_field.name);
}
if (leaf_fallback != nullptr) {
candidates.push_back(*leaf_fallback);
}
return candidates;
}
bool find_parquet_equality_field_prefix_by_name_path(
const FieldDescriptor* descriptor,
const std::vector<const schema::external::TField*>& table_path,
const std::string& leaf_fallback, ParquetEqualityFieldPath* result) {
DORIS_CHECK(descriptor != nullptr);
DORIS_CHECK(result != nullptr);
DORIS_CHECK(!table_path.empty());
const std::vector<FieldSchema>* children = nullptr;
for (size_t path_index = 0; path_index < table_path.size(); ++path_index) {
const auto* table_field = table_path[path_index];
DORIS_CHECK(table_field != nullptr);
const auto names = equality_field_name_candidates(
*table_field, path_index + 1 == table_path.size() ? &leaf_fallback : nullptr);
const FieldSchema* match = nullptr;
size_t match_index = 0;
const size_t child_count =
children == nullptr ? cast_set<size_t>(descriptor->size()) : children->size();
for (const auto& name : names) {
for (size_t child_index = 0; child_index < child_count; ++child_index) {
const auto* child = children == nullptr
? descriptor->get_column(cast_set<int>(child_index))
: &(*children)[child_index];
if (child != nullptr && iequal(child->name, name)) {
match = child;
match_index = child_index;
break;
}
}
if (match != nullptr) {
break;
}
}
if (match == nullptr) {
return false;
}
if (!result->fields.empty()) {
result->child_indexes.push_back(match_index);
}
result->fields.push_back(match);
children = &match->children;
}
return true;
}
struct OrcEqualityFieldPath {
std::vector<const orc::Type*> fields;
std::vector<std::string> names;
std::vector<size_t> child_indexes;
};
bool find_orc_equality_field_path_by_id(const orc::Type* root, int32_t field_id,
OrcEqualityFieldPath* result) {
DORIS_CHECK(root != nullptr);
DORIS_CHECK(result != nullptr);
const auto find = [field_id](const auto& self, const orc::Type* field,
const std::string& field_name,
OrcEqualityFieldPath* path) -> bool {
DORIS_CHECK(field != nullptr);
path->fields.push_back(field);
path->names.push_back(field_name);
if (field->hasAttributeKey(kIcebergOrcAttribute) &&
std::stoi(field->getAttributeValue(kIcebergOrcAttribute)) == field_id) {
return true;
}
for (size_t index = 0; index < field->getSubtypeCount(); ++index) {
path->child_indexes.push_back(index);
if (self(self, field->getSubtype(index), field->getFieldName(index), path)) {
return true;
}
path->child_indexes.pop_back();
}
path->fields.pop_back();
path->names.pop_back();
return false;
};
for (size_t index = 0; index < root->getSubtypeCount(); ++index) {
if (find(find, root->getSubtype(index), root->getFieldName(index), result)) {
return true;
}
}
return false;
}
bool find_orc_equality_field_prefix_by_id_path(
const orc::Type* root, const std::vector<const schema::external::TField*>& table_path,
OrcEqualityFieldPath* result) {
DORIS_CHECK(root != nullptr);
DORIS_CHECK(result != nullptr);
DORIS_CHECK(!table_path.empty());
const orc::Type* parent = root;
for (const auto* table_field : table_path) {
DORIS_CHECK(table_field != nullptr);
DORIS_CHECK(table_field->__isset.id);
const orc::Type* match = nullptr;
size_t match_index = 0;
for (size_t candidate_index = 0; candidate_index < parent->getSubtypeCount();
++candidate_index) {
const auto* candidate = parent->getSubtype(candidate_index);
if (candidate->hasAttributeKey(kIcebergOrcAttribute) &&
std::stoi(candidate->getAttributeValue(kIcebergOrcAttribute)) == table_field->id) {
match = candidate;
match_index = candidate_index;
break;
}
}
if (match == nullptr) {
const auto wrapper = TableSchemaChangeHelper::BuildTableInfoUtil::
find_unique_idless_orc_wrapper_index(*table_field, parent,
kIcebergOrcAttribute);
if (wrapper.has_value()) {
match_index = *wrapper;
match = parent->getSubtype(match_index);
}
}
if (match == nullptr) {
return false;
}
if (!result->fields.empty()) {
result->child_indexes.push_back(match_index);
}
result->fields.push_back(match);
result->names.push_back(parent->getFieldName(match_index));
parent = match;
}
return true;
}
bool find_orc_equality_field_prefix_by_name_path(
const orc::Type* root, const std::vector<const schema::external::TField*>& table_path,
const std::string& leaf_fallback, OrcEqualityFieldPath* result) {
DORIS_CHECK(root != nullptr);
DORIS_CHECK(result != nullptr);
DORIS_CHECK(!table_path.empty());
const orc::Type* parent = root;
for (size_t path_index = 0; path_index < table_path.size(); ++path_index) {
const auto* table_field = table_path[path_index];
DORIS_CHECK(table_field != nullptr);
const auto names = equality_field_name_candidates(
*table_field, path_index + 1 == table_path.size() ? &leaf_fallback : nullptr);
const orc::Type* match = nullptr;
size_t match_index = 0;
for (const auto& name : names) {
for (size_t child_index = 0; child_index < parent->getSubtypeCount(); ++child_index) {
if (iequal(parent->getFieldName(child_index), name)) {
match = parent->getSubtype(child_index);
match_index = child_index;
break;
}
}
if (match != nullptr) {
break;
}
}
if (match == nullptr) {
return false;
}
if (!result->fields.empty()) {
result->child_indexes.push_back(match_index);
}
result->fields.push_back(match);
result->names.push_back(parent->getFieldName(match_index));
parent = match;
}
return true;
}
} // namespace
const std::string IcebergOrcReader::ICEBERG_ORC_ATTRIBUTE = kIcebergOrcAttribute;
bool IcebergTableReader::_is_fully_dictionary_encoded(
const tparquet::ColumnMetaData& column_metadata) {
const auto is_dictionary_encoding = [](tparquet::Encoding::type encoding) {
return encoding == tparquet::Encoding::PLAIN_DICTIONARY ||
encoding == tparquet::Encoding::RLE_DICTIONARY;
};
const auto is_data_page = [](tparquet::PageType::type page_type) {
return page_type == tparquet::PageType::DATA_PAGE ||
page_type == tparquet::PageType::DATA_PAGE_V2;
};
const auto is_level_encoding = [](tparquet::Encoding::type encoding) {
return encoding == tparquet::Encoding::RLE || encoding == tparquet::Encoding::BIT_PACKED;
};
// A column chunk may have a dictionary page but still contain plain-encoded data pages.
// Only treat it as dictionary-coded when all data pages are dictionary encoded.
if (column_metadata.__isset.encoding_stats) {
bool has_data_page_stats = false;
for (const tparquet::PageEncodingStats& enc_stat : column_metadata.encoding_stats) {
if (is_data_page(enc_stat.page_type) && enc_stat.count > 0) {
has_data_page_stats = true;
if (!is_dictionary_encoding(enc_stat.encoding)) {
return false;
}
}
}
if (has_data_page_stats) {
return true;
}
}
bool has_dict_encoding = false;
bool has_nondict_encoding = false;
for (const tparquet::Encoding::type& encoding : column_metadata.encodings) {
if (is_dictionary_encoding(encoding)) {
has_dict_encoding = true;
}
if (!is_dictionary_encoding(encoding) && !is_level_encoding(encoding)) {
has_nondict_encoding = true;
break;
}
}
if (!has_dict_encoding || has_nondict_encoding) {
return false;
}
return true;
}
// ============================================================================
// IcebergParquetReader: on_before_init_reader (Parquet-specific schema matching)
// ============================================================================
// This format-specific setup mirrors the existing reader initialization sequence.
// NOLINTNEXTLINE(readability-function-cognitive-complexity,readability-function-size)
Status IcebergParquetReader::on_before_init_reader(ReaderInitContext* ctx) {
_column_descs = ctx->column_descs;
_fill_col_name_to_block_idx = ctx->col_name_to_block_idx;
_file_format = Fileformat::PARQUET;
// Get file metadata schema first (available because _open_file() already ran)
const FieldDescriptor* field_desc = nullptr;
RETURN_IF_ERROR(this->get_file_metadata_schema(&field_desc));
DCHECK(field_desc != nullptr);
// Build table_info_node by field_id or name matching.
// This must happen BEFORE column classification so we can use children_column_exists
// to check if a column exists in the file (by field ID, not name).
if (!get_scan_params().__isset.history_schema_info ||
get_scan_params().history_schema_info.empty()) [[unlikely]] {
RETURN_IF_ERROR(BuildTableInfoUtil::by_parquet_name(ctx->tuple_descriptor, *field_desc,
ctx->table_info_node));
} else {
RETURN_IF_ERROR(BuildTableInfoUtil::by_parquet_field_id_with_name_mapping(
get_scan_params().history_schema_info.front().root_field, *field_desc,
ctx->table_info_node, supports_iceberg_scan_semantics_v1(&get_scan_params())));
}
std::unordered_set<std::string> partition_col_names;
if (ctx->range->__isset.columns_from_path_keys) {
partition_col_names.insert(ctx->range->columns_from_path_keys.begin(),
ctx->range->columns_from_path_keys.end());
}
// Single pass: classify columns, detect $row_id, handle partition fallback.
bool has_partition_from_path = false;
for (const auto& desc : *ctx->column_descs) {
if (desc.category == ColumnCategory::SYNTHESIZED) {
if (desc.name == BeConsts::ICEBERG_ROWID_COL) {
this->register_synthesized_column_handler(
BeConsts::ICEBERG_ROWID_COL, [this](Block* block, size_t rows) -> Status {
return _fill_iceberg_row_id(block, rows);
});
continue;
} else if (desc.name.starts_with(BeConsts::GLOBAL_ROWID_COL)) {
auto topn_row_id_column_iter = _create_topn_row_id_column_iterator();
this->register_synthesized_column_handler(
desc.name,
[iter = std::move(topn_row_id_column_iter), this, &desc](
Block* block, size_t rows) -> Status {
return fill_topn_row_id(iter, desc.name, block, rows);
});
continue;
}
} else if (desc.category == ColumnCategory::PARTITION_KEY) {
bool has_partition_value = partition_col_names.contains(desc.name);
bool exists_in_file = ctx->table_info_node->children_column_exists(desc.name);
if (!has_partition_value || exists_in_file) {
// Keep PARTITION_KEY category stable for scan planning, but still read
// from file when the column exists there.
ctx->column_names.push_back(desc.name);
continue;
}
has_partition_from_path = true;
} else if (desc.category == ColumnCategory::REGULAR) {
ctx->column_names.push_back(desc.name);
} else if (desc.category == ColumnCategory::GENERATED) {
_init_row_lineage_columns();
if (desc.name == ROW_LINEAGE_ROW_ID) {
ctx->column_names.push_back(desc.name);
this->register_generated_column_handler(
ROW_LINEAGE_ROW_ID, [this](Block* block, size_t rows) -> Status {
return _fill_row_lineage_row_id(block, rows);
});
continue;
} else if (desc.name == ROW_LINEAGE_LAST_UPDATED_SEQ_NUMBER) {
ctx->column_names.push_back(desc.name);
this->register_generated_column_handler(
ROW_LINEAGE_LAST_UPDATED_SEQ_NUMBER,
[this](Block* block, size_t rows) -> Status {
return _fill_row_lineage_last_updated_sequence_number(block, rows);
});
continue;
}
}
}
// Set up partition value extraction if any partition columns need filling from path
if (has_partition_from_path) {
RETURN_IF_ERROR(_extract_partition_values(*ctx->range, ctx->tuple_descriptor,
_fill_partition_values,
&_fill_partition_value_is_null));
}
_all_required_col_names = ctx->column_names;
// Create column IDs from field descriptor
auto column_id_result =
_create_column_ids(field_desc, ctx->tuple_descriptor, ctx->table_info_node);
ctx->column_ids = std::move(column_id_result.column_ids);
ctx->filter_column_ids = std::move(column_id_result.filter_column_ids);
// Build field_id -> block_column_name mapping for equality delete filtering.
// This was previously done in init_reader() column matching (pre-CRTP refactoring).
for (const auto* slot : ctx->tuple_descriptor->slots()) {
_id_to_block_column_name.emplace(slot->col_unique_id(), slot->col_name());
}
// Process delete files (must happen before _do_init_reader so expand col IDs are included)
RETURN_IF_ERROR(_init_row_filters());
// Add expand column IDs for equality delete and remap expand column names
// to match master's behavior:
// - Use field_id to find the actual file column name in Parquet schema
// - Prefix with __equality_delete_column__ to avoid name conflicts
// - Correctly map table_col_name → file_col_name in table_info_node
const static std::string EQ_DELETE_PRE = "__equality_delete_column__";
bool all_file_columns_have_field_ids = true;
bool any_file_column_has_field_id = false;
for (int i = 0; i < field_desc->size(); ++i) {
const auto* field_schema = field_desc->get_column(i);
if (field_schema) {
if (field_schema->field_id < 0) {
all_file_columns_have_field_ids = false;
}
if (parquet_subtree_has_iceberg_id(*field_schema)) {
any_file_column_has_field_id = true;
}
}
}
const bool use_field_ids_for_hidden_keys =
supports_iceberg_scan_semantics_v1(&get_scan_params())
? any_file_column_has_field_id
: all_file_columns_have_field_ids;
const auto find_file_column_by_name = [&](const std::string& name) -> const FieldSchema* {
for (int j = 0; j < field_desc->size(); ++j) {
const auto* candidate = field_desc->get_column(j);
if (candidate != nullptr && iequal(candidate->name, name)) {
return candidate;
}
}
return nullptr;
};
// Rebuild _expand_col_names with proper file-column-based names
std::vector<std::string> new_expand_col_names;
DORIS_CHECK(_expand_col_names.size() == _expand_col_field_ids.size());
DORIS_CHECK(_expand_col_names.size() == _expand_columns.size());
for (size_t i = 0; i < _expand_col_names.size(); ++i) {
const auto& old_name = _expand_col_names[i];
const int32_t field_id = _expand_col_field_ids[i];
const FieldSchema* file_column = nullptr;
ParquetEqualityFieldPath file_path;
bool complete_file_path = false;
if (use_field_ids_for_hidden_keys) {
complete_file_path =
find_parquet_equality_field_path_by_id(field_desc, field_id, &file_path);
if (!complete_file_path && supports_iceberg_scan_semantics_v2(&get_scan_params())) {
const auto table_path = _find_schema_field_path(field_id);
if (!table_path.empty()) {
complete_file_path = find_parquet_equality_field_prefix_by_id_path(
field_desc, table_path, &file_path);
}
}
if (!file_path.fields.empty()) {
file_column = file_path.fields.front();
}
} else {
const auto table_path = _find_schema_field_path(field_id);
if (!table_path.empty()) {
complete_file_path = find_parquet_equality_field_prefix_by_name_path(
field_desc, table_path, old_name, &file_path);
if (!file_path.fields.empty()) {
file_column = file_path.fields.front();
}
} else {
file_column = find_file_column_by_name(old_name);
complete_file_path = file_column != nullptr;
}
}
std::string leaf_name = old_name;
if (!file_path.fields.empty()) {
leaf_name = file_path.fields.back()->name;
} else if (file_column != nullptr) {
leaf_name = file_column->name;
}
const std::string file_col_name = file_column == nullptr ? old_name : file_column->name;
std::string table_col_name = EQ_DELETE_PRE + std::to_string(field_id) + "_" + leaf_name;
// Update _id_to_block_column_name
if (field_id >= 0) {
_id_to_block_column_name[field_id] = table_col_name;
}
// Update _expand_columns name
_expand_columns[i].name = table_col_name;
if (file_column == nullptr) {
RETURN_IF_ERROR(_register_missing_equality_delete_column(field_id, table_col_name,
_expand_columns[i].type));
// The old data file predates this equality key. Keep it in the expand block so the
// synthesized-column hook can materialize its logical initial default before reader
// filtering, but do not advertise it to Parquet as a physical child.
new_expand_col_names.push_back(table_col_name);
continue;
}
new_expand_col_names.push_back(table_col_name);
if (!complete_file_path) {
ColumnPtr missing_value;
RETURN_IF_ERROR(_create_missing_equality_delete_value(
field_id, _expand_columns[i].type, file_path.fields.size(), &missing_value));
_nested_equality_delete_columns.push_back({
.field_id = field_id,
.block_name = table_col_name,
.leaf_type = _expand_columns[i].type,
.child_indexes = file_path.child_indexes,
.missing_value = std::move(missing_value),
});
_expand_columns[i].type = make_nullable(file_column->data_type);
_expand_columns[i].column = _expand_columns[i].type->create_column();
} else if (!file_path.child_indexes.empty()) {
_nested_equality_delete_columns.push_back({
.field_id = field_id,
.block_name = table_col_name,
.leaf_type = _expand_columns[i].type,
.child_indexes = file_path.child_indexes,
.missing_value = nullptr,
});
_expand_columns[i].type = make_nullable(file_column->data_type);
_expand_columns[i].column = _expand_columns[i].type->create_column();
}
// A hidden nested key is read through its containing top-level struct. V1 column IDs are
// pre-order ranges, so include the complete subtree before extracting the primitive leaf.
for (uint64_t column_id = file_column->get_column_id();
column_id <= file_column->get_max_column_id(); ++column_id) {
ctx->column_ids.insert(column_id);
}
// Register in table_info_node: table_col_name → file_col_name
ctx->column_names.push_back(table_col_name);
ctx->table_info_node->add_children(table_col_name, file_col_name,
TableSchemaChangeHelper::ConstNode::get_instance());
}
_expand_col_names = std::move(new_expand_col_names);
// Enable group filtering for Iceberg
_filter_groups = true;
return Status::OK();
}
// ============================================================================
// IcebergParquetReader: _create_column_ids
// ============================================================================
ColumnIdResult IcebergParquetReader::_create_column_ids(
const FieldDescriptor* field_desc, const TupleDescriptor* tuple_descriptor,
const std::shared_ptr<TableSchemaChangeHelper::Node>& table_info_node) {
auto* mutable_field_desc = const_cast<FieldDescriptor*>(field_desc);
mutable_field_desc->assign_ids();
std::unordered_map<int, const FieldSchema*> iceberg_id_to_field_schema_map;
for (int i = 0; i < field_desc->size(); ++i) {
const auto* field_schema = field_desc->get_column(i);
if (!field_schema) {
continue;
}
int iceberg_id = field_schema->field_id;
iceberg_id_to_field_schema_map[iceberg_id] = field_schema;
}
std::set<uint64_t> column_ids;
std::set<uint64_t> filter_column_ids;
auto process_access_paths = [](const FieldSchema* parquet_field,
const std::vector<TColumnAccessPath>& access_paths,
std::set<uint64_t>& out_ids) {
process_nested_access_paths(
parquet_field, access_paths, out_ids,
[](const FieldSchema* field) { return field->get_column_id(); },
[](const FieldSchema* field) { return field->get_max_column_id(); },
IcebergParquetNestedColumnUtils::extract_nested_column_ids);
};
// The Iceberg schema-mapping root is a StructNode whose registered children are the real
// table columns. When present, resolve each column by name through it so the column-id set
// stays consistent with the schema-mapping decision (BY_ID or BY_NAME/name-mapping);
// otherwise fall back to matching by Iceberg field id.
const auto* struct_node =
dynamic_cast<const TableSchemaChangeHelper::StructNode*>(table_info_node.get());
for (const auto* slot : tuple_descriptor->slots()) {
const FieldSchema* field_schema = nullptr;
if (struct_node != nullptr) {
// Synthesized/metadata slots (e.g. the TopN global row-id or the $row_id column) are
// never registered as children, so check membership before querying: calling
// children_column_exists() on an unregistered name DCHECK-aborts in debug builds and
// throws std::out_of_range from .at() in release builds.
if (struct_node->get_children().contains(slot->col_name()) &&
struct_node->children_column_exists(slot->col_name())) {
// Use the physical child selected by the schema-mapping pass. This keeps partial-id
// files in BY_NAME mode from binding a projected column through an unrelated stale
// field id.
const auto& file_column_name =
struct_node->children_file_column_name(slot->col_name());
for (int i = 0; i < field_desc->size(); ++i) {
const auto* candidate = field_desc->get_column(i);
if (candidate != nullptr && candidate->name == file_column_name) {
field_schema = candidate;
break;
}
}
DORIS_CHECK(field_schema != nullptr);
}
} else {
auto it = iceberg_id_to_field_schema_map.find(slot->col_unique_id());
if (it != iceberg_id_to_field_schema_map.end()) {
field_schema = it->second;
}
}
if (field_schema == nullptr) {
continue;
}
if ((slot->col_type() != TYPE_STRUCT && slot->col_type() != TYPE_ARRAY &&
slot->col_type() != TYPE_MAP)) {
column_ids.insert(field_schema->column_id);
if (slot->is_predicate()) {
filter_column_ids.insert(field_schema->column_id);
}
continue;
}
const auto& all_access_paths = slot->all_access_paths();
process_access_paths(field_schema, all_access_paths, column_ids);
const auto& predicate_access_paths = slot->predicate_access_paths();
if (!predicate_access_paths.empty()) {
process_access_paths(field_schema, predicate_access_paths, filter_column_ids);
}
}
return {std::move(column_ids), std::move(filter_column_ids)};
}
// ============================================================================
// IcebergParquetReader: _read_position_delete_file
// ============================================================================
Status IcebergParquetReader::_read_position_delete_file(const TFileRangeDesc* delete_range,
DeleteFile* position_delete) {
ParquetReader parquet_delete_reader(get_profile(), get_scan_params(), *delete_range,
READ_DELETE_FILE_BATCH_SIZE, &get_state()->timezone_obj(),
get_io_ctx(), get_state(), _meta_cache);
// The delete file range has size=-1 (read whole file). We must disable
// row group filtering before init; otherwise _do_init_reader returns EndOfFile
// when _filter_groups && _range_size < 0.
ParquetInitContext delete_ctx;
delete_ctx.filter_groups = false;
delete_ctx.column_names = delete_file_col_names;
delete_ctx.col_name_to_block_idx =
const_cast<std::unordered_map<std::string, uint32_t>*>(&DELETE_COL_NAME_TO_BLOCK_IDX);
RETURN_IF_ERROR(parquet_delete_reader.init_reader(&delete_ctx));
const tparquet::FileMetaData* meta_data = parquet_delete_reader.get_meta_data();
bool dictionary_coded = true;
for (const auto& row_group : meta_data->row_groups) {
const auto& column_chunk = row_group.columns[ICEBERG_FILE_PATH_INDEX];
if (!(column_chunk.__isset.meta_data && has_dict_page(column_chunk.meta_data))) {
dictionary_coded = false;
break;
}
}
DataTypePtr data_type_file_path = make_nullable(std::make_shared<DataTypeString>());
DataTypePtr data_type_pos = make_nullable(std::make_shared<DataTypeInt64>());
bool eof = false;
while (!eof) {
Block block = {
dictionary_coded
? ColumnWithTypeAndName {ColumnNullable::create(ColumnDictI32::create(),
ColumnUInt8::create()),
data_type_file_path, ICEBERG_FILE_PATH}
: ColumnWithTypeAndName {data_type_file_path, ICEBERG_FILE_PATH},
{data_type_pos, ICEBERG_ROW_POS}};
size_t read_rows = 0;
RETURN_IF_ERROR(parquet_delete_reader.get_next_block(&block, &read_rows, &eof));
if (read_rows <= 0) {
break;
}
RETURN_IF_ERROR(_gen_position_delete_file_range(block, position_delete, read_rows,
dictionary_coded));
}
return Status::OK();
};
// ============================================================================
// IcebergOrcReader: on_before_init_reader (ORC-specific schema matching)
// ============================================================================
// This format-specific setup mirrors the existing reader initialization sequence.
// NOLINTNEXTLINE(readability-function-cognitive-complexity,readability-function-size)
Status IcebergOrcReader::on_before_init_reader(ReaderInitContext* ctx) {
_column_descs = ctx->column_descs;
_fill_col_name_to_block_idx = ctx->col_name_to_block_idx;
_file_format = Fileformat::ORC;
// Get ORC file type first (available because _create_file_reader() already ran)
const orc::Type* orc_type_ptr = nullptr;
RETURN_IF_ERROR(this->get_file_type(&orc_type_ptr));
// Build table_info_node by field_id or name matching.
// This must happen BEFORE column classification so we can use children_column_exists
// to check if a column exists in the file (by field ID, not name).
if (!get_scan_params().__isset.history_schema_info ||
get_scan_params().history_schema_info.empty()) [[unlikely]] {
RETURN_IF_ERROR(BuildTableInfoUtil::by_orc_name(ctx->tuple_descriptor, orc_type_ptr,
ctx->table_info_node));
} else {
RETURN_IF_ERROR(BuildTableInfoUtil::by_orc_field_id_with_name_mapping(
get_scan_params().history_schema_info.front().root_field, orc_type_ptr,
ICEBERG_ORC_ATTRIBUTE, ctx->table_info_node,
supports_iceberg_scan_semantics_v1(&get_scan_params())));
}
std::unordered_set<std::string> partition_col_names;
if (ctx->range->__isset.columns_from_path_keys) {
partition_col_names.insert(ctx->range->columns_from_path_keys.begin(),
ctx->range->columns_from_path_keys.end());
}
// Single pass: classify columns, detect $row_id, handle partition fallback.
bool has_partition_from_path = false;
for (const auto& desc : *ctx->column_descs) {
if (desc.category == ColumnCategory::SYNTHESIZED) {
if (desc.name == BeConsts::ICEBERG_ROWID_COL) {
this->register_synthesized_column_handler(
BeConsts::ICEBERG_ROWID_COL, [this](Block* block, size_t rows) -> Status {
return _fill_iceberg_row_id(block, rows);
});
continue;
} else if (desc.name.starts_with(BeConsts::GLOBAL_ROWID_COL)) {
auto topn_row_id_column_iter = _create_topn_row_id_column_iterator();
this->register_synthesized_column_handler(
desc.name,
[iter = std::move(topn_row_id_column_iter), this, &desc](
Block* block, size_t rows) -> Status {
return fill_topn_row_id(iter, desc.name, block, rows);
});
continue;
}
} else if (desc.category == ColumnCategory::PARTITION_KEY) {
bool has_partition_value = partition_col_names.contains(desc.name);
bool exists_in_file = ctx->table_info_node->children_column_exists(desc.name);
if (!has_partition_value || exists_in_file) {
ctx->column_names.push_back(desc.name);
continue;
}
has_partition_from_path = true;
} else if (desc.category == ColumnCategory::REGULAR) {
ctx->column_names.push_back(desc.name);
} else if (desc.category == ColumnCategory::GENERATED) {
_init_row_lineage_columns();
if (desc.name == ROW_LINEAGE_ROW_ID) {
ctx->column_names.push_back(desc.name);
this->register_generated_column_handler(
ROW_LINEAGE_ROW_ID, [this](Block* block, size_t rows) -> Status {
return _fill_row_lineage_row_id(block, rows);
});
continue;
} else if (desc.name == ROW_LINEAGE_LAST_UPDATED_SEQ_NUMBER) {
ctx->column_names.push_back(desc.name);
this->register_generated_column_handler(
ROW_LINEAGE_LAST_UPDATED_SEQ_NUMBER,
[this](Block* block, size_t rows) -> Status {
return _fill_row_lineage_last_updated_sequence_number(block, rows);
});
continue;
}
}
}
if (has_partition_from_path) {
RETURN_IF_ERROR(_extract_partition_values(*ctx->range, ctx->tuple_descriptor,
_fill_partition_values,
&_fill_partition_value_is_null));
}
_all_required_col_names = ctx->column_names;
// Create column IDs from ORC type
auto column_id_result =
_create_column_ids(orc_type_ptr, ctx->tuple_descriptor, ctx->table_info_node);
ctx->column_ids = std::move(column_id_result.column_ids);
ctx->filter_column_ids = std::move(column_id_result.filter_column_ids);
// Build field_id -> block_column_name mapping for equality delete filtering.
for (const auto* slot : ctx->tuple_descriptor->slots()) {
_id_to_block_column_name.emplace(slot->col_unique_id(), slot->col_name());
}
// Process delete files (must happen before _do_init_reader so expand col IDs are included)
RETURN_IF_ERROR(_init_row_filters());
// Add expand column IDs for equality delete and remap expand column names
// (matching master's behavior with __equality_delete_column__ prefix)
const static std::string EQ_DELETE_PRE = "__equality_delete_column__";
bool all_file_columns_have_field_ids = true;
for (uint64_t i = 0; i < orc_type_ptr->getSubtypeCount(); ++i) {
const orc::Type* sub_type = orc_type_ptr->getSubtype(i);
if (!sub_type->hasAttributeKey(ICEBERG_ORC_ATTRIBUTE)) {
all_file_columns_have_field_ids = false;
}
}
const bool use_field_ids_for_hidden_keys =
supports_iceberg_scan_semantics_v1(&get_scan_params())
? orc_subtree_has_iceberg_id(orc_type_ptr, ICEBERG_ORC_ATTRIBUTE)
: all_file_columns_have_field_ids;
const auto find_file_column_by_name = [&](const std::string& name) -> const orc::Type* {
for (uint64_t j = 0; j < orc_type_ptr->getSubtypeCount(); ++j) {
if (iequal(orc_type_ptr->getFieldName(j), name)) {
return orc_type_ptr->getSubtype(j);
}
}
return nullptr;
};
std::vector<std::string> new_expand_col_names;
DORIS_CHECK(_expand_col_names.size() == _expand_col_field_ids.size());
DORIS_CHECK(_expand_col_names.size() == _expand_columns.size());
for (size_t i = 0; i < _expand_col_names.size(); ++i) {
const auto& old_name = _expand_col_names[i];
const int32_t field_id = _expand_col_field_ids[i];
const orc::Type* file_column = nullptr;
OrcEqualityFieldPath file_path;
bool complete_file_path = false;
if (use_field_ids_for_hidden_keys) {
complete_file_path =
find_orc_equality_field_path_by_id(orc_type_ptr, field_id, &file_path);
if (!complete_file_path && supports_iceberg_scan_semantics_v2(&get_scan_params())) {
const auto table_path = _find_schema_field_path(field_id);
if (!table_path.empty()) {
complete_file_path = find_orc_equality_field_prefix_by_id_path(
orc_type_ptr, table_path, &file_path);
}
}
if (!file_path.fields.empty()) {
file_column = file_path.fields.front();
}
} else {
const auto table_path = _find_schema_field_path(field_id);
if (!table_path.empty()) {
complete_file_path = find_orc_equality_field_prefix_by_name_path(
orc_type_ptr, table_path, old_name, &file_path);
if (!file_path.fields.empty()) {
file_column = file_path.fields.front();
}
} else {
file_column = find_file_column_by_name(old_name);
complete_file_path = file_column != nullptr;
}
}
std::string file_col_name = old_name;
std::string leaf_name = old_name;
if (!file_path.fields.empty()) {
file_col_name = file_path.names.front();
leaf_name = file_path.names.back();
} else if (file_column != nullptr) {
for (uint64_t j = 0; j < orc_type_ptr->getSubtypeCount(); ++j) {
if (orc_type_ptr->getSubtype(j) == file_column) {
file_col_name = orc_type_ptr->getFieldName(j);
leaf_name = file_col_name;
break;
}
}
}
std::string table_col_name = EQ_DELETE_PRE + std::to_string(field_id) + "_" + leaf_name;
if (field_id >= 0) {
_id_to_block_column_name[field_id] = table_col_name;
}
_expand_columns[i].name = table_col_name;
if (file_column == nullptr) {
RETURN_IF_ERROR(_register_missing_equality_delete_column(field_id, table_col_name,
_expand_columns[i].type));
// The old data file predates this equality key. Keep it in the expand block so the
// synthesized-column hook can materialize its logical initial default before ORC's
// block-size checks. Adding it to column_names/table_info_node would mark it as an
// existing ORC child and make OrcReader read a column that is not present in the file.
new_expand_col_names.push_back(table_col_name);
continue;
}
new_expand_col_names.push_back(table_col_name);
if (!complete_file_path) {
ColumnPtr missing_value;
RETURN_IF_ERROR(_create_missing_equality_delete_value(
field_id, _expand_columns[i].type, file_path.fields.size(), &missing_value));
_nested_equality_delete_columns.push_back({
.field_id = field_id,
.block_name = table_col_name,
.leaf_type = _expand_columns[i].type,
.child_indexes = file_path.child_indexes,
.missing_value = std::move(missing_value),
});
_expand_columns[i].type = make_nullable(convert_to_doris_type(file_column));
_expand_columns[i].column = _expand_columns[i].type->create_column();
} else if (!file_path.child_indexes.empty()) {
_nested_equality_delete_columns.push_back({
.field_id = field_id,
.block_name = table_col_name,
.leaf_type = _expand_columns[i].type,
.child_indexes = file_path.child_indexes,
.missing_value = nullptr,
});
_expand_columns[i].type = make_nullable(convert_to_doris_type(file_column));
_expand_columns[i].column = _expand_columns[i].type->create_column();
}
for (uint64_t column_id = file_column->getColumnId();
column_id <= file_column->getMaximumColumnId(); ++column_id) {
ctx->column_ids.insert(column_id);
}
ctx->column_names.push_back(table_col_name);
ctx->table_info_node->add_children(table_col_name, file_col_name,
TableSchemaChangeHelper::ConstNode::get_instance());
}
_expand_col_names = std::move(new_expand_col_names);
return Status::OK();
}
// ============================================================================
// IcebergOrcReader: _create_column_ids
// ============================================================================
ColumnIdResult IcebergOrcReader::_create_column_ids(
const orc::Type* orc_type, const TupleDescriptor* tuple_descriptor,
const std::shared_ptr<TableSchemaChangeHelper::Node>& table_info_node) {
std::unordered_map<int, const orc::Type*> iceberg_id_to_orc_type_map;
for (uint64_t i = 0; i < orc_type->getSubtypeCount(); ++i) {
const auto* orc_sub_type = orc_type->getSubtype(i);
if (!orc_sub_type) {
continue;
}
if (!orc_sub_type->hasAttributeKey(ICEBERG_ORC_ATTRIBUTE)) {
continue;
}
int iceberg_id = std::stoi(orc_sub_type->getAttributeValue(ICEBERG_ORC_ATTRIBUTE));
iceberg_id_to_orc_type_map[iceberg_id] = orc_sub_type;
}
std::set<uint64_t> column_ids;
std::set<uint64_t> filter_column_ids;
auto process_access_paths = [](const orc::Type* orc_field,
const std::vector<TColumnAccessPath>& access_paths,
std::set<uint64_t>& out_ids) {
process_nested_access_paths(
orc_field, access_paths, out_ids,
[](const orc::Type* type) { return type->getColumnId(); },
[](const orc::Type* type) { return type->getMaximumColumnId(); },
IcebergOrcNestedColumnUtils::extract_nested_column_ids);
};
// The Iceberg schema-mapping root is a StructNode whose registered children are the real
// table columns. When present, resolve each column by name through it so the column-id set
// stays consistent with the schema-mapping decision (BY_ID or BY_NAME/name-mapping);
// otherwise fall back to matching by Iceberg field id.
const auto* struct_node =
dynamic_cast<const TableSchemaChangeHelper::StructNode*>(table_info_node.get());
for (const auto* slot : tuple_descriptor->slots()) {
const orc::Type* orc_field = nullptr;
if (struct_node != nullptr) {
// Synthesized/metadata slots (e.g. the TopN global row-id or the $row_id column) are
// never registered as children, so check membership before querying: calling
// children_column_exists() on an unregistered name DCHECK-aborts in debug builds and
// throws std::out_of_range from .at() in release builds.
if (struct_node->get_children().contains(slot->col_name()) &&
struct_node->children_column_exists(slot->col_name())) {
// Select the physical child resolved by the shared schema-mapping pass. Hidden
// equality keys and projected columns must obey the same BY_NAME decision for
// partial-id ORC files.
const auto& file_column_name =
struct_node->children_file_column_name(slot->col_name());
for (uint64_t i = 0; i < orc_type->getSubtypeCount(); ++i) {
if (orc_type->getFieldName(i) == file_column_name) {
orc_field = orc_type->getSubtype(i);
break;
}
}
DORIS_CHECK(orc_field != nullptr);
}
} else {
auto it = iceberg_id_to_orc_type_map.find(slot->col_unique_id());
if (it != iceberg_id_to_orc_type_map.end()) {
orc_field = it->second;
}
}
if (orc_field == nullptr) {
continue;
}
if ((slot->col_type() != TYPE_STRUCT && slot->col_type() != TYPE_ARRAY &&
slot->col_type() != TYPE_MAP)) {
column_ids.insert(orc_field->getColumnId());
if (slot->is_predicate()) {
filter_column_ids.insert(orc_field->getColumnId());
}
continue;
}
const auto& all_access_paths = slot->all_access_paths();
process_access_paths(orc_field, all_access_paths, column_ids);
const auto& predicate_access_paths = slot->predicate_access_paths();
if (!predicate_access_paths.empty()) {
process_access_paths(orc_field, predicate_access_paths, filter_column_ids);
}
}
return {std::move(column_ids), std::move(filter_column_ids)};
}
// ============================================================================
// IcebergOrcReader: _read_position_delete_file
// ============================================================================
Status IcebergOrcReader::_read_position_delete_file(const TFileRangeDesc* delete_range,
DeleteFile* position_delete) {
OrcReader orc_delete_reader(get_profile(), get_state(), get_scan_params(), *delete_range,
READ_DELETE_FILE_BATCH_SIZE, get_state()->timezone(), get_io_ctx(),
_meta_cache);
OrcInitContext delete_ctx;
delete_ctx.column_names = delete_file_col_names;
delete_ctx.col_name_to_block_idx =
const_cast<std::unordered_map<std::string, uint32_t>*>(&DELETE_COL_NAME_TO_BLOCK_IDX);
RETURN_IF_ERROR(orc_delete_reader.init_reader(&delete_ctx));
bool eof = false;
DataTypePtr data_type_file_path {new DataTypeString};
DataTypePtr data_type_pos {new DataTypeInt64};
while (!eof) {
Block block = {{data_type_file_path, ICEBERG_FILE_PATH}, {data_type_pos, ICEBERG_ROW_POS}};
size_t read_rows = 0;
RETURN_IF_ERROR(orc_delete_reader.get_next_block(&block, &read_rows, &eof));
RETURN_IF_ERROR(_gen_position_delete_file_range(block, position_delete, read_rows, false));
}
return Status::OK();
}
} // namespace doris