Skip to content
Open
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
4 changes: 3 additions & 1 deletion be/src/cloud/cloud_schema_change_job.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -244,6 +244,9 @@ Status CloudSchemaChangeJob::process_alter_tablet(const TAlterTabletReqV2& reque
_new_tablet_schema = _new_tablet->tablet_schema();

ReadSchemaSPtr read_schema = std::make_shared<ReadSchema>(_base_tablet_schema->columns());
RETURN_IF_ERROR(read_schema->init_from_tablet_schema(*_base_tablet_schema,
/*merge_by_sequence_mapping=*/false,
/*map_row_binlog_columns=*/false));

// delete handlers to filter out deleted rows
DeleteHandler delete_handler;
Expand All @@ -262,7 +265,6 @@ Status CloudSchemaChangeJob::process_alter_tablet(const TAlterTabletReqV2& reque
// reader_context is stack variables, it's lifetime MUST keep the same with rs_readers
RowsetReaderContext reader_context;
reader_context.reader_type = ReaderType::READER_ALTER_TABLE;
reader_context.tablet_schema = _base_tablet_schema;
reader_context.need_ordered_result = true;
reader_context.delete_handler = &delete_handler;
reader_context.read_schema = read_schema;
Expand Down
19 changes: 10 additions & 9 deletions be/src/exec/common/variant_util.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -117,7 +117,7 @@ PathInData make_full_subcolumn_path(const TabletColumnPtr& parent_column, std::s
return builder.append(parent_column->name_lower_case(), false).append("", false).build();
}

void append_empty_key_subcolumn_from_stats(TabletSchema::PathsSetInfo& paths_set_info,
void append_empty_key_subcolumn_from_stats(VariantCompactionPaths& paths_set_info,
const TabletColumnPtr& parent_column,
TabletSchemaSPtr& output_schema) {
if (!paths_set_info.sub_path_set.contains("") || paths_set_info.sparse_path_set.contains("") ||
Expand Down Expand Up @@ -1002,7 +1002,7 @@ Status VariantCompactionUtil::aggregate_variant_extended_info(
// get the subpaths and sparse paths for the variant column
void VariantCompactionUtil::get_subpaths(int32_t max_subcolumns_count,
const PathToNoneNullValues& stats,
TabletSchema::PathsSetInfo& paths_set_info) {
VariantCompactionPaths& paths_set_info) {
// max_subcolumns_count is 0 means no limit
if (max_subcolumns_count > 0 && stats.size() > max_subcolumns_count) {
std::vector<std::pair<size_t, std::string_view>> paths_with_sizes;
Expand Down Expand Up @@ -1148,7 +1148,7 @@ Status VariantCompactionUtil::check_path_stats(const std::vector<RowsetSharedPtr
Status VariantCompactionUtil::get_compaction_typed_columns(
const TabletSchemaSPtr& target, const std::unordered_set<std::string>& typed_paths,
const TabletColumnPtr parent_column, TabletSchemaSPtr& output_schema,
TabletSchema::PathsSetInfo& paths_set_info) {
VariantCompactionPaths& paths_set_info) {
if (parent_column->variant_enable_typed_paths_to_sparse()) {
return Status::OK();
}
Expand All @@ -1169,7 +1169,7 @@ Status VariantCompactionUtil::get_compaction_typed_columns(
Status VariantCompactionUtil::get_compaction_nested_columns(
const std::unordered_set<PathInData, PathInData::Hash>& nested_paths,
const PathToDataTypes& path_to_data_types, const TabletColumnPtr parent_column,
TabletSchemaSPtr& output_schema, TabletSchema::PathsSetInfo& paths_set_info) {
TabletSchemaSPtr& output_schema, VariantCompactionPaths& paths_set_info) {
const auto& parent_indexes = output_schema->inverted_indexs(parent_column->unique_id());
for (const auto& path : nested_paths) {
const auto& find_data_types = path_to_data_types.find(path);
Expand Down Expand Up @@ -1200,7 +1200,7 @@ Status VariantCompactionUtil::get_compaction_nested_columns(
}

void VariantCompactionUtil::get_compaction_subcolumns_from_subpaths(
TabletSchema::PathsSetInfo& paths_set_info, const TabletColumnPtr parent_column,
VariantCompactionPaths& paths_set_info, const TabletColumnPtr parent_column,
const TabletSchemaSPtr& target, const PathToDataTypes& path_to_data_types,
const std::unordered_set<std::string>& sparse_paths, TabletSchemaSPtr& output_schema) {
auto& path_set = paths_set_info.sub_path_set;
Expand Down Expand Up @@ -1265,7 +1265,7 @@ void VariantCompactionUtil::get_compaction_subcolumns_from_subpaths(
}

void VariantCompactionUtil::get_compaction_subcolumns_from_data_types(
TabletSchema::PathsSetInfo& paths_set_info, const TabletColumnPtr parent_column,
VariantCompactionPaths& paths_set_info, const TabletColumnPtr parent_column,
const TabletSchemaSPtr& target, const PathToDataTypes& path_to_data_types,
TabletSchemaSPtr& output_schema) {
const auto& parent_indexes = target->inverted_indexs(parent_column->unique_id());
Expand Down Expand Up @@ -1305,7 +1305,8 @@ void VariantCompactionUtil::get_compaction_subcolumns_from_data_types(
// ordinary extracted subcolumns. NG typed paths still use get_compaction_typed_columns(), keeping
// typed-column rules out of the NG-specific regular-path filtering.
Status VariantCompactionUtil::get_extended_compaction_schema(
const std::vector<RowsetSharedPtr>& rowsets, TabletSchemaSPtr& target) {
const std::vector<RowsetSharedPtr>& rowsets, TabletSchemaSPtr& target,
VariantCompactionPathsMap& paths) {
std::unordered_map<int32_t, VariantExtendedInfo> uid_to_variant_extended_info;
const bool needs_variant_extended_info =
std::ranges::any_of(target->columns(), [](const TabletColumnPtr& column) {
Expand All @@ -1322,7 +1323,7 @@ Status VariantCompactionUtil::get_extended_compaction_schema(
// build the output schema
TabletSchemaSPtr output_schema = std::make_shared<TabletSchema>();
output_schema->shawdow_copy_without_columns(*target);
std::unordered_map<int32_t, TabletSchema::PathsSetInfo> uid_to_paths_set_info;
VariantCompactionPathsMap uid_to_paths_set_info;
const auto ng_root_uids =
collect_nested_group_compaction_root_uids(target, uid_to_variant_extended_info);
for (const TabletColumnPtr& column : target->columns()) {
Expand Down Expand Up @@ -1417,7 +1418,7 @@ Status VariantCompactionUtil::get_extended_compaction_schema(

target = output_schema;
// used to merge & filter path to sparse column during reading in compaction
target->set_path_set_info(std::move(uid_to_paths_set_info));
paths = std::move(uid_to_paths_set_info);
VLOG_DEBUG << "dump schema " << target->dump_full_schema();
return Status::OK();
}
Expand Down
24 changes: 14 additions & 10 deletions be/src/exec/common/variant_util.h
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@
#include "core/string_ref.h"
#include "core/types.h"
#include "exprs/aggregate/aggregate_function.h"
#include "storage/segment/variant/variant_compaction_paths.h"
#include "storage/tablet/tablet_fwd.h"
#include "storage/tablet/tablet_schema.h"

Expand Down Expand Up @@ -186,7 +187,7 @@ class VariantCompactionUtil {
public:
// get the subpaths and sparse paths for the variant column
static void get_subpaths(int32_t max_subcolumns_count, const PathToNoneNullValues& path_stats,
TabletSchema::PathsSetInfo& paths_set_info);
VariantCompactionPaths& paths_set_info);

// collect extended info from the variant column
static Status aggregate_variant_extended_info(
Expand All @@ -198,9 +199,11 @@ class VariantCompactionUtil {
const RowsetSharedPtr& rs,
std::unordered_map<int32_t, PathToNoneNullValues>* uid_to_path_stats);

// Build the temporary schema for compaction, this will reduce the memory usage of compacting variant columns
// Build the temporary schema for compaction, this will reduce the memory usage of compacting
// variant columns. `paths` receives that schema's variant path layout.
static Status get_extended_compaction_schema(const std::vector<RowsetSharedPtr>& rowsets,
TabletSchemaSPtr& target);
TabletSchemaSPtr& target,
VariantCompactionPathsMap& paths);

// Used to collect all the subcolumns types of variant column from rowsets
static TabletSchemaSPtr calculate_variant_extended_schema(
Expand All @@ -218,25 +221,26 @@ class VariantCompactionUtil {
size_t num_rows);

static void get_compaction_subcolumns_from_subpaths(
TabletSchema::PathsSetInfo& paths_set_info, const TabletColumnPtr parent_column,
VariantCompactionPaths& paths_set_info, const TabletColumnPtr parent_column,
const TabletSchemaSPtr& target, const PathToDataTypes& path_to_data_types,
const std::unordered_set<std::string>& sparse_paths, TabletSchemaSPtr& output_schema);

static void get_compaction_subcolumns_from_data_types(
TabletSchema::PathsSetInfo& paths_set_info, const TabletColumnPtr parent_column,
const TabletSchemaSPtr& target, const PathToDataTypes& path_to_data_types,
TabletSchemaSPtr& output_schema);
static void get_compaction_subcolumns_from_data_types(VariantCompactionPaths& paths_set_info,
const TabletColumnPtr parent_column,
const TabletSchemaSPtr& target,
const PathToDataTypes& path_to_data_types,
TabletSchemaSPtr& output_schema);

static Status get_compaction_typed_columns(const TabletSchemaSPtr& target,
const std::unordered_set<std::string>& typed_paths,
const TabletColumnPtr parent_column,
TabletSchemaSPtr& output_schema,
TabletSchema::PathsSetInfo& paths_set_info);
VariantCompactionPaths& paths_set_info);

static Status get_compaction_nested_columns(
const std::unordered_set<PathInData, PathInData::Hash>& nested_paths,
const PathToDataTypes& path_to_data_types, const TabletColumnPtr parent_column,
TabletSchemaSPtr& output_schema, TabletSchema::PathsSetInfo& paths_set_info);
TabletSchemaSPtr& output_schema, VariantCompactionPaths& paths_set_info);
};

} // namespace doris::variant_util
29 changes: 17 additions & 12 deletions be/src/exec/rowid_fetcher.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -138,15 +138,11 @@ struct IteratorItem {
StorageReadOptions storage_read_options;
};

static void set_slot_access_paths(const SlotDescriptor& slot, const TabletSchema& schema,
// read_column is what slot resolves to, or its variant parent column for a subpath slot.
static void set_slot_access_paths(const SlotDescriptor& slot, const TabletColumn& read_column,
StorageReadOptions& storage_read_options) {
int32_t unique_id = slot.col_unique_id();
const int field_index =
unique_id >= 0 ? schema.field_index(unique_id) : schema.field_index(slot.col_name());
if (field_index >= 0) {
const auto& column = schema.column(field_index);
unique_id = column.unique_id() >= 0 ? column.unique_id() : column.parent_unique_id();
}
const int32_t unique_id =
read_column.unique_id() >= 0 ? read_column.unique_id() : read_column.parent_unique_id();
if (unique_id < 0) {
return;
}
Expand Down Expand Up @@ -1043,10 +1039,19 @@ Status RowIdStorageReader::read_doris_format_row(
iterator_item.storage_read_options.io_ctx.file_cache_miss_policy =
file_cache_miss_policy;
}
set_slot_access_paths(slots[x], full_read_schema, iterator_item.storage_read_options);
RETURN_IF_ERROR(segment->seek_and_read_by_rowid(
full_read_schema, &slots[x], row_ids, column,
iterator_item.storage_read_options, iterator_item.iterator));
int32_t index = slots[x].col_unique_id() >= 0
? full_read_schema.field_index(slots[x].col_unique_id())
: full_read_schema.field_index(slots[x].col_name());
if (index < 0) {
return Status::InternalError(
"field name is invalid. field={}, field_name_to_index={}",
slots[x].col_name(), full_read_schema.get_all_field_names());
}
const auto& read_column = full_read_schema.column(index);
set_slot_access_paths(slots[x], read_column, iterator_item.storage_read_options);
RETURN_IF_ERROR(segment->seek_and_read_by_rowid(read_column, &slots[x], row_ids, column,
iterator_item.storage_read_options,
iterator_item.iterator));
}
}
return Status::OK();
Expand Down
13 changes: 11 additions & 2 deletions be/src/service/point_query_executor.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -595,18 +595,27 @@ Status PointQueryExecutor::_lookup_row_data() {
return seg->id() == row_loc.segment_id;
});
const auto& segment = *it;
const auto tablet_schema = _tablet->tablet_schema();
for (int cid : _reusable->missing_col_uids()) {
int pos = _reusable->get_col_uid_to_idx().at(cid);
std::vector<segment_v2::rowid_t> row_ids {
static_cast<segment_v2::rowid_t>(row_loc.row_id)};
auto& column = result_columns[pos];
std::unique_ptr<ColumnIterator> iter;
SlotDescriptor* slot = _reusable->tuple_desc()->slots()[pos];
int32_t index = slot->col_unique_id() >= 0
? tablet_schema->field_index(slot->col_unique_id())
: tablet_schema->field_index(slot->col_name());
if (index < 0) {
return Status::InternalError(
"field name is invalid. field={}, field_name_to_index={}",
slot->col_name(), tablet_schema->get_all_field_names());
}
StorageReadOptions storage_read_options;
storage_read_options.stats = &_read_stats;
storage_read_options.io_ctx = io_ctx;
RETURN_IF_ERROR(segment->seek_and_read_by_rowid(*_tablet->tablet_schema(), slot,
row_ids, column,
RETURN_IF_ERROR(segment->seek_and_read_by_rowid(tablet_schema->column(index),
slot, row_ids, column,
storage_read_options, iter));
}
}
Expand Down
10 changes: 8 additions & 2 deletions be/src/storage/compaction/compaction.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -574,8 +574,10 @@ Status CompactionMixin::build_basic_info(bool is_ordered_compaction) {
// so get_extended_compaction_schema will extended the schema for variant columns
// for ordered compaction, we don't need to extend the schema for variant columns
if (_enable_vertical_compact_variant_subcolumns && !is_ordered_compaction) {
auto paths = std::make_shared<VariantCompactionPathsMap>();
RETURN_IF_ERROR(variant_util::VariantCompactionUtil::get_extended_compaction_schema(
_input_rowsets, _cur_tablet_schema));
_input_rowsets, _cur_tablet_schema, *paths));
_cur_variant_compaction_paths = std::move(paths);
}
return Status::OK();
}
Expand Down Expand Up @@ -1787,6 +1789,7 @@ Status CompactionMixin::construct_output_rowset_writer(RowsetWriterContext& ctx)
ctx.rowset_state = VISIBLE;
ctx.segments_overlap = _trigger_quick_merge_by_binlog ? OVERLAPPING : NONOVERLAPPING;
ctx.tablet_schema = _cur_tablet_schema;
ctx.variant_compaction_paths = _cur_variant_compaction_paths;
ctx.newest_write_timestamp = _newest_write_timestamp;
ctx.write_type = DataWriteType::TYPE_COMPACTION;
ctx.compaction_type = compaction_type();
Expand Down Expand Up @@ -2092,8 +2095,10 @@ Status CloudCompactionMixin::build_basic_info() {
// if enable_vertical_compact_variant_subcolumns is true, we need to compact the variant subcolumns in seperate column groups
// so get_extended_compaction_schema will extended the schema for variant columns
if (_enable_vertical_compact_variant_subcolumns) {
auto paths = std::make_shared<VariantCompactionPathsMap>();
RETURN_IF_ERROR(variant_util::VariantCompactionUtil::get_extended_compaction_schema(
_input_rowsets, _cur_tablet_schema));
_input_rowsets, _cur_tablet_schema, *paths));
_cur_variant_compaction_paths = std::move(paths);
}
return Status::OK();
}
Expand Down Expand Up @@ -2370,6 +2375,7 @@ Status CloudCompactionMixin::construct_output_rowset_writer(RowsetWriterContext&
ctx.rowset_state = VISIBLE;
ctx.segments_overlap = NONOVERLAPPING;
ctx.tablet_schema = _cur_tablet_schema;
ctx.variant_compaction_paths = _cur_variant_compaction_paths;
ctx.newest_write_timestamp = _newest_write_timestamp;
ctx.write_type = DataWriteType::TYPE_COMPACTION;
ctx.compaction_type = compaction_type();
Expand Down
2 changes: 2 additions & 0 deletions be/src/storage/compaction/compaction.h
Original file line number Diff line number Diff line change
Expand Up @@ -208,6 +208,8 @@ class Compaction {
int64_t _newest_write_timestamp {-1};
std::unique_ptr<RowIdConversion> _rowid_conversion = nullptr;
TabletSchemaSPtr _cur_tablet_schema;
// The variant path layout of _cur_tablet_schema; empty unless it is an extended schema.
VariantCompactionPathsSPtr _cur_variant_compaction_paths;

std::unique_ptr<RuntimeProfile> _profile;

Expand Down
Loading
Loading