diff --git a/be/src/cloud/cloud_cumulative_compaction.cpp b/be/src/cloud/cloud_cumulative_compaction.cpp index cd61ea9c4261d2..f4ce0f81a02f65 100644 --- a/be/src/cloud/cloud_cumulative_compaction.cpp +++ b/be/src/cloud/cloud_cumulative_compaction.cpp @@ -17,6 +17,8 @@ #include "cloud/cloud_cumulative_compaction.h" +#include +#include #include #include @@ -26,12 +28,17 @@ #include "cloud/config.h" #include "common/config.h" #include "common/logging.h" +#include "common/metrics/doris_metrics.h" #include "common/status.h" #include "cpp/sync_point.h" #include "service/backend_options.h" #include "storage/compaction/compaction.h" #include "storage/compaction/cumulative_compaction_policy.h" #include "storage/compaction/cumulative_compaction_time_series_policy.h" +#include "storage/merger.h" +#include "storage/rowset/rowset_reader.h" +#include "storage/rowset/rowset_writer.h" +#include "storage/tablet/tablet_schema.h" #include "util/debug_points.h" #include "util/trace.h" #include "util/uuid_generator.h" @@ -40,6 +47,74 @@ namespace doris { #include "common/compile_check_begin.h" using namespace ErrorCode; +namespace cloud { + +bool is_single_rowset_compaction_candidate(const RowsetSharedPtr& rowset) { + const auto& rowset_meta = rowset->rowset_meta(); + const int64_t overlap_unit_count = + rowset_meta->segments_overlap() == NONOVERLAPPING_WITHIN_GROUP + ? static_cast(rowset_meta->segment_group_sizes().size()) + : rowset->num_segments(); + return !rowset_meta->has_delete_predicate() && rowset_meta->is_segments_overlapping() && + overlap_unit_count >= config::cloud_single_rowset_compaction_min_segments; +} + +bool should_use_single_rowset_grouped_compaction(const std::vector& input_rowsets, + const TabletSchema& tablet_schema, + std::string_view compaction_policy) { + return compaction_policy == CUMULATIVE_SIZE_BASED_POLICY && + tablet_schema.num_key_columns() > 0 && tablet_schema.cluster_key_uids().empty() && + config::enable_cloud_single_rowset_compaction && input_rowsets.size() == 1 && + is_single_rowset_compaction_candidate(input_rowsets.front()); +} + +std::vector build_segment_group_merge_ranges(const RowsetMeta& rowset_meta, + int64_t segment_group_size) { + DORIS_CHECK_GT(segment_group_size, 1); + DORIS_CHECK_GT(rowset_meta.num_segments(), 0); + + std::vector ranges; + if (rowset_meta.segments_overlap() == NONOVERLAPPING_WITHIN_GROUP) { + const auto& input_segment_group_sizes = rowset_meta.segment_group_sizes(); + const int64_t input_group_count = cast_set(input_segment_group_sizes.size()); + DORIS_CHECK_GT(input_group_count, 0); + ranges.reserve(cast_set((input_group_count + segment_group_size - 1) / + segment_group_size)); + + int64_t segment_end = 0; + for (int64_t group_start = 0; group_start < input_group_count; + group_start += segment_group_size) { + const int64_t group_end = std::min(group_start + segment_group_size, input_group_count); + const int64_t segment_start = segment_end; + for (int64_t group_index = group_start; group_index < group_end; ++group_index) { + const int32_t input_group_size = + input_segment_group_sizes.Get(cast_set(group_index)); + DORIS_CHECK_GT(input_group_size, 0); + segment_end += input_group_size; + } + + ranges.push_back({.segment_start = segment_start, + .segment_end = segment_end, + .merge_way_num = group_end - group_start}); + } + DORIS_CHECK_EQ(segment_end, rowset_meta.num_segments()); + } else { + ranges.reserve(cast_set((rowset_meta.num_segments() + segment_group_size - 1) / + segment_group_size)); + for (int64_t segment_start = 0; segment_start < rowset_meta.num_segments(); + segment_start += segment_group_size) { + const int64_t segment_end = + std::min(segment_start + segment_group_size, rowset_meta.num_segments()); + ranges.push_back({.segment_start = segment_start, + .segment_end = segment_end, + .merge_way_num = segment_end - segment_start}); + } + } + return ranges; +} + +} // namespace cloud + bvar::Adder cumu_output_size("cumu_compaction", "output_size"); bvar::LatencyRecorder g_cu_compaction_hold_delete_bitmap_lock_time_ms( "cu_compaction_hold_delete_bitmap_lock_time_ms"); @@ -250,6 +325,18 @@ Status CloudCumulativeCompaction::execute_compact() { return st; } +bool CloudCumulativeCompaction::should_calculate_new_cumulative_point( + int64_t input_cumulative_point) const { + if (!_single_rowset_compaction_segment_group_size.has_value()) { + return true; + } + + DORIS_CHECK_EQ(_input_rowsets.size(), 1); + DORIS_CHECK(_output_rowset != nullptr); + return _input_rowsets.front()->start_version() == input_cumulative_point && + _output_rowset->rowset_meta()->segments_overlap() == NONOVERLAPPING; +} + Status CloudCumulativeCompaction::modify_rowsets() { // calculate new cumulative point int64_t input_cumulative_point; @@ -267,17 +354,19 @@ Status CloudCumulativeCompaction::modify_rowsets() { } auto compaction_policy = cloud_tablet()->tablet_meta()->compaction_policy(); int64_t new_cumulative_point = input_cumulative_point; - if (!_enable_parallel_cumu_compaction && input_tablet_state == TABLET_NOTREADY && - _output_rowset->start_version() > input_cumulative_point) { - // Historical rowsets are absent from a schema-change target until conversion finishes. - DORIS_CHECK_LE(input_cumulative_point, input_alter_version); - DORIS_CHECK_GT(_output_rowset->start_version(), input_alter_version); - } else if (!_enable_parallel_cumu_compaction || - _output_rowset->start_version() == input_cumulative_point) { - new_cumulative_point = - _engine.cumu_compaction_policy(compaction_policy) - ->new_cumulative_point(cloud_tablet(), _output_rowset, _last_delete_version, - input_cumulative_point); + if (should_calculate_new_cumulative_point(input_cumulative_point)) { + if (!_enable_parallel_cumu_compaction && input_tablet_state == TABLET_NOTREADY && + _output_rowset->start_version() > input_cumulative_point) { + // Historical rowsets are absent from a schema-change target until conversion finishes. + DORIS_CHECK_LE(input_cumulative_point, input_alter_version); + DORIS_CHECK_GT(_output_rowset->start_version(), input_alter_version); + } else if (!_enable_parallel_cumu_compaction || + _output_rowset->start_version() == input_cumulative_point) { + new_cumulative_point = + _engine.cumu_compaction_policy(compaction_policy) + ->new_cumulative_point(cloud_tablet(), _output_rowset, + _last_delete_version, input_cumulative_point); + } } // commit compaction job cloud::TabletJobInfoPB job; @@ -608,6 +697,7 @@ Status CloudCumulativeCompaction::advance_cumulative_point_before_pick( Status CloudCumulativeCompaction::pick_rowsets_to_compact() { _input_rowsets.clear(); + _single_rowset_compaction_segment_group_size.reset(); int64_t min_conflict_version = _min_conflict_version; int64_t max_conflict_version = _max_conflict_version; @@ -646,6 +736,20 @@ Status CloudCumulativeCompaction::pick_rowsets_to_compact() { config::cumulative_compaction_min_deltas, &_input_rowsets, &_last_delete_version, &compaction_score); + const int64_t segment_group_size = + config::cloud_single_rowset_compaction_segment_group_size; + if (config::enable_cloud_single_rowset_compaction && segment_group_size > 1) { + for (const auto& rowset : _input_rowsets) { + if (cloud::should_use_single_rowset_grouped_compaction( + {rowset}, *cloud_tablet()->tablet_schema(), compaction_policy)) { + auto grouped_input_rowset = rowset; + _input_rowsets = {std::move(grouped_input_rowset)}; + _single_rowset_compaction_segment_group_size = segment_group_size; + return Status::OK(); + } + } + } + if (_input_rowsets.empty()) { return Status::Error( "no suitable versions: input rowsets empty"); @@ -709,6 +813,89 @@ Status CloudCumulativeCompaction::pick_rowsets_to_compact() { return Status::OK(); } +Status CloudCumulativeCompaction::prepare_merge_input_rowsets(MergeInputRowsetsResult* result) { + if (!_single_rowset_compaction_segment_group_size.has_value()) { + return Status::OK(); + } + + const int64_t segment_group_size = *_single_rowset_compaction_segment_group_size; + DORIS_CHECK_GT(segment_group_size, 1); + result->is_segment_grouped = true; + result->segment_group_size = segment_group_size; + return Status::OK(); +} + +Status CloudCumulativeCompaction::do_merge_input_rowsets( + const std::vector& input_rs_readers, + MergeInputRowsetsResult* result) { + if (!result->is_segment_grouped) { + return Compaction::do_merge_input_rowsets(input_rs_readers, result); + } + + const int64_t segment_group_size = result->segment_group_size; + const auto& input_rowset = _input_rowsets.front(); + const auto segment_ranges = cloud::build_segment_group_merge_ranges( + *input_rowset->rowset_meta(), segment_group_size); + for (size_t range_index = 0; range_index < segment_ranges.size(); ++range_index) { + const auto& range = segment_ranges[range_index]; + std::vector output_segment_num_rows; + RETURN_IF_ERROR(_output_rs_writer->get_segment_num_rows(&output_segment_num_rows)); + const size_t output_segment_start = output_segment_num_rows.size(); + + RowsetReaderSharedPtr rs_reader; + RETURN_IF_ERROR(input_rowset->create_reader(&rs_reader)); + std::vector group_readers; + group_readers.push_back(std::move(rs_reader)); + + Merger::Statistics group_stats; + group_stats.rowid_conversion = _stats.rowid_conversion; + RETURN_IF_ERROR(execute_merge(group_readers, range.merge_way_num, &group_stats, + std::make_pair(range.segment_start, range.segment_end), + {.total_ranges = cast_set(segment_ranges.size()), + .range_index = cast_set(range_index)})); + + _stats.output_rows += group_stats.output_rows; + _stats.merged_rows += group_stats.merged_rows; + _stats.filtered_rows += group_stats.filtered_rows; + _stats.bytes_read_from_local += group_stats.bytes_read_from_local; + _stats.bytes_read_from_remote += group_stats.bytes_read_from_remote; + _stats.cached_bytes_total += group_stats.cached_bytes_total; + _stats.cloud_local_read_time += group_stats.cloud_local_read_time; + _stats.cloud_remote_read_time += group_stats.cloud_remote_read_time; + + RETURN_IF_ERROR(_output_rs_writer->get_segment_num_rows(&output_segment_num_rows)); + DORIS_CHECK_GE(output_segment_num_rows.size(), output_segment_start); + const int32_t output_group_size = + cast_set(output_segment_num_rows.size() - output_segment_start); + if (output_group_size > 0) { + result->output_segment_group_sizes.push_back(output_group_size); + } + } + return Status::OK(); +} + +void CloudCumulativeCompaction::update_output_rowset_after_build( + const MergeInputRowsetsResult& result) { + if (!result.is_segment_grouped) { + return; + } + if (result.output_segment_group_sizes.size() > 1) { + _output_rowset->rowset_meta()->set_segments_overlap(NONOVERLAPPING_WITHIN_GROUP); + _output_rowset->rowset_meta()->set_segment_group_sizes(result.output_segment_group_sizes); + } + + const auto& input_rowset = _input_rowsets.front(); + LOG_INFO("finish single rowset grouped compaction, tablet_id={}, version=[{}-{}]", + _tablet->tablet_id(), input_rowset->start_version(), input_rowset->end_version()) + .tag("job_id", _uuid) + .tag("input_segments", input_rowset->num_segments()) + .tag("segment_group_size", result.segment_group_size) + .tag("output_segments", _output_rowset->num_segments()) + .tag("output_groups", result.output_segment_group_sizes.size()) + .tag("output_segment_group_sizes", + fmt::format("[{}]", fmt::join(result.output_segment_group_sizes, ", "))); +} + void CloudCumulativeCompaction::update_cumulative_point(int64_t input_cumulative_point, int64_t output_cumulative_point) { DORIS_CHECK_LT(input_cumulative_point, output_cumulative_point); diff --git a/be/src/cloud/cloud_cumulative_compaction.h b/be/src/cloud/cloud_cumulative_compaction.h index 28f88c6432f55a..004bd7a5db0b65 100644 --- a/be/src/cloud/cloud_cumulative_compaction.h +++ b/be/src/cloud/cloud_cumulative_compaction.h @@ -20,15 +20,39 @@ #include #include #include +#include +#include #include "cloud/cloud_storage_engine.h" #include "cloud/cloud_tablet.h" #include "storage/compaction/compaction.h" #include "storage/compaction_task_tracker.h" +#include "storage/tablet/tablet_fwd.h" namespace doris { #include "common/compile_check_begin.h" +class RowsetMeta; + +namespace cloud { + +struct SegmentGroupMergeRange { + int64_t segment_start; + int64_t segment_end; + int64_t merge_way_num; +}; + +bool is_single_rowset_compaction_candidate(const RowsetSharedPtr& rowset); + +bool should_use_single_rowset_grouped_compaction(const std::vector& input_rowsets, + const TabletSchema& tablet_schema, + std::string_view compaction_policy); + +std::vector build_segment_group_merge_ranges(const RowsetMeta& rowset_meta, + int64_t segment_group_size); + +} // namespace cloud + class CloudCumulativeCompaction : public CloudCompactionMixin { public: CloudCumulativeCompaction(CloudStorageEngine& engine, CloudTabletSPtr tablet); @@ -54,6 +78,15 @@ class CloudCumulativeCompaction : public CloudCompactionMixin { Status pick_rowsets_to_compact(); + Status prepare_merge_input_rowsets(MergeInputRowsetsResult* result) override; + + Status do_merge_input_rowsets(const std::vector& input_rs_readers, + MergeInputRowsetsResult* result) override; + + void update_output_rowset_after_build(const MergeInputRowsetsResult& result) override; + + bool should_calculate_new_cumulative_point(int64_t input_cumulative_point) const; + std::string_view compaction_name() const override { return "CloudCumulativeCompaction"; } protected: @@ -76,6 +109,7 @@ class CloudCumulativeCompaction : public CloudCompactionMixin { int64_t _cumulative_compaction_cnt = 0; int64_t _picked_cumulative_point = 0; Version _last_delete_version {-1, -1}; + std::optional _single_rowset_compaction_segment_group_size; }; #include "common/compile_check_end.h" diff --git a/be/src/cloud/cloud_cumulative_compaction_policy.cpp b/be/src/cloud/cloud_cumulative_compaction_policy.cpp index 0f7371307c2322..2be0dde84e3cbc 100644 --- a/be/src/cloud/cloud_cumulative_compaction_policy.cpp +++ b/be/src/cloud/cloud_cumulative_compaction_policy.cpp @@ -18,6 +18,7 @@ #include "cloud/cloud_cumulative_compaction_policy.h" #include +#include #include #include #include @@ -213,6 +214,9 @@ int64_t CloudSizeBasedCumulativeCompactionPolicy::pick_input_rowsets( auto rs_begin = input_rowsets->begin(); size_t new_compaction_score = *compaction_score; + const bool can_handle_exhausted_input = + (config::prioritize_query_perf_in_compaction && tablet->keys_type() != DUP_KEYS) || + *compaction_score >= static_cast(max_compaction_score); while (rs_begin != input_rowsets->end()) { auto& rs_meta = (*rs_begin)->rowset_meta(); int64_t current_level = _level_size(rs_meta->total_disk_size()); @@ -222,9 +226,16 @@ int64_t CloudSizeBasedCumulativeCompactionPolicy::pick_input_rowsets( if (current_level <= remain_level) { break; } + + auto next = std::next(rs_begin); + // Keep the last suffix rowset for the singleton checks unless the exhausted-input + // fallback below can select a useful input. + if (next == input_rowsets->end() && !can_handle_exhausted_input) { + break; + } total_size -= rs_meta->total_disk_size(); new_compaction_score -= rs_meta->get_compaction_score(); - ++rs_begin; + rs_begin = next; } if (rs_begin == input_rowsets->end()) { // No suitable level size found in `input_rowsets` if (config::prioritize_query_perf_in_compaction && tablet->keys_type() != DUP_KEYS) { diff --git a/be/src/cloud/cloud_schema_change_job.cpp b/be/src/cloud/cloud_schema_change_job.cpp index 32b4386e0595e8..b63f6ee53bb29a 100644 --- a/be/src/cloud/cloud_schema_change_job.cpp +++ b/be/src/cloud/cloud_schema_change_job.cpp @@ -381,7 +381,13 @@ Status CloudSchemaChangeJob::_convert_historical_rowsets(const SchemaChangeParam context.txn_expiration = _expiration; context.version = rs_reader->version(); context.rowset_state = VISIBLE; - context.segments_overlap = rs_reader->rowset()->rowset_meta()->segments_overlap(); + const auto input_segments_overlap = rs_reader->rowset()->rowset_meta()->segments_overlap(); + // Cloud schema change rewrites remote rowsets, so the input group layout is no longer + // applicable. Fall back to the conservative overlap state without assuming that the + // rewritten segments are globally ordered. + context.segments_overlap = input_segments_overlap == NONOVERLAPPING_WITHIN_GROUP + ? OVERLAPPING + : input_segments_overlap; context.tablet_schema = _new_tablet->tablet_schema(); context.newest_write_timestamp = rs_reader->newest_write_timestamp(); context.storage_resource = _cloud_storage_engine.get_storage_resource(sc_params.vault_id); diff --git a/be/src/cloud/cloud_snapshot_mgr.cpp b/be/src/cloud/cloud_snapshot_mgr.cpp index bd4e3fee76757f..1f096c7b87a483 100644 --- a/be/src/cloud/cloud_snapshot_mgr.cpp +++ b/be/src/cloud/cloud_snapshot_mgr.cpp @@ -271,6 +271,8 @@ Status CloudSnapshotMgr::_create_rowset_meta( new_rowset_meta_pb->set_creation_time(time(nullptr)); new_rowset_meta_pb->set_num_segments(source_meta_pb.num_segments()); new_rowset_meta_pb->set_rowset_state(source_meta_pb.rowset_state()); + new_rowset_meta_pb->mutable_segment_group_sizes()->CopyFrom( + source_meta_pb.segment_group_sizes()); new_rowset_meta_pb->clear_segments_key_bounds(); for (const auto& key_bound : source_meta_pb.segments_key_bounds()) { diff --git a/be/src/cloud/config.cpp b/be/src/cloud/config.cpp index cd180c2e5b3b87..1de5a5cc60b257 100644 --- a/be/src/cloud/config.cpp +++ b/be/src/cloud/config.cpp @@ -57,6 +57,9 @@ DEFINE_mInt32(max_base_compaction_task_num_per_disk, "2"); DEFINE_mBool(prioritize_query_perf_in_compaction, "false"); DEFINE_mInt32(compaction_max_rowset_count, "10000"); DEFINE_mInt64(compaction_txn_max_size_bytes, "7340032"); // 7MB +DEFINE_mBool(enable_cloud_single_rowset_compaction, "false"); +DEFINE_mInt32(cloud_single_rowset_compaction_min_segments, "512"); +DEFINE_mInt32(cloud_single_rowset_compaction_segment_group_size, "64"); DEFINE_mInt32(refresh_s3_info_interval_s, "60"); DEFINE_mInt32(vacuum_stale_rowsets_interval_s, "300"); diff --git a/be/src/cloud/config.h b/be/src/cloud/config.h index 5ef03b9df446d5..270ca667688b38 100644 --- a/be/src/cloud/config.h +++ b/be/src/cloud/config.h @@ -97,6 +97,9 @@ DECLARE_mInt32(max_base_compaction_task_num_per_disk); DECLARE_mBool(prioritize_query_perf_in_compaction); DECLARE_mInt32(compaction_max_rowset_count); DECLARE_mInt64(compaction_txn_max_size_bytes); +DECLARE_mBool(enable_cloud_single_rowset_compaction); +DECLARE_mInt32(cloud_single_rowset_compaction_min_segments); +DECLARE_mInt32(cloud_single_rowset_compaction_segment_group_size); // CloudStorageEngine config DECLARE_mInt32(refresh_s3_info_interval_s); diff --git a/be/src/cloud/pb_convert.cpp b/be/src/cloud/pb_convert.cpp index a90f93feeda77f..b8d386b201ba39 100644 --- a/be/src/cloud/pb_convert.cpp +++ b/be/src/cloud/pb_convert.cpp @@ -80,6 +80,7 @@ void doris_rowset_meta_to_cloud(RowsetMetaCloudPB* out, const RowsetMetaPB& in) } out->set_txn_expiration(in.txn_expiration()); out->set_segments_overlap_pb(in.segments_overlap_pb()); + out->mutable_segment_group_sizes()->CopyFrom(in.segment_group_sizes()); if (in.has_segments_key_bounds_truncated()) { out->set_segments_key_bounds_truncated(in.segments_key_bounds_truncated()); } @@ -163,6 +164,7 @@ void doris_rowset_meta_to_cloud(RowsetMetaCloudPB* out, RowsetMetaPB&& in) { } out->set_txn_expiration(in.txn_expiration()); out->set_segments_overlap_pb(in.segments_overlap_pb()); + out->mutable_segment_group_sizes()->Swap(in.mutable_segment_group_sizes()); if (in.has_segments_key_bounds_truncated()) { out->set_segments_key_bounds_truncated(in.segments_key_bounds_truncated()); } @@ -258,6 +260,7 @@ void cloud_rowset_meta_to_doris(RowsetMetaPB* out, const RowsetMetaCloudPB& in) } out->set_txn_expiration(in.txn_expiration()); out->set_segments_overlap_pb(in.segments_overlap_pb()); + out->mutable_segment_group_sizes()->CopyFrom(in.segment_group_sizes()); if (in.has_segments_key_bounds_truncated()) { out->set_segments_key_bounds_truncated(in.segments_key_bounds_truncated()); } @@ -341,6 +344,7 @@ void cloud_rowset_meta_to_doris(RowsetMetaPB* out, RowsetMetaCloudPB&& in) { } out->set_txn_expiration(in.txn_expiration()); out->set_segments_overlap_pb(in.segments_overlap_pb()); + out->mutable_segment_group_sizes()->Swap(in.mutable_segment_group_sizes()); if (in.has_segments_key_bounds_truncated()) { out->set_segments_key_bounds_truncated(in.segments_key_bounds_truncated()); } diff --git a/be/src/storage/compaction/compaction.cpp b/be/src/storage/compaction/compaction.cpp index b2aa69ffe6dee6..aa61df40dcc998 100644 --- a/be/src/storage/compaction/compaction.cpp +++ b/be/src/storage/compaction/compaction.cpp @@ -256,6 +256,9 @@ int64_t Compaction::merge_way_num() { } Status Compaction::merge_input_rowsets() { + MergeInputRowsetsResult result; + RETURN_IF_ERROR(prepare_merge_input_rowsets(&result)); + std::vector input_rs_readers; input_rs_readers.reserve(_input_rowsets.size()); for (auto& rowset : _input_rowsets) { @@ -280,37 +283,10 @@ Status Compaction::merge_input_rowsets() { _stats.rowid_conversion = _rowid_conversion.get(); } - int64_t way_num = merge_way_num(); - - Status res; { SCOPED_TIMER(_merge_rowsets_latency_timer); // 1. Merge segment files and write bkd inverted index - if (_is_vertical) { - if (!_tablet->tablet_schema()->cluster_key_uids().empty()) { - RETURN_IF_ERROR(update_delete_bitmap()); - } - auto progress_cb = [compaction_id = this->_compaction_id](int64_t total, - int64_t completed) { - CompactionTaskTracker::instance()->update_progress(compaction_id, total, completed); - }; - res = Merger::vertical_merge_rowsets(_tablet, compaction_type(), *_cur_tablet_schema, - input_rs_readers, _output_rs_writer.get(), - cast_set(get_avg_segment_rows()), - way_num, &_stats, progress_cb); - } else { - if (!_tablet->tablet_schema()->cluster_key_uids().empty()) { - return Status::InternalError( - "mow table with cluster keys does not support non vertical compaction"); - } - res = Merger::vmerge_rowsets(_tablet, compaction_type(), *_cur_tablet_schema, - input_rs_readers, _output_rs_writer.get(), &_stats); - } - - _tablet->last_compaction_status = res; - if (!res.ok()) { - return res; - } + RETURN_IF_ERROR(do_merge_input_rowsets(input_rs_readers, &result)); // 2. Merge the remaining inverted index files of the string type RETURN_IF_ERROR(do_inverted_index_compaction()); } @@ -335,6 +311,7 @@ Status Compaction::merge_input_rowsets() { //RETURN_IF_ERROR(_engine.meta_mgr().commit_rowset(*_output_rowset->rowset_meta().get())); set_delete_predicate_for_output_rowset(); + update_output_rowset_after_build(result); _local_read_bytes_total = _stats.bytes_read_from_local; _remote_read_bytes_total = _stats.bytes_read_from_remote; @@ -351,6 +328,45 @@ Status Compaction::merge_input_rowsets() { return check_correctness(); } +Status Compaction::do_merge_input_rowsets( + const std::vector& input_rs_readers, + MergeInputRowsetsResult* /*result*/) { + return execute_merge(input_rs_readers, merge_way_num(), &_stats); +} + +Status Compaction::execute_merge(const std::vector& input_rs_readers, + int64_t merge_way_num, Merger::Statistics* stats, + std::optional> segment_range, + VerticalMergeProgressContext progress) { + Status status; + if (_is_vertical) { + if (!_tablet->tablet_schema()->cluster_key_uids().empty() && !segment_range.has_value()) { + RETURN_IF_ERROR(update_delete_bitmap()); + } + auto progress_cb = [compaction_id = this->_compaction_id, progress](int64_t total, + int64_t completed) { + CompactionTaskTracker::instance()->update_progress( + compaction_id, total * progress.total_ranges, + total * progress.range_index + completed); + }; + status = Merger::vertical_merge_rowsets(_tablet, compaction_type(), *_cur_tablet_schema, + input_rs_readers, _output_rs_writer.get(), + cast_set(get_avg_segment_rows()), + merge_way_num, stats, progress_cb, segment_range); + } else { + if (!_tablet->tablet_schema()->cluster_key_uids().empty()) { + return Status::InternalError( + "mow table with cluster keys does not support non vertical compaction"); + } + status = Merger::vmerge_rowsets(_tablet, compaction_type(), *_cur_tablet_schema, + input_rs_readers, _output_rs_writer.get(), stats, + segment_range); + } + + _tablet->last_compaction_status = status; + return status; +} + void Compaction::set_delete_predicate_for_output_rowset() { // Now we support delete in cumu compaction, to make all data in rowsets whose version // is below output_version to be delete in the future base compaction, we should carry diff --git a/be/src/storage/compaction/compaction.h b/be/src/storage/compaction/compaction.h index 4bc0b4e3ad20c1..c24cde1aa61676 100644 --- a/be/src/storage/compaction/compaction.h +++ b/be/src/storage/compaction/compaction.h @@ -30,6 +30,7 @@ #endif #include #include +#include #include #include "cloud/cloud_tablet.h" @@ -102,8 +103,35 @@ class Compaction { void set_delete_predicate_for_output_rowset(); protected: + struct MergeInputRowsetsResult { + bool is_segment_grouped = false; + int64_t segment_group_size = 0; + std::vector output_segment_group_sizes; + }; + Status merge_input_rowsets(); + virtual Status prepare_merge_input_rowsets(MergeInputRowsetsResult* /*result*/) { + return Status::OK(); + } + + virtual Status do_merge_input_rowsets( + const std::vector& input_rs_readers, + MergeInputRowsetsResult* result); + + virtual void update_output_rowset_after_build(const MergeInputRowsetsResult& /*result*/) {} + + // Maps one segment-range merge's column-group progress into the whole compaction task. + struct VerticalMergeProgressContext { + int64_t total_ranges; + int64_t range_index; + }; + + Status execute_merge(const std::vector& input_rs_readers, + int64_t merge_way_num, Merger::Statistics* stats, + std::optional> segment_range = std::nullopt, + VerticalMergeProgressContext progress = {1, 0}); + // merge inverted index files Status do_inverted_index_compaction(); @@ -250,6 +278,8 @@ class CloudCompactionMixin : public Compaction { virtual Status garbage_collection(); + Status construct_output_rowset_writer(RowsetWriterContext& ctx) override; + // Helper function to apply truncation and log the result // Returns the number of rowsets that were truncated size_t apply_txn_size_truncation_and_log(const std::string& compaction_name); @@ -266,8 +296,6 @@ class CloudCompactionMixin : public Compaction { virtual Status rebuild_tablet_schema() { return Status::OK(); } private: - Status construct_output_rowset_writer(RowsetWriterContext& ctx) override; - Status set_storage_resource_from_input_rowsets(RowsetWriterContext& ctx); Status execute_compact_impl(int64_t permits); diff --git a/be/src/storage/iterator/vertical_block_reader.cpp b/be/src/storage/iterator/vertical_block_reader.cpp index 1fe1274dac55f1..697d838b14466b 100644 --- a/be/src/storage/iterator/vertical_block_reader.cpp +++ b/be/src/storage/iterator/vertical_block_reader.cpp @@ -56,6 +56,40 @@ VerticalBlockReader::~VerticalBlockReader() { } } +void VerticalBlockReader::_append_grouped_iterator_init_flags( + const RowsetMeta& rowset_meta, std::pair segment_offsets, + size_t added_iterators, std::vector* iterator_init_flag) { + auto [segment_start, segment_end] = segment_offsets; + if (segment_start == segment_end) { + segment_start = 0; + segment_end = rowset_meta.num_segments(); + } + + DORIS_CHECK_GE(segment_start, 0); + DORIS_CHECK_LE(segment_start, segment_end); + DORIS_CHECK_LE(segment_end, rowset_meta.num_segments()); + DORIS_CHECK_EQ(added_iterators, static_cast(segment_end - segment_start)); + DORIS_CHECK_GT(rowset_meta.segment_group_sizes().size(), 1); + + const size_t flags_start = iterator_init_flag->size(); + int64_t group_start = 0; + for (const auto group_size : rowset_meta.segment_group_sizes()) { + DORIS_CHECK_GT(group_size, 0); + const int64_t group_end = group_start + group_size; + const int64_t selected_start = std::max(group_start, segment_start); + const int64_t selected_end = std::min(group_end, segment_end); + if (selected_start < selected_end) { + iterator_init_flag->push_back(true); + iterator_init_flag->insert(iterator_init_flag->end(), + static_cast(selected_end - selected_start - 1), + false); + } + group_start = group_end; + } + DORIS_CHECK_EQ(group_start, rowset_meta.num_segments()); + DORIS_CHECK_EQ(iterator_init_flag->size() - flags_start, added_iterators); +} + Status VerticalBlockReader::next_block_with_aggregation(Block* block, bool* eof) { auto res = (this->*_next_block_func)(block, eof); if (!config::is_cloud_mode()) { @@ -80,23 +114,35 @@ Status VerticalBlockReader::_get_segment_iterators(const ReaderParams& read_para return res; } for (const auto& rs_split : read_params.rs_splits) { + RETURN_IF_ERROR(rs_split.rs_reader->init(&_reader_context, rs_split)); + const auto rowset = rs_split.rs_reader->rowset(); // segment iterator will be inited here // In vertical compaction, every group will load segment so we should cache // segment to avoid tot many s3 head request - bool use_cache = !rs_split.rs_reader->rowset()->is_local(); + bool use_cache = !rowset->is_local(); + size_t iterators_start = segment_iters->size(); RETURN_IF_ERROR(rs_split.rs_reader->get_segment_iterators(&_reader_context, segment_iters, use_cache)); - // if segments overlapping, all segment iterator should be inited in - // heap merge iterator. If segments are none overlapping, only first segment of this - // rowset will be inited and push to heap, other segment will be inited later when current - // segment reached it's end. - // Use this iterator_init_flag so we can load few segments in HeapMergeIterator to save memory - if (rs_split.rs_reader->rowset()->is_segments_overlapping()) { - for (int i = 0; i < rs_split.rs_reader->rowset()->num_segments(); ++i) { - iterator_init_flag->push_back(true); + // If segments overlap, all segment iterators should be initialized in the heap merge + // iterator. For grouped segments, initialize the first iterator of each group because + // segments are ordered within a group but groups may overlap. If segments do not overlap, + // only the first segment of the rowset is initialized; later segments are initialized when + // the current segment reaches its end. + // Use this iterator_init_flag so we can load few segments in HeapMergeIterator to save + // memory. + size_t added_iterators = segment_iters->size() - iterators_start; + if (rowset->is_segments_overlapping()) { + const auto& rowset_meta = *rowset->rowset_meta(); + if (rowset_meta.segments_overlap() == NONOVERLAPPING_WITHIN_GROUP) { + _append_grouped_iterator_init_flags(rowset_meta, rs_split.segment_offsets, + added_iterators, iterator_init_flag); + } else { + for (size_t i = 0; i < added_iterators; ++i) { + iterator_init_flag->push_back(true); + } } } else { - for (int i = 0; i < rs_split.rs_reader->rowset()->num_segments(); ++i) { + for (size_t i = 0; i < added_iterators; ++i) { if (i == 0) { iterator_init_flag->push_back(true); continue; @@ -104,8 +150,8 @@ Status VerticalBlockReader::_get_segment_iterators(const ReaderParams& read_para iterator_init_flag->push_back(false); } } - for (int i = 0; i < rs_split.rs_reader->rowset()->num_segments(); ++i) { - rowset_ids->push_back(rs_split.rs_reader->rowset()->rowset_id()); + for (size_t i = 0; i < added_iterators; ++i) { + rowset_ids->push_back(rowset->rowset_id()); } rs_split.rs_reader->reset_read_options(); } diff --git a/be/src/storage/iterator/vertical_block_reader.h b/be/src/storage/iterator/vertical_block_reader.h index baba2a215f8d19..2907ef52b2058e 100644 --- a/be/src/storage/iterator/vertical_block_reader.h +++ b/be/src/storage/iterator/vertical_block_reader.h @@ -43,8 +43,10 @@ namespace doris { struct RowsetId; class RowSourcesBuffer; +class RowsetMeta; struct RowBatch; struct VerticalCompactionContextStats; +class VerticalBlockReaderTestAccessor; class VerticalBlockReader final : public TabletReader { public: @@ -72,6 +74,13 @@ class VerticalBlockReader final : public TabletReader { static uint64_t nextId; private: + friend class VerticalBlockReaderTestAccessor; + + static void _append_grouped_iterator_init_flags(const RowsetMeta& rowset_meta, + std::pair segment_offsets, + size_t added_iterators, + std::vector* iterator_init_flag); + // Directly read row from rowset and pass to upper caller. No need to do aggregation. // This is usually used for DUPLICATE KEY tables Status _direct_next_block(Block* block, bool* eof); diff --git a/be/src/storage/merger.cpp b/be/src/storage/merger.cpp index f5364d487877aa..2d7712afe86819 100644 --- a/be/src/storage/merger.cpp +++ b/be/src/storage/merger.cpp @@ -68,7 +68,8 @@ namespace doris { Status Merger::vmerge_rowsets(BaseTabletSPtr tablet, ReaderType reader_type, const TabletSchema& cur_tablet_schema, const std::vector& src_rowset_readers, - RowsetWriter* dst_rowset_writer, Statistics* stats_output) { + RowsetWriter* dst_rowset_writer, Statistics* stats_output, + std::optional> segment_range) { if (!cur_tablet_schema.cluster_key_uids().empty()) { return Status::InternalError( "mow table with cluster keys does not support non vertical compaction"); @@ -81,7 +82,11 @@ Status Merger::vmerge_rowsets(BaseTabletSPtr tablet, ReaderType reader_type, TabletReadSource read_source; read_source.rs_splits.reserve(src_rowset_readers.size()); for (const RowsetReaderSharedPtr& rs_reader : src_rowset_readers) { - read_source.rs_splits.emplace_back(rs_reader); + auto& rs_split = read_source.rs_splits.emplace_back(rs_reader); + if (segment_range.has_value()) { + DCHECK_EQ(src_rowset_readers.size(), 1); + rs_split.segment_offsets = segment_range.value(); + } } read_source.fill_delete_predicates(); reader_params.set_read_source(std::move(read_source)); @@ -252,7 +257,7 @@ Status Merger::vertical_compact_one_group( RowsetWriter* dst_rowset_writer, uint32_t max_rows_per_segment, Statistics* stats_output, std::vector key_group_cluster_key_idxes, int64_t batch_size, CompactionSampleInfo* sample_info, VerticalCompactionContextStats* context_stats, - bool enable_sparse_optimization) { + bool enable_sparse_optimization, std::optional> segment_range) { // build tablet reader VLOG_NOTICE << "vertical compact one group, max_rows_per_segment=" << max_rows_per_segment; VerticalBlockReader reader(row_source_buf, context_stats); @@ -266,7 +271,11 @@ Status Merger::vertical_compact_one_group( TabletReadSource read_source; read_source.rs_splits.reserve(src_rowset_readers.size()); for (const RowsetReaderSharedPtr& rs_reader : src_rowset_readers) { - read_source.rs_splits.emplace_back(rs_reader); + auto& rs_split = read_source.rs_splits.emplace_back(rs_reader); + if (segment_range.has_value()) { + DCHECK_EQ(src_rowset_readers.size(), 1); + rs_split.segment_offsets = segment_range.value(); + } } read_source.fill_delete_predicates(); reader_params.set_read_source(std::move(read_source)); @@ -494,7 +503,8 @@ Status Merger::vertical_merge_rowsets(BaseTabletSPtr tablet, ReaderType reader_t RowsetWriter* dst_rowset_writer, uint32_t max_rows_per_segment, int64_t merge_way_num, Statistics* stats_output, - VerticalCompactionProgressCallback progress_cb) { + VerticalCompactionProgressCallback progress_cb, + std::optional> segment_range) { LOG(INFO) << "Start to do vertical compaction, tablet_id: " << tablet->tablet_id(); VerticalCompactionContextStats context_stats; Defer log_context_stats {[&] { @@ -537,29 +547,33 @@ Status Merger::vertical_merge_rowsets(BaseTabletSPtr tablet, ReaderType reader_t progress_cb(column_groups.size(), 0); } - // Calculate total rows for density calculation after compaction + // Segment-range vertical compaction only sees part of a rowset. Do not use partial rows to + // update tablet-level density or drive sparse optimization. int64_t total_rows = 0; - for (const auto& rs_reader : src_rowset_readers) { - total_rows += rs_reader->rowset()->rowset_meta()->num_rows(); - } - // Use historical density for sparse wide table optimization // density = (total_cells - null_cells) / total_cells, smaller means more sparse // When density <= threshold, enable sparse optimization // threshold = 0 means disable, 1 means always enable (default) bool enable_sparse_optimization = false; - if (config::sparse_column_compaction_threshold_percent > 0 && - tablet->keys_type() == KeysType::UNIQUE_KEYS) { - double density = tablet->compaction_density.load(); - enable_sparse_optimization = density <= config::sparse_column_compaction_threshold_percent; + if (!segment_range.has_value()) { + for (const auto& rs_reader : src_rowset_readers) { + total_rows += rs_reader->rowset()->rowset_meta()->num_rows(); + } - LOG(INFO) << "Vertical compaction sparse optimization check: tablet_id=" - << tablet->tablet_id() << ", density=" << density - << ", threshold=" << config::sparse_column_compaction_threshold_percent - << ", total_rows=" << total_rows - << ", num_columns=" << tablet_schema.num_columns() - << ", total_cells=" << total_rows * tablet_schema.num_columns() - << ", enable_sparse_optimization=" << enable_sparse_optimization; + if (config::sparse_column_compaction_threshold_percent > 0 && + tablet->keys_type() == KeysType::UNIQUE_KEYS) { + double density = tablet->compaction_density.load(); + enable_sparse_optimization = + density <= config::sparse_column_compaction_threshold_percent; + + LOG(INFO) << "Vertical compaction sparse optimization check: tablet_id=" + << tablet->tablet_id() << ", density=" << density + << ", threshold=" << config::sparse_column_compaction_threshold_percent + << ", total_rows=" << total_rows + << ", num_columns=" << tablet_schema.num_columns() + << ", total_cells=" << total_rows * tablet_schema.num_columns() + << ", enable_sparse_optimization=" << enable_sparse_optimization; + } } RowSourcesBuffer row_sources_buf(tablet->tablet_id(), dst_rowset_writer->context().tablet_path, @@ -594,7 +608,7 @@ Status Merger::vertical_merge_rowsets(BaseTabletSPtr tablet, ReaderType reader_t } } } - if (need_footer_collection) { + if (!segment_range.has_value() && need_footer_collection) { for (const auto& rs_reader : src_rowset_readers) { auto beta_rowset = std::dynamic_pointer_cast(rs_reader->rowset()); if (!beta_rowset) { @@ -608,7 +622,8 @@ Status Merger::vertical_merge_rowsets(BaseTabletSPtr tablet, ReaderType reader_t << ", rowset_id: " << beta_rowset->rowset_id() << ", status: " << st; continue; } - for (const auto& segment : segments) { + for (int64_t segment_idx = 0; segment_idx < segments.size(); ++segment_idx) { + const auto& segment = segments[segment_idx]; int64_t row_count = segment->num_rows(); auto collect_st = segment->traverse_column_meta_pbs( [&](const segment_v2::ColumnMetaPB& meta) { @@ -630,7 +645,7 @@ Status Merger::vertical_merge_rowsets(BaseTabletSPtr tablet, ReaderType reader_t // Pre-compute per-row estimate for each column group from footer data. std::vector group_per_row_from_footer(column_groups.size(), 0); - std::vector group_footer_fallback(column_groups.size(), false); + std::vector group_footer_fallback(column_groups.size(), segment_range.has_value()); for (size_t i = 0; i < column_groups.size(); ++i) { int64_t group_per_row = 0; bool need_fallback = false; @@ -708,7 +723,7 @@ Status Merger::vertical_merge_rowsets(BaseTabletSPtr tablet, ReaderType reader_t tablet, reader_type, tablet_schema, is_key, column_groups[i], &row_sources_buf, src_rowset_readers, dst_rowset_writer, max_rows_per_segment, group_stats_ptr, key_group_cluster_key_idxes, batch_size, &sample_info, &context_stats, - enable_sparse_optimization); + enable_sparse_optimization, segment_range); { std::unique_lock lock(sample_info_lock); sample_infos[i] = sample_info; @@ -739,7 +754,7 @@ Status Merger::vertical_merge_rowsets(BaseTabletSPtr tablet, ReaderType reader_t // Calculate and store density for next compaction's sparse optimization threshold // density = (total_cells - total_null_count) / total_cells // Smaller density means more sparse - { + if (!segment_range.has_value()) { std::unique_lock lock(sample_info_lock); int64_t total_null_count = 0; for (const auto& info : sample_infos) { diff --git a/be/src/storage/merger.h b/be/src/storage/merger.h index 5185c97e6b0dc9..a53d795500cd64 100644 --- a/be/src/storage/merger.h +++ b/be/src/storage/merger.h @@ -18,6 +18,8 @@ #pragma once #include +#include +#include #include #include "common/status.h" @@ -63,15 +65,17 @@ class Merger { // return OK and set statistics into `*stats_output`. // return others on error - static Status vmerge_rowsets(BaseTabletSPtr tablet, ReaderType reader_type, - const TabletSchema& cur_tablet_schema, - const std::vector& src_rowset_readers, - RowsetWriter* dst_rowset_writer, Statistics* stats_output); + static Status vmerge_rowsets( + BaseTabletSPtr tablet, ReaderType reader_type, const TabletSchema& cur_tablet_schema, + const std::vector& src_rowset_readers, + RowsetWriter* dst_rowset_writer, Statistics* stats_output, + std::optional> segment_range = std::nullopt); static Status vertical_merge_rowsets( BaseTabletSPtr tablet, ReaderType reader_type, const TabletSchema& tablet_schema, const std::vector& src_rowset_readers, RowsetWriter* dst_rowset_writer, uint32_t max_rows_per_segment, int64_t merge_way_num, - Statistics* stats_output, VerticalCompactionProgressCallback progress_cb = nullptr); + Statistics* stats_output, VerticalCompactionProgressCallback progress_cb = nullptr, + std::optional> segment_range = std::nullopt); // for vertical compaction static void vertical_split_columns(const TabletSchema& tablet_schema, @@ -86,7 +90,8 @@ class Merger { RowsetWriter* dst_rowset_writer, uint32_t max_rows_per_segment, Statistics* stats_output, std::vector key_group_cluster_key_idxes, int64_t batch_size, CompactionSampleInfo* sample_info, - VerticalCompactionContextStats* context_stats, bool enable_sparse_optimization = false); + VerticalCompactionContextStats* context_stats, bool enable_sparse_optimization = false, + std::optional> segment_range = std::nullopt); // for segcompaction static Status vertical_compact_one_group( diff --git a/be/src/storage/rowid_conversion.h b/be/src/storage/rowid_conversion.h index 4a516f0b60e365..213ee4c4ef3722 100644 --- a/be/src/storage/rowid_conversion.h +++ b/be/src/storage/rowid_conversion.h @@ -21,6 +21,7 @@ #include #include "common/cast_set.h" +#include "common/check.h" #include "runtime/thread_context.h" #include "storage/olap_common.h" #include "storage/utils.h" @@ -41,6 +42,15 @@ class RowIdConversion { // resize segment rowid map to its rows num Status init_segment_map(const RowsetId& src_rowset_id, const std::vector& num_rows) { for (size_t i = 0; i < num_rows.size(); i++) { + auto src_segment = std::pair {src_rowset_id, cast_set(i)}; + auto iter = _segment_to_id_map.find(src_segment); + // Each segment-group reader initializes all source segments, so reuse existing maps. + if (iter != _segment_to_id_map.end()) { + DORIS_CHECK_LT(iter->second, _segments_rowid_map.size()); + DORIS_CHECK_EQ(_segments_rowid_map[iter->second].size(), num_rows[i]); + continue; + } + constexpr size_t RESERVED_MEMORY = 10 * 1024 * 1024; // 10M if (doris::GlobalMemoryArbitrator::is_exceed_hard_mem_limit(RESERVED_MEMORY)) { return Status::MemoryLimitExceeded(fmt::format( @@ -60,9 +70,10 @@ class RowIdConversion { ->consumption())); } - uint32_t id = static_cast(_segments_rowid_map.size()); - _segment_to_id_map.emplace(std::pair {src_rowset_id, i}, id); - _id_to_segment_map.emplace_back(src_rowset_id, i); + uint32_t id = cast_set(_segments_rowid_map.size()); + auto insert_result = _segment_to_id_map.emplace(src_segment, id); + DORIS_CHECK(insert_result.second); + _id_to_segment_map.push_back(src_segment); std::vector> vec( num_rows[i], std::pair(UINT32_MAX, UINT32_MAX)); diff --git a/be/src/storage/rowset/rowset_meta.h b/be/src/storage/rowset/rowset_meta.h index 51b11494862434..c35a8195c608b6 100644 --- a/be/src/storage/rowset/rowset_meta.h +++ b/be/src/storage/rowset/rowset_meta.h @@ -159,6 +159,22 @@ class RowsetMeta : public MetadataAdder { auto& get_num_segment_rows() const { return _rowset_meta_pb.num_segment_rows(); } + void set_segment_group_sizes(const std::vector& segment_group_sizes) { + DORIS_CHECK_GT(segment_group_sizes.size(), 1); + int64_t segment_count = 0; + for (const auto group_size : segment_group_sizes) { + DORIS_CHECK_GT(group_size, 0); + segment_count += group_size; + } + DORIS_CHECK_EQ(segment_count, num_segments()); + _rowset_meta_pb.mutable_segment_group_sizes()->Assign(segment_group_sizes.cbegin(), + segment_group_sizes.cend()); + } + + void clear_segment_group_sizes() { _rowset_meta_pb.clear_segment_group_sizes(); } + + const auto& segment_group_sizes() const { return _rowset_meta_pb.segment_group_sizes(); } + int64_t total_disk_size() const { return _rowset_meta_pb.total_disk_size(); } void set_total_disk_size(int64_t total_disk_size) { @@ -293,15 +309,17 @@ class RowsetMeta : public MetadataAdder { // 1. the rowset contains more than one segment // 2. the rowset's start version == end version (non-singleton rowset was generated by compaction process // which always produces non-overlapped segments) - // 3. segments_overlap() flag is not NONOVERLAPPING (OVERLAP_UNKNOWN and OVERLAPPING are OK) + // 3. segments_overlap() flag is not NONOVERLAPPING (OVERLAP_UNKNOWN, OVERLAPPING, + // and NONOVERLAPPING_WITHIN_GROUP are considered overlapping) bool is_segments_overlapping() const { return num_segments() > 1 && is_singleton_delta() && segments_overlap() != NONOVERLAPPING; } bool produced_by_compaction() const { - return has_version() && - (start_version() < end_version() || - (start_version() == end_version() && segments_overlap() == NONOVERLAPPING)); + return has_version() && (start_version() < end_version() || + (start_version() == end_version() && + (segments_overlap() == NONOVERLAPPING || + segments_overlap() == NONOVERLAPPING_WITHIN_GROUP))); } // get the compaction score of this rowset. diff --git a/be/src/storage/rowset/vertical_beta_rowset_writer.cpp b/be/src/storage/rowset/vertical_beta_rowset_writer.cpp index c6e7a12b91bc6a..4655726ee0ccc6 100644 --- a/be/src/storage/rowset/vertical_beta_rowset_writer.cpp +++ b/be/src/storage/rowset/vertical_beta_rowset_writer.cpp @@ -139,8 +139,7 @@ Status VerticalBetaRowsetWriter::_flush_columns(segment_v2::SegmentWriter* se key_bounds.set_min_key(min_key.to_string()); key_bounds.set_max_key(max_key.to_string()); this->_segments_encoded_key_bounds.emplace_back(std::move(key_bounds)); - this->_segment_num_rows.resize(_cur_writer_idx + 1); - this->_segment_num_rows[_cur_writer_idx] = _segment_writers[_cur_writer_idx]->row_count(); + this->_segment_num_rows.emplace_back(segment_writer->row_count()); } return Status::OK(); } @@ -204,6 +203,7 @@ template requires std::is_base_of_v Status VerticalBetaRowsetWriter::final_flush() { for (auto& segment_writer : _segment_writers) { + DCHECK(segment_writer); uint64_t segment_size = 0; //uint64_t footer_position = 0; segment_v2::SegmentIndexFileCacheInfo index_file_cache_info; @@ -219,6 +219,8 @@ Status VerticalBetaRowsetWriter::final_flush() { segment_writer.reset(); _record_segment_index_file_cache_preload(segment_id, index_file_cache_info); } + _segment_writers.clear(); + _cur_writer_idx = 0; return Status::OK(); } diff --git a/be/test/cloud/cloud_compaction_test.cpp b/be/test/cloud/cloud_compaction_test.cpp index 62141d5ad883c5..0088086d8034a2 100644 --- a/be/test/cloud/cloud_compaction_test.cpp +++ b/be/test/cloud/cloud_compaction_test.cpp @@ -24,6 +24,7 @@ #include #include #include +#include #include #include @@ -51,6 +52,20 @@ namespace doris { class TabletMap; +namespace { + +void expect_segment_group_merge_ranges(const std::vector& actual, + const std::vector& expected) { + ASSERT_EQ(actual.size(), expected.size()); + for (size_t i = 0; i < expected.size(); ++i) { + EXPECT_EQ(actual[i].segment_start, expected[i].segment_start); + EXPECT_EQ(actual[i].segment_end, expected[i].segment_end); + EXPECT_EQ(actual[i].merge_way_num, expected[i].merge_way_num); + } +} + +} // namespace + class CloudCompactionTest : public testing::Test { CloudCompactionTest() : _engine(CloudStorageEngine(EngineOptions {})) {} void SetUp() override { @@ -324,7 +339,7 @@ TEST_F(CloudCompactionTest, generate_cloud_compaction_tasks_clears_metrics_witho } static RowsetSharedPtr create_rowset(Version version, int num_segments, bool overlapping, - int data_size) { + int data_size, int num_key_columns = 1) { auto rs_meta = std::make_shared(); rs_meta->set_rowset_type(BETA_ROWSET); // important rs_meta->_rowset_meta_pb.set_start_version(version.first); @@ -332,6 +347,19 @@ static RowsetSharedPtr create_rowset(Version version, int num_segments, bool ove rs_meta->set_num_segments(num_segments); rs_meta->set_segments_overlap(overlapping ? OVERLAPPING : NONOVERLAPPING); rs_meta->set_total_disk_size(data_size); + TabletSchemaPB tablet_schema_pb; + tablet_schema_pb.set_keys_type(DUP_KEYS); + for (int i = 0; i < num_key_columns + 1; ++i) { + ColumnPB* column = tablet_schema_pb.add_column(); + column->set_unique_id(i); + column->set_name("c" + std::to_string(i)); + column->set_type("INT"); + column->set_is_key(i < num_key_columns); + column->set_is_nullable(false); + } + auto tablet_schema = std::make_shared(); + tablet_schema->init_from_pb(tablet_schema_pb); + rs_meta->set_tablet_schema(tablet_schema); RowsetSharedPtr rowset; Status st = RowsetFactory::create_rowset(nullptr, "", rs_meta, &rowset); if (!st.ok()) { @@ -1084,6 +1112,345 @@ TEST_F(CloudCompactionTest, should_cache_compaction_output) { LOG(INFO) << "should_cache_compaction_output done"; } +TEST_F(CloudCompactionTest, single_rowset_grouped_compaction_execution_path_conditions) { + auto old_enable = config::enable_cloud_single_rowset_compaction; + auto old_min_segments = config::cloud_single_rowset_compaction_min_segments; + auto old_group_size = config::cloud_single_rowset_compaction_segment_group_size; + Defer restore_config {[&] { + config::enable_cloud_single_rowset_compaction = old_enable; + config::cloud_single_rowset_compaction_min_segments = old_min_segments; + config::cloud_single_rowset_compaction_segment_group_size = old_group_size; + }}; + config::enable_cloud_single_rowset_compaction = true; + config::cloud_single_rowset_compaction_min_segments = 4; + config::cloud_single_rowset_compaction_segment_group_size = 2; + + RowsetSharedPtr candidate = create_rowset(Version(2, 2), 4, true, 1024); + ASSERT_TRUE(candidate != nullptr); + const auto& tablet_schema = *candidate->tablet_schema(); + EXPECT_TRUE(cloud::is_single_rowset_compaction_candidate(candidate)); + EXPECT_TRUE(cloud::should_use_single_rowset_grouped_compaction({candidate}, tablet_schema, + CUMULATIVE_SIZE_BASED_POLICY)); + EXPECT_FALSE(cloud::should_use_single_rowset_grouped_compaction({candidate}, tablet_schema, + CUMULATIVE_TIME_SERIES_POLICY)); + + TabletSchemaPB cluster_key_schema_pb; + tablet_schema.to_schema_pb(&cluster_key_schema_pb); + cluster_key_schema_pb.set_keys_type(UNIQUE_KEYS); + cluster_key_schema_pb.add_cluster_key_uids(1); + TabletSchema cluster_key_schema; + cluster_key_schema.init_from_pb(cluster_key_schema_pb); + EXPECT_FALSE(cloud::should_use_single_rowset_grouped_compaction({candidate}, cluster_key_schema, + CUMULATIVE_SIZE_BASED_POLICY)); + + config::enable_cloud_single_rowset_compaction = false; + EXPECT_TRUE(cloud::is_single_rowset_compaction_candidate(candidate)); + EXPECT_FALSE(cloud::should_use_single_rowset_grouped_compaction({candidate}, tablet_schema, + CUMULATIVE_SIZE_BASED_POLICY)); + config::enable_cloud_single_rowset_compaction = true; + + RowsetSharedPtr non_overlapping = create_rowset(Version(3, 3), 4, false, 1024); + ASSERT_TRUE(non_overlapping != nullptr); + EXPECT_FALSE(cloud::is_single_rowset_compaction_candidate(non_overlapping)); + EXPECT_FALSE(cloud::should_use_single_rowset_grouped_compaction( + {non_overlapping}, tablet_schema, CUMULATIVE_SIZE_BASED_POLICY)); + + RowsetSharedPtr too_few_segments = create_rowset(Version(4, 4), 3, true, 1024); + ASSERT_TRUE(too_few_segments != nullptr); + EXPECT_FALSE(cloud::is_single_rowset_compaction_candidate(too_few_segments)); + EXPECT_FALSE(cloud::should_use_single_rowset_grouped_compaction( + {too_few_segments}, tablet_schema, CUMULATIVE_SIZE_BASED_POLICY)); + + RowsetSharedPtr grouped_candidate = create_rowset(Version(5, 5), 8, true, 1024); + ASSERT_TRUE(grouped_candidate != nullptr); + grouped_candidate->rowset_meta()->set_segments_overlap(NONOVERLAPPING_WITHIN_GROUP); + grouped_candidate->rowset_meta()->set_segment_group_sizes({2, 2, 2, 2}); + EXPECT_TRUE(cloud::is_single_rowset_compaction_candidate(grouped_candidate)); + + RowsetSharedPtr grouped_with_too_few_groups = create_rowset(Version(6, 6), 8, true, 1024); + ASSERT_TRUE(grouped_with_too_few_groups != nullptr); + grouped_with_too_few_groups->rowset_meta()->set_segments_overlap(NONOVERLAPPING_WITHIN_GROUP); + grouped_with_too_few_groups->rowset_meta()->set_segment_group_sizes({3, 3, 2}); + EXPECT_FALSE(cloud::is_single_rowset_compaction_candidate(grouped_with_too_few_groups)); + + RowsetSharedPtr no_key_columns = create_rowset(Version(7, 7), 4, true, 1024, 0); + ASSERT_TRUE(no_key_columns != nullptr); + EXPECT_TRUE(cloud::is_single_rowset_compaction_candidate(no_key_columns)); + EXPECT_FALSE(cloud::should_use_single_rowset_grouped_compaction( + {no_key_columns}, *no_key_columns->tablet_schema(), CUMULATIVE_SIZE_BASED_POLICY)); + + RowsetSharedPtr with_delete_predicate = create_rowset(Version(8, 8), 4, true, 1024); + ASSERT_TRUE(with_delete_predicate != nullptr); + DeletePredicatePB delete_predicate; + auto* in_predicate = delete_predicate.add_in_predicates(); + in_predicate->set_column_name("c1"); + in_predicate->add_values("1"); + with_delete_predicate->rowset_meta()->set_delete_predicate(std::move(delete_predicate)); + EXPECT_FALSE(cloud::is_single_rowset_compaction_candidate(with_delete_predicate)); + EXPECT_FALSE(cloud::should_use_single_rowset_grouped_compaction( + {with_delete_predicate}, tablet_schema, CUMULATIVE_SIZE_BASED_POLICY)); + + RowsetSharedPtr another_candidate = create_rowset(Version(9, 9), 4, true, 1024); + ASSERT_TRUE(another_candidate != nullptr); + EXPECT_FALSE(cloud::should_use_single_rowset_grouped_compaction( + {candidate, another_candidate}, tablet_schema, CUMULATIVE_SIZE_BASED_POLICY)); + EXPECT_FALSE(cloud::should_use_single_rowset_grouped_compaction({}, tablet_schema, + CUMULATIVE_SIZE_BASED_POLICY)); + + CloudTabletSPtr tablet = std::make_shared(_engine, _tablet_meta); + CloudCumulativeCompaction compaction(_engine, tablet); + compaction._input_rowsets = {candidate}; + compaction._cur_tablet_schema = candidate->tablet_schema(); + compaction._single_rowset_compaction_segment_group_size = + config::cloud_single_rowset_compaction_segment_group_size; + Compaction::MergeInputRowsetsResult result; + ASSERT_TRUE(compaction.prepare_merge_input_rowsets(&result).ok()); + EXPECT_TRUE(compaction._single_rowset_compaction_segment_group_size.has_value()); + EXPECT_TRUE(result.is_segment_grouped); + EXPECT_EQ(result.segment_group_size, config::cloud_single_rowset_compaction_segment_group_size); + + _tablet_meta->set_compaction_policy(std::string(CUMULATIVE_TIME_SERIES_POLICY)); + CloudCumulativeCompaction time_series_compaction(_engine, tablet); + time_series_compaction._input_rowsets = {candidate}; + time_series_compaction._cur_tablet_schema = candidate->tablet_schema(); + Compaction::MergeInputRowsetsResult time_series_result; + ASSERT_TRUE(time_series_compaction.prepare_merge_input_rowsets(&time_series_result).ok()); + EXPECT_FALSE(time_series_compaction._single_rowset_compaction_segment_group_size.has_value()); + EXPECT_FALSE(time_series_result.is_segment_grouped); +} + +TEST_F(CloudCompactionTest, single_rowset_grouped_compaction_builds_logical_group_ranges) { + RowsetMeta overlapping_meta; + overlapping_meta.set_num_segments(5); + overlapping_meta.set_segments_overlap(OVERLAPPING); + + const auto overlapping_ranges = cloud::build_segment_group_merge_ranges(overlapping_meta, 2); + expect_segment_group_merge_ranges(overlapping_ranges, + {{.segment_start = 0, .segment_end = 2, .merge_way_num = 2}, + {.segment_start = 2, .segment_end = 4, .merge_way_num = 2}, + {.segment_start = 4, .segment_end = 5, .merge_way_num = 1}}); + + const auto single_overlapping_range = + cloud::build_segment_group_merge_ranges(overlapping_meta, 10); + expect_segment_group_merge_ranges(single_overlapping_range, + {{.segment_start = 0, .segment_end = 5, .merge_way_num = 5}}); + + overlapping_meta.set_segments_overlap(NONOVERLAPPING); + const auto nonoverlapping_ranges = cloud::build_segment_group_merge_ranges(overlapping_meta, 2); + expect_segment_group_merge_ranges(nonoverlapping_ranges, + {{.segment_start = 0, .segment_end = 2, .merge_way_num = 2}, + {.segment_start = 2, .segment_end = 4, .merge_way_num = 2}, + {.segment_start = 4, .segment_end = 5, .merge_way_num = 1}}); + + overlapping_meta.set_segments_overlap(OVERLAP_UNKNOWN); + const auto unknown_overlap_ranges = + cloud::build_segment_group_merge_ranges(overlapping_meta, 2); + expect_segment_group_merge_ranges(unknown_overlap_ranges, + {{.segment_start = 0, .segment_end = 2, .merge_way_num = 2}, + {.segment_start = 2, .segment_end = 4, .merge_way_num = 2}, + {.segment_start = 4, .segment_end = 5, .merge_way_num = 1}}); + + RowsetMeta grouped_meta; + grouped_meta.set_num_segments(5); + grouped_meta.set_segments_overlap(NONOVERLAPPING_WITHIN_GROUP); + grouped_meta.set_segment_group_sizes({2, 2, 1}); + + const auto grouped_ranges = cloud::build_segment_group_merge_ranges(grouped_meta, 2); + expect_segment_group_merge_ranges(grouped_ranges, + {{.segment_start = 0, .segment_end = 4, .merge_way_num = 2}, + {.segment_start = 4, .segment_end = 5, .merge_way_num = 1}}); +} + +TEST_F(CloudCompactionTest, single_rowset_grouped_compaction_builds_group_range_boundaries) { + RowsetMeta grouped_meta; + grouped_meta.set_num_segments(5); + grouped_meta.set_segments_overlap(NONOVERLAPPING_WITHIN_GROUP); + grouped_meta.set_segment_group_sizes({2, 2, 1}); + + const auto single_range = cloud::build_segment_group_merge_ranges(grouped_meta, 10); + expect_segment_group_merge_ranges(single_range, + {{.segment_start = 0, .segment_end = 5, .merge_way_num = 3}}); + + grouped_meta.set_num_segments(10); + grouped_meta.set_segment_group_sizes({1, 2, 3, 4}); + const auto exact_ranges = cloud::build_segment_group_merge_ranges(grouped_meta, 2); + expect_segment_group_merge_ranges( + exact_ranges, {{.segment_start = 0, .segment_end = 3, .merge_way_num = 2}, + {.segment_start = 3, .segment_end = 10, .merge_way_num = 2}}); + + grouped_meta.set_num_segments(15); + grouped_meta.set_segment_group_sizes({3, 1, 4, 2, 5}); + const auto irregular_ranges = cloud::build_segment_group_merge_ranges(grouped_meta, 2); + expect_segment_group_merge_ranges( + irregular_ranges, {{.segment_start = 0, .segment_end = 4, .merge_way_num = 2}, + {.segment_start = 4, .segment_end = 10, .merge_way_num = 2}, + {.segment_start = 10, .segment_end = 15, .merge_way_num = 1}}); +} + +TEST_F(CloudCompactionTest, single_rowset_grouped_compaction_rejects_invalid_range_input) { + RowsetMeta rowset_meta; + rowset_meta.set_num_segments(5); + rowset_meta.set_segments_overlap(OVERLAPPING); + EXPECT_DEATH(static_cast(cloud::build_segment_group_merge_ranges(rowset_meta, 1)), ""); + + RowsetMeta empty_rowset_meta; + empty_rowset_meta.set_segments_overlap(OVERLAPPING); + EXPECT_DEATH(static_cast(cloud::build_segment_group_merge_ranges(empty_rowset_meta, 2)), + ""); + + RowsetMetaPB invalid_group_layout_pb; + invalid_group_layout_pb.set_rowset_id(1); + invalid_group_layout_pb.set_num_segments(5); + invalid_group_layout_pb.set_segments_overlap_pb(NONOVERLAPPING_WITHIN_GROUP); + + RowsetMeta empty_group_layout; + ASSERT_TRUE(empty_group_layout.init_from_pb(invalid_group_layout_pb)); + EXPECT_DEATH(static_cast(cloud::build_segment_group_merge_ranges(empty_group_layout, 2)), + ""); + + invalid_group_layout_pb.add_segment_group_sizes(2); + invalid_group_layout_pb.add_segment_group_sizes(2); + RowsetMeta invalid_group_layout; + ASSERT_TRUE(invalid_group_layout.init_from_pb(invalid_group_layout_pb)); + EXPECT_DEATH( + static_cast(cloud::build_segment_group_merge_ranges(invalid_group_layout, 2)), + ""); + + invalid_group_layout_pb.clear_segment_group_sizes(); + invalid_group_layout_pb.add_segment_group_sizes(2); + invalid_group_layout_pb.add_segment_group_sizes(0); + invalid_group_layout_pb.add_segment_group_sizes(3); + RowsetMeta zero_sized_group_layout; + ASSERT_TRUE(zero_sized_group_layout.init_from_pb(invalid_group_layout_pb)); + EXPECT_DEATH( + static_cast(cloud::build_segment_group_merge_ranges(zero_sized_group_layout, 2)), + ""); +} + +TEST_F(CloudCompactionTest, single_rowset_grouped_compaction_calculates_cumulative_point) { + CloudTabletSPtr tablet = std::make_shared(_engine, _tablet_meta); + CloudCumulativeCompaction compaction(_engine, tablet); + compaction._input_rowsets = {create_rowset(Version(2, 2), 5, true, 1024)}; + compaction._output_rowset = create_rowset(Version(2, 2), 5, false, 1024); + ASSERT_TRUE(compaction._input_rowsets.front() != nullptr); + ASSERT_TRUE(compaction._output_rowset != nullptr); + compaction._single_rowset_compaction_segment_group_size = 2; + + EXPECT_TRUE(compaction.should_calculate_new_cumulative_point(2)); + EXPECT_FALSE(compaction.should_calculate_new_cumulative_point(1)); + + compaction._output_rowset->rowset_meta()->set_segments_overlap(NONOVERLAPPING_WITHIN_GROUP); + compaction._output_rowset->rowset_meta()->set_segment_group_sizes({2, 2, 1}); + EXPECT_FALSE(compaction.should_calculate_new_cumulative_point(2)); + + compaction._single_rowset_compaction_segment_group_size.reset(); + EXPECT_TRUE(compaction.should_calculate_new_cumulative_point(1)); +} + +TEST_F(CloudCompactionTest, single_rowset_grouped_compaction_uses_selection_snapshot) { + auto old_enable = config::enable_cloud_single_rowset_compaction; + auto old_min_segments = config::cloud_single_rowset_compaction_min_segments; + auto old_group_size = config::cloud_single_rowset_compaction_segment_group_size; + Defer restore_config {[&] { + config::enable_cloud_single_rowset_compaction = old_enable; + config::cloud_single_rowset_compaction_min_segments = old_min_segments; + config::cloud_single_rowset_compaction_segment_group_size = old_group_size; + }}; + config::enable_cloud_single_rowset_compaction = true; + config::cloud_single_rowset_compaction_min_segments = 4; + config::cloud_single_rowset_compaction_segment_group_size = 2; + + std::vector rowsets; + auto grouped_rowset = create_rowset(Version(2, 2), 4, true, 1024); + ASSERT_TRUE(grouped_rowset != nullptr); + rowsets.push_back(grouped_rowset); + for (int64_t version = 3; version <= 13; ++version) { + auto rowset = create_rowset(Version(version, version), 1, false, 1024); + ASSERT_TRUE(rowset != nullptr); + rowsets.push_back(std::move(rowset)); + } + + TabletSchemaPB tablet_schema_pb; + grouped_rowset->tablet_schema()->to_schema_pb(&tablet_schema_pb); + _tablet_meta->mutable_tablet_schema()->init_from_pb(tablet_schema_pb); + _tablet_meta->set_compaction_policy(std::string(CUMULATIVE_SIZE_BASED_POLICY)); + CloudTabletSPtr tablet = std::make_shared(_engine, _tablet_meta); + { + std::unique_lock wlock(tablet->get_header_lock()); + tablet->add_rowsets(std::move(rowsets), false, wlock, false); + } + + for (const int32_t invalid_group_size : {0, 1}) { + config::cloud_single_rowset_compaction_segment_group_size = invalid_group_size; + CloudCumulativeCompaction invalid_config_compaction(_engine, tablet); + ASSERT_TRUE(invalid_config_compaction.pick_rowsets_to_compact().ok()); + EXPECT_FALSE( + invalid_config_compaction._single_rowset_compaction_segment_group_size.has_value()); + } + config::cloud_single_rowset_compaction_segment_group_size = 2; + + CloudCumulativeCompaction compaction(_engine, tablet); + ASSERT_TRUE(compaction.pick_rowsets_to_compact().ok()); + ASSERT_EQ(compaction._input_rowsets.size(), 1); + EXPECT_EQ(compaction._input_rowsets.front(), grouped_rowset); + ASSERT_TRUE(compaction._single_rowset_compaction_segment_group_size.has_value()); + EXPECT_EQ(*compaction._single_rowset_compaction_segment_group_size, 2); + + config::enable_cloud_single_rowset_compaction = false; + config::cloud_single_rowset_compaction_min_segments = 5; + config::cloud_single_rowset_compaction_segment_group_size = 3; + + Compaction::MergeInputRowsetsResult result; + ASSERT_TRUE(compaction.prepare_merge_input_rowsets(&result).ok()); + EXPECT_TRUE(compaction._single_rowset_compaction_segment_group_size.has_value()); + EXPECT_TRUE(result.is_segment_grouped); + EXPECT_EQ(result.segment_group_size, 2); +} + +TEST_F(CloudCompactionTest, single_rowset_grouped_compaction_honors_notready_policy_filter) { + auto old_enable = config::enable_cloud_single_rowset_compaction; + auto old_min_segments = config::cloud_single_rowset_compaction_min_segments; + auto old_enable_empty_rowset_compaction = config::enable_empty_rowset_compaction; + Defer restore_config {[&] { + config::enable_cloud_single_rowset_compaction = old_enable; + config::cloud_single_rowset_compaction_min_segments = old_min_segments; + config::enable_empty_rowset_compaction = old_enable_empty_rowset_compaction; + }}; + config::enable_cloud_single_rowset_compaction = true; + config::cloud_single_rowset_compaction_min_segments = 4; + config::enable_empty_rowset_compaction = false; + + std::vector rowsets; + // Keep enough older inputs mergeable after the NOTREADY policy filters versions 11 through 20. + for (int64_t version = 2; version <= 19; ++version) { + auto rowset = create_rowset(Version(version, version), 1, false, 1024); + ASSERT_TRUE(rowset != nullptr); + rowsets.push_back(std::move(rowset)); + } + auto filtered_grouped_rowset = create_rowset(Version(20, 20), 4, true, 1024); + ASSERT_TRUE(filtered_grouped_rowset != nullptr); + rowsets.push_back(filtered_grouped_rowset); + + TabletSchemaPB tablet_schema_pb; + filtered_grouped_rowset->tablet_schema()->to_schema_pb(&tablet_schema_pb); + _tablet_meta->mutable_tablet_schema()->init_from_pb(tablet_schema_pb); + _tablet_meta->set_compaction_policy(std::string(CUMULATIVE_SIZE_BASED_POLICY)); + _tablet_meta->set_tablet_state(TABLET_NOTREADY); + CloudTabletSPtr tablet = std::make_shared(_engine, _tablet_meta); + tablet->set_alter_version(1); + { + std::unique_lock wlock(tablet->get_header_lock()); + tablet->add_rowsets(std::move(rowsets), false, wlock, false); + } + + CloudCumulativeCompaction compaction(_engine, tablet); + ASSERT_TRUE(compaction.pick_rowsets_to_compact().ok()); + ASSERT_EQ(compaction._input_rowsets.size(), 9); + EXPECT_EQ(compaction._input_rowsets.front()->version(), Version(2, 2)); + EXPECT_EQ(compaction._input_rowsets.back()->version(), Version(10, 10)); + EXPECT_FALSE(compaction._single_rowset_compaction_segment_group_size.has_value()); +} + TEST_F(CloudCompactionTest, test_truncate_rowsets_by_txn_size_empty_input) { std::vector rowsets; int64_t kept_size = 100; diff --git a/be/test/cloud/cloud_cumulative_compaction_policy_test.cpp b/be/test/cloud/cloud_cumulative_compaction_policy_test.cpp index 6e9fa1e1aaa13e..ad651374fdac9b 100644 --- a/be/test/cloud/cloud_cumulative_compaction_policy_test.cpp +++ b/be/test/cloud/cloud_cumulative_compaction_policy_test.cpp @@ -185,6 +185,63 @@ TEST_F(TestCloudSizeBasedCumulativeCompactionPolicy, EXPECT_EQ(Version(-1, -1), last_delete_version); } +TEST_F(TestCloudSizeBasedCumulativeCompactionPolicy, + overlapping_output_uses_policy_to_advance_cumulative_point) { + _tablet_meta->set_tablet_state(TABLET_RUNNING); + _tablet_meta->set_time_series_compaction_level_threshold(1); + CloudTablet tablet(_engine, _tablet_meta); + tablet._base_size = 100; + + auto output_rowset = create_rowset(Version(3, 3), 5, true, 256 * 1024 * 1024); + Version last_delete_version {-1, -1}; + + CloudSizeBasedCumulativeCompactionPolicy size_based_policy; + EXPECT_EQ(4, size_based_policy.new_cumulative_point(&tablet, output_rowset, last_delete_version, + 2)); + + CloudTimeSeriesCumulativeCompactionPolicy time_series_policy; + EXPECT_EQ(4, time_series_policy.new_cumulative_point(&tablet, output_rowset, + last_delete_version, 2)); +} + +TEST_F(TestCloudSizeBasedCumulativeCompactionPolicy, + pick_input_rowsets_small_single_overlapping_rowset_not_trimmed_empty) { + TTabletSchema schema; + schema.keys_type = TKeysType::DUP_KEYS; + TabletMetaSharedPtr tablet_meta(new TabletMeta(1, 2, 15673, 15674, 4, 5, schema, 6, {{7, 8}}, + UniqueId(9, 10), TTabletType::TABLET_TYPE_DISK, + TCompressionType::LZ4F)); + tablet_meta->set_tablet_state(TABLET_RUNNING); + CloudTablet tablet(_engine, tablet_meta); + tablet._base_size = 1024L * 1024 * 1024; + CloudSizeBasedCumulativeCompactionPolicy policy; + + // Match the regression case: one small rowset, six segments, min score 2, max score 1000. + auto rowset = create_rowset(Version(2, 2), 6, true, 6370); + ASSERT_NE(nullptr, rowset); + for (auto overlap : {OVERLAPPING, NONOVERLAPPING_WITHIN_GROUP, NONOVERLAPPING}) { + SCOPED_TRACE(static_cast(overlap)); + rowset->rowset_meta()->set_segments_overlap(overlap); + rowset->rowset_meta()->clear_segment_group_sizes(); + if (overlap == NONOVERLAPPING_WITHIN_GROUP) { + rowset->rowset_meta()->set_segment_group_sizes({2, 2, 2}); + } + std::vector input_rowsets; + Version last_delete_version {-1, -1}; + size_t compaction_score = 0; + policy.pick_input_rowsets(&tablet, {rowset}, 1000, 2, &input_rowsets, + &last_delete_version, &compaction_score, true); + if (overlap == NONOVERLAPPING) { + EXPECT_TRUE(input_rowsets.empty()); + EXPECT_EQ(0, compaction_score); + } else { + ASSERT_EQ(1, input_rowsets.size()); + EXPECT_EQ(rowset, input_rowsets.front()); + EXPECT_EQ(6, compaction_score); + } + } +} + // Test case: Empty rowset compaction with skip_trim TEST_F(TestCloudSizeBasedCumulativeCompactionPolicy, pick_input_rowsets_empty_rowset_compaction) { // Save original config values diff --git a/be/test/cloud/cloud_snapshot_mgr_test.cpp b/be/test/cloud/cloud_snapshot_mgr_test.cpp index 5b097be04d4685..64e039b46f5131 100644 --- a/be/test/cloud/cloud_snapshot_mgr_test.cpp +++ b/be/test/cloud/cloud_snapshot_mgr_test.cpp @@ -101,6 +101,9 @@ TEST_F(CloudSnapshotMgrTest, TestConvertRowsets) { rowset_meta->set_tablet_id(1000); rowset_meta->set_txn_id(2000); rowset_meta->set_num_segments(3); + rowset_meta->set_segments_overlap_pb(NONOVERLAPPING_WITHIN_GROUP); + rowset_meta->add_segment_group_sizes(1); + rowset_meta->add_segment_group_sizes(2); rowset_meta->set_num_rows(100); rowset_meta->set_start_version(100); rowset_meta->set_end_version(101); @@ -148,6 +151,10 @@ TEST_F(CloudSnapshotMgrTest, TestConvertRowsets) { EXPECT_EQ(output_meta_pb.rs_metas(0).tablet_schema().index(0).index_id(), 1001); EXPECT_EQ(output_meta_pb.rs_metas(0).tablet_schema().index(1).index_id(), 1002); EXPECT_EQ(output_meta_pb.rs_metas(0).num_segments(), 3); + EXPECT_EQ(output_meta_pb.rs_metas(0).segments_overlap_pb(), NONOVERLAPPING_WITHIN_GROUP); + ASSERT_EQ(output_meta_pb.rs_metas(0).segment_group_sizes_size(), 2); + EXPECT_EQ(output_meta_pb.rs_metas(0).segment_group_sizes(0), 1); + EXPECT_EQ(output_meta_pb.rs_metas(0).segment_group_sizes(1), 2); EXPECT_EQ(output_meta_pb.rs_metas(0).num_rows(), 100); EXPECT_EQ(output_meta_pb.rs_metas(0).start_version(), 100); EXPECT_EQ(output_meta_pb.rs_metas(0).end_version(), 101); diff --git a/be/test/storage/cloud_file_cache_write_index_only_test.cpp b/be/test/storage/cloud_file_cache_write_index_only_test.cpp index 49f43e80cf16b0..b7baae652028f5 100644 --- a/be/test/storage/cloud_file_cache_write_index_only_test.cpp +++ b/be/test/storage/cloud_file_cache_write_index_only_test.cpp @@ -699,7 +699,6 @@ TEST_F(CloudFileCacheWriteIndexOnlyTest, auto writer_result = RowsetFactory::create_rowset_writer(*_engine, context, true); ASSERT_TRUE(writer_result.has_value()) << writer_result.error(); auto rowset_writer = std::move(writer_result).value(); - std::vector key_column_ids = {0}; auto key_block = create_column_block(tablet_schema, key_column_ids, 8, 1); auto st = rowset_writer->add_columns(&key_block, key_column_ids, true, 4, false); diff --git a/be/test/storage/iterator/vertical_block_reader_test.cpp b/be/test/storage/iterator/vertical_block_reader_test.cpp new file mode 100644 index 00000000000000..52ee75b2d2fcba --- /dev/null +++ b/be/test/storage/iterator/vertical_block_reader_test.cpp @@ -0,0 +1,136 @@ +// 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 "storage/iterator/vertical_block_reader.h" + +#include + +#include "storage/rowset/rowset_meta.h" + +namespace doris { + +class VerticalBlockReaderTestAccessor { +public: + static void append_grouped_iterator_init_flags(const RowsetMeta& rowset_meta, + std::pair segment_offsets, + size_t added_iterators, + std::vector* iterator_init_flags) { + VerticalBlockReader::_append_grouped_iterator_init_flags( + rowset_meta, segment_offsets, added_iterators, iterator_init_flags); + } + + static std::vector grouped_iterator_init_flags( + const RowsetMeta& rowset_meta, std::pair segment_offsets, + size_t added_iterators) { + std::vector iterator_init_flags; + append_grouped_iterator_init_flags(rowset_meta, segment_offsets, added_iterators, + &iterator_init_flags); + return iterator_init_flags; + } +}; + +TEST(VerticalBlockReaderTest, GroupedIteratorInitFlagsForWholeRowset) { + RowsetMeta rowset_meta; + rowset_meta.set_num_segments(5); + rowset_meta.set_segment_group_sizes({2, 2, 1}); + + EXPECT_EQ(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {0, 0}, 5), + std::vector({true, false, true, false, true})); +} + +TEST(VerticalBlockReaderTest, GroupedIteratorInitFlagsForSegmentRange) { + RowsetMeta rowset_meta; + rowset_meta.set_num_segments(8); + rowset_meta.set_segment_group_sizes({3, 2, 3}); + + EXPECT_EQ(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {0, 3}, 3), + std::vector({true, false, false})); + EXPECT_EQ(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {1, 3}, 2), + std::vector({true, false})); + EXPECT_EQ(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {1, 2}, 1), + std::vector({true})); + EXPECT_EQ(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {1, 4}, 3), + std::vector({true, false, true})); + EXPECT_EQ(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {3, 5}, 2), + std::vector({true, false})); + EXPECT_EQ(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {4, 7}, 3), + std::vector({true, true, false})); + EXPECT_EQ(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {5, 8}, 3), + std::vector({true, false, false})); + EXPECT_EQ(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {7, 8}, 1), + std::vector({true})); +} + +TEST(VerticalBlockReaderTest, GroupedIteratorInitFlagsAppendToExistingFlags) { + RowsetMeta rowset_meta; + rowset_meta.set_num_segments(5); + rowset_meta.set_segment_group_sizes({2, 2, 1}); + std::vector iterator_init_flags {false, true}; + + VerticalBlockReaderTestAccessor::append_grouped_iterator_init_flags(rowset_meta, {1, 4}, 3, + &iterator_init_flags); + EXPECT_EQ(iterator_init_flags, std::vector({false, true, true, true, false})); + + VerticalBlockReaderTestAccessor::append_grouped_iterator_init_flags(rowset_meta, {4, 5}, 1, + &iterator_init_flags); + EXPECT_EQ(iterator_init_flags, std::vector({false, true, true, true, false, true})); +} + +TEST(VerticalBlockReaderTest, GroupedIteratorInitFlagsRejectInvalidInput) { + RowsetMeta rowset_meta; + rowset_meta.set_num_segments(5); + rowset_meta.set_segment_group_sizes({2, 2, 1}); + + RowsetMetaPB invalid_group_layout_pb; + invalid_group_layout_pb.set_rowset_id(1); + invalid_group_layout_pb.set_num_segments(5); + invalid_group_layout_pb.add_segment_group_sizes(2); + invalid_group_layout_pb.add_segment_group_sizes(2); + RowsetMeta invalid_group_layout; + ASSERT_TRUE(invalid_group_layout.init_from_pb(invalid_group_layout_pb)); + +#ifndef NDEBUG + GTEST_FLAG_SET(death_test_style, "threadsafe"); + EXPECT_DEATH( + { + static_cast(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags( + rowset_meta, {1, 4}, 2)); + }, + "Check failed"); + EXPECT_DEATH( + { + static_cast(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags( + rowset_meta, {0, 6}, 6)); + }, + "Check failed"); + EXPECT_DEATH( + { + static_cast(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags( + invalid_group_layout, {0, 0}, 5)); + }, + "Check failed"); +#else + EXPECT_ANY_THROW(static_cast( + VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {1, 4}, 2))); + EXPECT_ANY_THROW(static_cast( + VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {0, 6}, 6))); + EXPECT_ANY_THROW(static_cast(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags( + invalid_group_layout, {0, 0}, 5))); +#endif +} + +} // namespace doris diff --git a/be/test/storage/pb_convert_test.cpp b/be/test/storage/pb_convert_test.cpp index 9e03b023a4b74c..d7548b5b11cd15 100644 --- a/be/test/storage/pb_convert_test.cpp +++ b/be/test/storage/pb_convert_test.cpp @@ -374,6 +374,30 @@ TEST(PbConvert, test_rvalue_overloads) { << "\n diff=" << print(set_diff(tablet_meta_cloud_set_fields, tablet_meta_out_set_fields)); } +TEST(PbConvert, test_rowset_segment_group_sizes_conversion) { + RowsetMetaPB rs; + rs.set_segments_overlap_pb(NONOVERLAPPING_WITHIN_GROUP); + rs.add_segment_group_sizes(3); + rs.add_segment_group_sizes(2); + rs.add_segment_group_sizes(5); + + RowsetMetaCloudPB rs_cloud; + doris_rowset_meta_to_cloud(&rs_cloud, rs); + ASSERT_EQ(rs_cloud.segments_overlap_pb(), NONOVERLAPPING_WITHIN_GROUP); + ASSERT_EQ(rs_cloud.segment_group_sizes_size(), 3); + EXPECT_EQ(rs_cloud.segment_group_sizes(0), 3); + EXPECT_EQ(rs_cloud.segment_group_sizes(1), 2); + EXPECT_EQ(rs_cloud.segment_group_sizes(2), 5); + + RowsetMetaPB rs_out; + cloud_rowset_meta_to_doris(&rs_out, rs_cloud); + ASSERT_EQ(rs_out.segments_overlap_pb(), NONOVERLAPPING_WITHIN_GROUP); + ASSERT_EQ(rs_out.segment_group_sizes_size(), 3); + EXPECT_EQ(rs_out.segment_group_sizes(0), 3); + EXPECT_EQ(rs_out.segment_group_sizes(1), 2); + EXPECT_EQ(rs_out.segment_group_sizes(2), 5); +} + TEST(PbConvert, test_return_value_overloads) { // rowset meta RowsetMetaPB rs; diff --git a/be/test/storage/rowid_conversion_test.cpp b/be/test/storage/rowid_conversion_test.cpp index c8b3bad9336a04..2df224f2766c95 100644 --- a/be/test/storage/rowid_conversion_test.cpp +++ b/be/test/storage/rowid_conversion_test.cpp @@ -29,11 +29,16 @@ #include #include +#include #include #include #include #include +#include "cloud/cloud_cumulative_compaction.h" +#include "cloud/cloud_storage_engine.h" +#include "cloud/cloud_tablet.h" +#include "common/config.h" #include "common/status.h" #include "core/block/block.h" #include "core/block/column_with_type_and_name.h" @@ -43,7 +48,10 @@ #include "io/io_common.h" #include "json2pb/json_to_pb.h" #include "runtime/exec_env.h" +#include "storage/compaction_task_tracker.h" #include "storage/delete/delete_handler.h" +#include "storage/index/index_writer.h" +#include "storage/index/inverted/inverted_index_desc.h" #include "storage/merger.h" #include "storage/options.h" #include "storage/rowset/beta_rowset.h" @@ -54,10 +62,12 @@ #include "storage/rowset/rowset_reader_context.h" #include "storage/rowset/rowset_writer.h" #include "storage/rowset/rowset_writer_context.h" +#include "storage/segment/segment.h" #include "storage/storage_engine.h" #include "storage/tablet/tablet.h" #include "storage/tablet/tablet_meta.h" #include "storage/tablet/tablet_schema.h" +#include "util/defer_op.h" #include "util/uid_util.h" namespace doris { @@ -79,6 +89,14 @@ class TestRowIdConversion : public testing::TestWithParamcreate_directory(absolute_dir + "/tablet_path") .ok()); + + std::vector tmp_paths; + tmp_paths.emplace_back(absolute_dir, 1024000000); + auto tmp_file_dirs = std::make_unique(tmp_paths); + st = tmp_file_dirs->init(); + ASSERT_TRUE(st.ok()) << st; + ExecEnv::GetInstance()->set_tmp_file_dir(std::move(tmp_file_dirs)); + doris::EngineOptions options; auto engine = std::make_unique(options); engine_ref = engine.get(); @@ -86,12 +104,14 @@ class TestRowIdConversion : public testing::TestWithParamset_tmp_file_dir(nullptr); EXPECT_TRUE(io::global_local_filesystem()->delete_directory(absolute_dir).ok()); engine_ref = nullptr; ExecEnv::GetInstance()->set_storage_engine(nullptr); } - TabletSchemaSPtr create_schema(KeysType keys_type = DUP_KEYS) { + TabletSchemaSPtr create_schema(KeysType keys_type = DUP_KEYS, + bool with_inverted_index = false) { TabletSchemaSPtr tablet_schema = std::make_shared(); TabletSchemaPB tablet_schema_pb; tablet_schema_pb.set_keys_type(keys_type); @@ -135,6 +155,15 @@ class TestRowIdConversion : public testing::TestWithParamset_is_bf_column(false); } + if (with_inverted_index) { + tablet_schema_pb.set_inverted_index_storage_format(InvertedIndexStorageFormatPB::V2); + auto* index = tablet_schema_pb.add_index(); + index->set_index_id(1); + index->set_index_name("c2_idx"); + index->set_index_type(IndexType::INVERTED); + index->add_col_unique_id(2); + } + tablet_schema->init_from_pb(tablet_schema_pb); return tablet_schema; } @@ -185,7 +214,7 @@ class TestRowIdConversion : public testing::TestWithParam>> input_data; + for (int64_t segment_id = 0; segment_id < num_segments; ++segment_id) { + std::vector> segment_data; + for (int64_t row_id = 0; row_id < rows_per_segment; ++row_id) { + int64_t key = segment_id * rows_per_segment + row_id; + segment_data.emplace_back(key, key + 1); + } + input_data.push_back(std::move(segment_data)); + } + + CloudStorageEngine cloud_engine(EngineOptions {}); + for (bool is_vertical : {false, true}) { + SCOPED_TRACE(is_vertical ? "vertical merge" : "horizontal merge"); + + TabletSchemaSPtr tablet_schema = create_schema(UNIQUE_KEYS, true); + tablet_schema->set_schema_version(schema_version); + tablet_schema->set_db_id(1000); + RowsetSharedPtr input_rowset = create_rowset(tablet_schema, OVERLAPPING, input_data, 2); + ASSERT_TRUE(input_rowset != nullptr); + + TabletSharedPtr local_tablet = create_tablet(*tablet_schema, true); + auto writer_context = create_rowset_writer_context( + tablet_schema, NONOVERLAPPING, rows_per_segment, input_rowset->version()); + writer_context.tablet_id = local_tablet->tablet_id(); + writer_context.index_id = local_tablet->index_id(); + writer_context.partition_id = local_tablet->partition_id(); + writer_context.tablet_schema_hash = local_tablet->schema_hash(); + writer_context.tablet_uid = local_tablet->tablet_uid(); + writer_context.newest_write_timestamp = newest_write_timestamp; + writer_context.compaction_level = compaction_level; + writer_context.enable_unique_key_merge_on_write = true; + auto writer_result = + RowsetFactory::create_rowset_writer(*engine_ref, writer_context, is_vertical); + ASSERT_TRUE(writer_result.has_value()) << writer_result.error(); + + auto cloud_tablet = std::make_shared( + cloud_engine, std::make_shared(*local_tablet->tablet_meta())); + CloudCumulativeCompaction compaction(cloud_engine, cloud_tablet); + compaction._input_rowsets = {input_rowset}; + compaction._cur_tablet_schema = tablet_schema; + compaction._output_rs_writer = std::move(writer_result).value(); + compaction._is_vertical = is_vertical; + compaction._input_row_num = input_rowset->num_rows(); + compaction._input_rowsets_data_size = input_rowset->data_disk_size(); + compaction._stats.rowid_conversion = compaction._rowid_conversion.get(); + + auto* compaction_task_tracker = CompactionTaskTracker::instance(); + CompactionTaskInfo task_info; + task_info.compaction_id = compaction.compaction_id(); + compaction_task_tracker->register_task(std::move(task_info)); + Defer remove_tracker_task { + [compaction_task_tracker, compaction_id = compaction.compaction_id()] { + compaction_task_tracker->remove_task(compaction_id); + }}; + + Compaction::MergeInputRowsetsResult merge_result; + merge_result.is_segment_grouped = true; + merge_result.segment_group_size = segment_group_size; + ASSERT_TRUE(compaction.do_merge_input_rowsets({}, &merge_result).ok()); + const int64_t segment_group_count = + (num_segments + segment_group_size - 1) / segment_group_size; + ASSERT_EQ(merge_result.output_segment_group_sizes.size(), + static_cast(segment_group_count)); + for (const auto output_group_size : merge_result.output_segment_group_sizes) { + EXPECT_GT(output_group_size, 0); + } + if (is_vertical) { + constexpr int32_t default_num_columns_per_group = 5; + const int32_t num_columns_per_group = + config::vertical_compaction_num_columns_per_group != + default_num_columns_per_group + ? config::vertical_compaction_num_columns_per_group + : cloud_tablet->tablet_meta() + ->vertical_compaction_num_columns_per_group(); + std::vector> column_groups; + std::vector key_group_cluster_key_idxes; + Merger::vertical_split_columns(*tablet_schema, &column_groups, + &key_group_cluster_key_idxes, num_columns_per_group); + + const auto tracked_tasks = compaction_task_tracker->get_all_tasks(); + const auto task_it = std::find_if( + tracked_tasks.begin(), tracked_tasks.end(), [&](const auto& tracked_task) { + return tracked_task.compaction_id == compaction.compaction_id(); + }); + ASSERT_NE(task_it, tracked_tasks.end()); + const int64_t expected_total_groups = + static_cast(column_groups.size()) * segment_group_count; + EXPECT_EQ(task_it->vertical_total_groups, expected_total_groups); + EXPECT_EQ(task_it->vertical_completed_groups, expected_total_groups); + } + + RowsetSharedPtr output_rowset; + ASSERT_EQ(Status::OK(), compaction._output_rs_writer->build(output_rowset)); + ASSERT_TRUE(output_rowset != nullptr); + compaction._output_rowset = output_rowset; + compaction.update_output_rowset_after_build(merge_result); + EXPECT_EQ(compaction._stats.output_rows, input_rowset->num_rows()); + if (is_vertical) { + EXPECT_GT(output_rowset->num_segments(), + static_cast(merge_result.output_segment_group_sizes.size())); + } + EXPECT_EQ(output_rowset->rowset_meta()->get_num_segment_rows().size(), + output_rowset->num_segments()); + + EXPECT_EQ(output_rowset->rowset_meta()->segments_overlap(), NONOVERLAPPING_WITHIN_GROUP); + const auto output_rowset_pb = output_rowset->rowset_meta()->get_rowset_pb(); + ASSERT_EQ(output_rowset_pb.segment_group_sizes_size(), + merge_result.output_segment_group_sizes.size()); + int64_t output_segment_count = 0; + for (int i = 0; i < output_rowset_pb.segment_group_sizes_size(); ++i) { + EXPECT_EQ(output_rowset_pb.segment_group_sizes(i), + merge_result.output_segment_group_sizes[static_cast(i)]); + output_segment_count += output_rowset_pb.segment_group_sizes(i); + } + EXPECT_EQ(output_segment_count, output_rowset->num_segments()); + + RowsetReaderContext reader_context; + reader_context.tablet_schema = tablet_schema; + reader_context.need_ordered_result = false; + std::vector return_columns = {0, 1}; + reader_context.return_columns = &return_columns; + RowsetReaderSharedPtr output_reader; + create_and_init_rowset_reader(output_rowset.get(), reader_context, &output_reader); + + std::vector> output_data; + Status read_status; + do { + Block output_block = tablet_schema->create_block(); + read_status = output_reader->next_batch(&output_block); + const auto& columns = output_block.get_columns_with_type_and_name(); + ASSERT_EQ(columns.size(), return_columns.size()); + for (size_t row_id = 0; row_id < output_block.rows(); ++row_id) { + output_data.emplace_back(columns[0].column->get_int(row_id), + columns[1].column->get_int(row_id)); + } + } while (read_status.ok()); + ASSERT_TRUE(read_status.is()) << read_status; + ASSERT_EQ(output_data.size(), input_rowset->num_rows()); + + auto beta_rowset = std::dynamic_pointer_cast(output_rowset); + ASSERT_TRUE(beta_rowset != nullptr); + const auto& rowset_meta = output_rowset->rowset_meta(); + const auto rowset_meta_pb = rowset_meta->get_rowset_pb(); + const auto& segment_num_rows_from_meta = rowset_meta->get_num_segment_rows(); + + EXPECT_EQ(rowset_meta_pb.rowset_id_v2(), writer_context.rowset_id.to_string()); + EXPECT_EQ(rowset_meta->rowset_id().to_string(), writer_context.rowset_id.to_string()); + + EXPECT_EQ(rowset_meta->rowset_type(), BETA_ROWSET); + EXPECT_EQ(rowset_meta->rowset_state(), VISIBLE); + EXPECT_TRUE(rowset_meta->has_version()); + EXPECT_EQ(rowset_meta->version(), input_rowset->version()); + EXPECT_FALSE(rowset_meta->empty()); + EXPECT_EQ(rowset_meta->num_rows(), input_rowset->num_rows()); + EXPECT_EQ(rowset_meta->segments_overlap(), NONOVERLAPPING_WITHIN_GROUP); + EXPECT_TRUE(rowset_meta->is_segments_overlapping()); + + ASSERT_TRUE(rowset_meta->tablet_schema() != nullptr); + EXPECT_EQ(*rowset_meta->tablet_schema(), *tablet_schema); + EXPECT_TRUE(rowset_meta_pb.has_tablet_schema()); + TabletSchemaPB expected_tablet_schema_pb; + tablet_schema->to_schema_pb(&expected_tablet_schema_pb); + EXPECT_EQ(rowset_meta_pb.tablet_schema().SerializeAsString(), + expected_tablet_schema_pb.SerializeAsString()); + EXPECT_EQ(rowset_meta_pb.schema_version(), schema_version); + EXPECT_TRUE(rowset_meta_pb.has_has_variant_type_in_schema()); + EXPECT_FALSE(rowset_meta_pb.has_variant_type_in_schema()); + + std::vector output_segments; + ASSERT_TRUE(beta_rowset->load_segments(&output_segments).ok()); + ASSERT_EQ(rowset_meta->num_segments(), output_segments.size()); + ASSERT_EQ(segment_num_rows_from_meta.size(), output_segments.size()); + + const auto& segment_key_bounds_from_meta = rowset_meta->get_segments_key_bounds(); + EXPECT_FALSE(rowset_meta->is_segments_key_bounds_aggregated()); + EXPECT_FALSE(rowset_meta->is_segments_key_bounds_truncated()); + ASSERT_EQ(segment_key_bounds_from_meta.size(), output_segments.size()); + + const auto& inverted_index_file_info_from_meta = rowset_meta->inverted_index_file_info(); + EXPECT_TRUE(rowset_meta_pb.enable_inverted_index_file_info()); + ASSERT_EQ(inverted_index_file_info_from_meta.size(), output_segments.size()); + + EXPECT_FALSE(rowset_meta_pb.enable_segments_file_size()); + EXPECT_TRUE(rowset_meta_pb.segments_file_size().empty()); + + int64_t actual_data_disk_size = 0; + int64_t actual_index_disk_size = 0; + int64_t actual_num_rows = 0; + for (size_t segment_id = 0; segment_id < output_segments.size(); ++segment_id) { + EXPECT_EQ(output_segments[segment_id]->id(), segment_id); + EXPECT_EQ(segment_num_rows_from_meta[segment_id], + output_segments[segment_id]->num_rows()) + << "segment_id=" << segment_id; + EXPECT_EQ(segment_key_bounds_from_meta[segment_id].min_key(), + output_segments[segment_id]->min_key()) + << "segment_id=" << segment_id; + EXPECT_EQ(segment_key_bounds_from_meta[segment_id].max_key(), + output_segments[segment_id]->max_key()) + << "segment_id=" << segment_id; + + const auto& index_file_info = inverted_index_file_info_from_meta[segment_id]; + ASSERT_TRUE(index_file_info.has_index_size()) << "segment_id=" << segment_id; + const auto segment_path = output_rowset->segment_path(segment_id); + ASSERT_TRUE(segment_path.has_value()) << segment_path.error(); + int64_t segment_file_size = 0; + const auto segment_file_size_status = + rowset_meta->fs()->file_size(segment_path.value(), &segment_file_size); + ASSERT_TRUE(segment_file_size_status.ok()) << segment_file_size_status; + actual_data_disk_size += segment_file_size; + actual_num_rows += output_segments[segment_id]->num_rows(); + + const auto index_file_path = + segment_v2::InvertedIndexDescriptor::get_index_file_path_v2( + segment_v2::InvertedIndexDescriptor::get_index_file_path_prefix( + segment_path.value())); + int64_t index_file_size = 0; + const auto index_file_size_status = + rowset_meta->fs()->file_size(index_file_path, &index_file_size); + ASSERT_TRUE(index_file_size_status.ok()) << index_file_size_status; + EXPECT_EQ(index_file_info.index_size(), index_file_size) << "segment_id=" << segment_id; + actual_index_disk_size += index_file_size; + } + EXPECT_EQ(rowset_meta->num_rows(), actual_num_rows); + EXPECT_EQ(rowset_meta->data_disk_size(), actual_data_disk_size); + EXPECT_EQ(rowset_meta->index_disk_size(), actual_index_disk_size); + EXPECT_EQ(rowset_meta->total_disk_size(), + rowset_meta->data_disk_size() + rowset_meta->index_disk_size()); + + std::vector output_segment_num_rows; + OlapReaderStatistics reader_stats; + ASSERT_TRUE( + beta_rowset->get_segment_num_rows(&output_segment_num_rows, false, &reader_stats) + .ok()); + + RowIdConversion& rowid_conversion = *compaction._stats.rowid_conversion; + EXPECT_EQ(rowid_conversion.get_src_segment_to_id_map().size(), num_segments); + EXPECT_EQ(rowid_conversion.get_rowid_conversion_map().size(), num_segments); + EXPECT_EQ(rowid_conversion.get_rowid_conversion_map().size(), + rowid_conversion.get_src_segment_to_id_map().size()); + for (int64_t segment_id = 0; segment_id < num_segments; ++segment_id) { + for (int64_t row_id = 0; row_id < rows_per_segment; ++row_id) { + RowLocation src(input_rowset->rowset_id(), segment_id, row_id); + RowLocation dst; + ASSERT_EQ(rowid_conversion.get(src, &dst), 0) + << "segment_id=" << segment_id << ", row_id=" << row_id; + ASSERT_LT(dst.segment_id, output_segment_num_rows.size()); + ASSERT_LT(dst.row_id, output_segment_num_rows[dst.segment_id]); + + size_t output_row_id = dst.row_id; + for (uint32_t output_segment_id = 0; output_segment_id < dst.segment_id; + ++output_segment_id) { + output_row_id += output_segment_num_rows[output_segment_id]; + } + ASSERT_LT(output_row_id, output_data.size()); + EXPECT_EQ(output_data[output_row_id], input_data[segment_id][row_id]); + } + } + + if (is_vertical) { + auto second_writer_context = create_rowset_writer_context( + tablet_schema, NONOVERLAPPING, rows_per_segment, output_rowset->version()); + second_writer_context.tablet_id = local_tablet->tablet_id(); + second_writer_context.index_id = local_tablet->index_id(); + second_writer_context.partition_id = local_tablet->partition_id(); + second_writer_context.tablet_schema_hash = local_tablet->schema_hash(); + second_writer_context.tablet_uid = local_tablet->tablet_uid(); + second_writer_context.newest_write_timestamp = newest_write_timestamp; + second_writer_context.compaction_level = compaction_level; + second_writer_context.enable_unique_key_merge_on_write = true; + auto second_writer_result = + RowsetFactory::create_rowset_writer(*engine_ref, second_writer_context, true); + ASSERT_TRUE(second_writer_result.has_value()) << second_writer_result.error(); + + CloudCumulativeCompaction second_compaction(cloud_engine, cloud_tablet); + second_compaction._input_rowsets = {output_rowset}; + second_compaction._cur_tablet_schema = tablet_schema; + second_compaction._output_rs_writer = std::move(second_writer_result).value(); + second_compaction._is_vertical = true; + second_compaction._input_row_num = output_rowset->num_rows(); + second_compaction._input_rowsets_data_size = output_rowset->data_disk_size(); + second_compaction._stats.rowid_conversion = second_compaction._rowid_conversion.get(); + + Compaction::MergeInputRowsetsResult second_merge_result; + second_merge_result.is_segment_grouped = true; + second_merge_result.segment_group_size = output_rowset->num_segments(); + ASSERT_TRUE(second_compaction.do_merge_input_rowsets({}, &second_merge_result).ok()); + ASSERT_EQ(second_merge_result.output_segment_group_sizes.size(), 1); + + RowsetSharedPtr second_output_rowset; + ASSERT_EQ(Status::OK(), + second_compaction._output_rs_writer->build(second_output_rowset)); + ASSERT_TRUE(second_output_rowset != nullptr); + second_compaction._output_rowset = second_output_rowset; + second_compaction.update_output_rowset_after_build(second_merge_result); + EXPECT_EQ(second_output_rowset->rowset_meta()->segments_overlap(), NONOVERLAPPING); + EXPECT_TRUE(second_output_rowset->rowset_meta()->segment_group_sizes().empty()); + EXPECT_EQ(second_output_rowset->num_rows(), output_rowset->num_rows()); + EXPECT_GT(second_output_rowset->num_segments(), 1); + + auto second_beta_rowset = std::dynamic_pointer_cast(second_output_rowset); + ASSERT_TRUE(second_beta_rowset != nullptr); + std::vector second_output_segments; + ASSERT_TRUE(second_beta_rowset->load_segments(&second_output_segments).ok()); + ASSERT_EQ(second_output_segments.size(), second_output_rowset->num_segments()); + for (size_t segment_id = 0; segment_id < second_output_segments.size(); ++segment_id) { + const auto& segment = second_output_segments[segment_id]; + EXPECT_LE(segment->min_key(), segment->max_key()) << "segment_id=" << segment_id; + if (segment_id > 0) { + EXPECT_LT(second_output_segments[segment_id - 1]->max_key(), segment->min_key()) + << "previous_segment_id=" << segment_id - 1 + << ", segment_id=" << segment_id; + } + } + + RowsetReaderContext second_reader_context; + second_reader_context.tablet_schema = tablet_schema; + second_reader_context.need_ordered_result = false; + second_reader_context.return_columns = &return_columns; + RowsetReaderSharedPtr second_output_reader; + create_and_init_rowset_reader(second_output_rowset.get(), second_reader_context, + &second_output_reader); + + std::vector> second_output_data; + Status second_read_status; + do { + Block output_block = tablet_schema->create_block(); + second_read_status = second_output_reader->next_batch(&output_block); + const auto& columns = output_block.get_columns_with_type_and_name(); + for (size_t row_id = 0; row_id < output_block.rows(); ++row_id) { + second_output_data.emplace_back(columns[0].column->get_int(row_id), + columns[1].column->get_int(row_id)); + } + } while (second_read_status.ok()); + ASSERT_TRUE(second_read_status.is()) << second_read_status; + EXPECT_EQ(second_output_data, output_data); + } + } +} + INSTANTIATE_TEST_SUITE_P( Parameters, TestRowIdConversion, ::testing::ValuesIn(std::vector> { diff --git a/be/test/storage/rowset/rowset_meta_test.cpp b/be/test/storage/rowset/rowset_meta_test.cpp index 954b8e9ed22515..010791c69b9972 100644 --- a/be/test/storage/rowset/rowset_meta_test.cpp +++ b/be/test/storage/rowset/rowset_meta_test.cpp @@ -124,6 +124,27 @@ TEST_F(RowsetMetaTest, TestRowsetIdInit) { EXPECT_EQ(id.to_string(), "72057594037927935"); } +TEST_F(RowsetMetaTest, SegmentGroupsKeepRowsetOverlappingSemantics) { + RowsetMeta rowset_meta; + rowset_meta.set_version({10, 10}); + rowset_meta.set_num_segments(5); + rowset_meta.set_segments_overlap(NONOVERLAPPING_WITHIN_GROUP); + rowset_meta.set_segment_group_sizes({2, 3}); + + EXPECT_TRUE(rowset_meta.is_segments_overlapping()); + EXPECT_TRUE(rowset_meta.produced_by_compaction()); + EXPECT_EQ(rowset_meta.get_compaction_score(), 5); + EXPECT_EQ(rowset_meta.get_merge_way_num(), 5); + + auto rowset_meta_pb = rowset_meta.get_rowset_pb(); + ASSERT_EQ(rowset_meta_pb.segment_group_sizes_size(), 2); + EXPECT_EQ(rowset_meta_pb.segment_group_sizes(0), 2); + EXPECT_EQ(rowset_meta_pb.segment_group_sizes(1), 3); + + rowset_meta.clear_segment_group_sizes(); + EXPECT_EQ(rowset_meta.get_rowset_pb().segment_group_sizes_size(), 0); +} + TEST_F(RowsetMetaTest, TestNumSegmentRowsSetAndGet) { RowsetMeta rowset_meta; EXPECT_TRUE(rowset_meta.init_from_json(_json_rowset_meta)); diff --git a/gensrc/proto/olap_file.proto b/gensrc/proto/olap_file.proto index 498f0238efa303..4866aa955296ac 100644 --- a/gensrc/proto/olap_file.proto +++ b/gensrc/proto/olap_file.proto @@ -53,6 +53,9 @@ enum SegmentsOverlapPB { OVERLAP_UNKNOWN = 0; // this enum is added since Doris v0.11, so previous rowset's segment is unknown OVERLAPPING = 1; NONOVERLAPPING = 2; + // Segments are non-overlapping within each group, but groups may overlap with each other. + // The group layout is described by RowsetMetaPB.segment_group_sizes. + NONOVERLAPPING_WITHIN_GROUP = 3; } message KeyBoundsPB { @@ -145,6 +148,10 @@ message RowsetMetaPB { // Only applies to non-MOW rowsets to reduce meta size on cloud FDB. optional bool segments_key_bounds_aggregated = 57; + // Valid only when segments_overlap_pb is NONOVERLAPPING_WITHIN_GROUP. + // Each value is the number of consecutive output segments in one non-overlapping group. + repeated int32 segment_group_sizes = 59; + // For cloud // for data recycling optional int64 txn_expiration = 1000; @@ -254,6 +261,10 @@ message RowsetMetaCloudPB { // Only applies to non-MOW rowsets to reduce meta size on cloud FDB. optional bool segments_key_bounds_aggregated = 57; + // Valid only when segments_overlap_pb is NONOVERLAPPING_WITHIN_GROUP. + // Each value is the number of consecutive output segments in one non-overlapping group. + repeated int32 segment_group_sizes = 59; + // cloud // the field is a vector, rename it repeated int64 segments_file_size = 100; diff --git a/regression-test/suites/cloud_p0/compaction/test_cloud_single_rowset_grouped_compaction.groovy b/regression-test/suites/cloud_p0/compaction/test_cloud_single_rowset_grouped_compaction.groovy new file mode 100644 index 00000000000000..b98d51c243822a --- /dev/null +++ b/regression-test/suites/cloud_p0/compaction/test_cloud_single_rowset_grouped_compaction.groovy @@ -0,0 +1,563 @@ +// 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. + +import java.util.Base64 + +suite("test_cloud_single_rowset_grouped_compaction", "nonConcurrent") { + if (!isCloudMode()) { + logger.info("not cloud mode, skip this test") + return + } + + int initialInputSegmentsPerGroup = 2 + long compactionTimeoutMs = 90000L + def customBeConfig = [ + doris_scanner_row_bytes: 1, + enable_cloud_single_rowset_compaction: true, + cloud_single_rowset_compaction_min_segments: 2, + cloud_single_rowset_compaction_segment_group_size: initialInputSegmentsPerGroup, + cumulative_compaction_min_deltas: 2, + enable_aggregate_non_mow_key_bounds: false, + disable_auto_compaction: true, + vertical_compaction_max_segment_size: 1073741824, + compaction_batch_size: -1, + enable_rowid_conversion_correctness_check: false + ] + + setBeConfigTemporary(customBeConfig) { + def metaServiceEndpoint = context.config.metaServiceHttpAddress + + def getRowsetMeta = { tabletId, int version -> + def rowsetMeta = null + getSegmentFilesFromMs(metaServiceEndpoint, tabletId, version) { responseCode, body -> + assertEquals(200, responseCode) + rowsetMeta = parseJson(body) + } + assertNotNull(rowsetMeta) + return rowsetMeta + } + + def compareEncodedKeys = { String leftBase64, String rightBase64 -> + byte[] left = Base64.getDecoder().decode(leftBase64) + byte[] right = Base64.getDecoder().decode(rightBase64) + int commonLength = Math.min(left.length, right.length) + for (int i = 0; i < commonLength; ++i) { + int leftByte = Byte.toUnsignedInt(left[i]) + int rightByte = Byte.toUnsignedInt(right[i]) + if (leftByte != rightByte) { + return leftByte <=> rightByte + } + } + return left.length <=> right.length + } + + def logCaseSection = { String description -> + logger.info("========================================================================") + logger.info(description) + logger.info("========================================================================") + } + + def showTablet = { tableName, beHost, bePort, tabletId -> + sql "SELECT COUNT(*) FROM ${tableName}" + def (code, out, err) = be_show_tablet_status(beHost, bePort, tabletId) + logger.info("Show tablet status: code=${code}, out=${out}, err=${err}") + assertEquals(0, code) + return parseJson(out.trim()) + } + + def rowsetByVersion = { tabletJson, int version -> + def rowset = tabletJson.rowsets.find { it.startsWith("[${version}-${version}] ") } + assertNotNull(rowset) + return rowset + } + + def parseRowsetInfo = { rowset -> + def matcher = rowset =~ /\[[0-9]+-[0-9]+\]\s+([0-9]+)\s+DATA\s+([A-Z_]+)/ + assertTrue(matcher.find(), "unexpected rowset format: ${rowset}") + return [segments: matcher.group(1).toInteger(), overlap: matcher.group(2)] + } + + def waitForCompaction = { beHost, bePort, tabletId, timeoutMs -> + long deadline = System.currentTimeMillis() + timeoutMs + def lastStatus = null + while (System.currentTimeMillis() < deadline) { + def (code, out, err) = be_get_compaction_status(beHost, bePort, tabletId) + logger.info("Get compaction status: code=${code}, out=${out}, err=${err}") + assertEquals(0, code) + lastStatus = parseJson(out.trim()) + assertEquals("success", lastStatus.status.toLowerCase()) + if (!lastStatus.run_status) { + return + } + Thread.sleep(1000) + } + assertTrue(false, "compaction did not finish on ${beHost}:${bePort}, " + + "tablet=${tabletId}, timeoutMs=${timeoutMs}, last=${lastStatus}") + } + + def runCumulativeCompaction = { beHost, bePort, tabletId -> + def (code, out, err) = be_run_cumulative_compaction(beHost, bePort, tabletId) + logger.info("Run compaction: code=${code}, out=${out}, err=${err}") + assertEquals(0, code) + def compactJson = parseJson(out.trim()) + assertEquals("success", compactJson.status.toLowerCase()) + waitForCompaction(beHost, bePort, tabletId, compactionTimeoutMs) + } + + def readPointRows = { String tableName -> + def pointResult = sql "SELECT k, v FROM ${tableName} WHERE k = 100 ORDER BY v" + return pointResult.collect { row -> row.collect { column -> column.toString() } } + } + + def checkPointRows = { String tableName, List> expectedPointRows -> + def pointRows = readPointRows(tableName) + if (expectedPointRows != null) { + assertEquals(expectedPointRows, pointRows) + } + return pointRows + } + + def checkSingleRowsetGroupedCompaction = { + String tableName, String keyType, String valueColumn, String extraProperties, + int inputSegmentsPerGroup, int rowsPerLoadRound, int expectedRows, + List> expectedPointRows, + boolean expectMultipleOutputSegmentsPerGroup -> + sql "DROP TABLE IF EXISTS ${tableName}" + sql """ + CREATE TABLE ${tableName} ( + k INT, + ${valueColumn} + ) + ${keyType}(k) + DISTRIBUTED BY HASH(k) BUCKETS 1 + PROPERTIES ( + "replication_num" = "1", + "disable_auto_compaction" = "true"${extraProperties} + ) + """ + + StringBuilder content = new StringBuilder() + for (int i = 0; i < 2; i++) { + (1..rowsPerLoadRound).each { + content.append("${it},${it + i}\n") + } + } + streamLoad { + table "${tableName}" + set "column_separator", "," + inputStream new ByteArrayInputStream(content.toString().getBytes()) + time 30000 + check { result, exception, startTime, endTime -> + if (exception != null) { + throw exception + } + def json = parseJson(result) + assertEquals("success", json.Status.toLowerCase()) + assertEquals(rowsPerLoadRound * 2, json.NumberTotalRows) + assertEquals(0, json.NumberFilteredRows) + } + } + sql "sync" + + def pointRowsBeforeCompaction = checkPointRows(tableName, expectedPointRows) + if (expectedPointRows == null) { + assertEquals(1, pointRowsBeforeCompaction.size()) + } + + def tablets = sql_return_maparray "SHOW TABLETS FROM ${tableName}" + assertEquals(1, tablets.size()) + def tabletId = tablets[0].TabletId + def backendId = tablets[0].BackendId + def backends = sql_return_maparray "SHOW BACKENDS" + def backend = backends.find { it.BackendId == backendId } + assertNotNull(backend) + + def before = showTablet(tableName, backend.Host, backend.HttpPort, tabletId) + def inputRowset = rowsetByVersion(before, 2) + def inputInfo = parseRowsetInfo(inputRowset) + assertEquals("OVERLAPPING", inputInfo.overlap) + assertTrue(inputInfo.segments > inputSegmentsPerGroup, inputRowset) + + set_be_param("cloud_single_rowset_compaction_segment_group_size", + inputSegmentsPerGroup.toString()) + runCumulativeCompaction(backend.Host, backend.HttpPort, tabletId) + + def after = showTablet(tableName, backend.Host, backend.HttpPort, tabletId) + def outputRowset = rowsetByVersion(after, 2) + def outputInfo = parseRowsetInfo(outputRowset) + assertEquals("NONOVERLAPPING_WITHIN_GROUP", outputInfo.overlap) + def expectedOutputSegments = + (inputInfo.segments + inputSegmentsPerGroup - 1) + .intdiv(inputSegmentsPerGroup) + if (expectMultipleOutputSegmentsPerGroup) { + assertTrue(outputInfo.segments > expectedOutputSegments, outputRowset) + } else { + assertEquals(expectedOutputSegments, outputInfo.segments) + } + def countResult = sql "SELECT COUNT(*) FROM ${tableName}" + assertEquals(expectedRows, countResult[0][0]) + def pointRowsAfterCompaction = checkPointRows(tableName, expectedPointRows) + if (expectedPointRows == null && pointRowsAfterCompaction != pointRowsBeforeCompaction) { + logger.warn("Point query result changed after single rowset grouped compaction" + + ", table=${tableName}, before=${pointRowsBeforeCompaction}" + + ", after=${pointRowsAfterCompaction}") + } + + // Compact the grouped rowset as one range. This exercises VerticalBlockReader's + // NONOVERLAPPING_WITHIN_GROUP iterator initialization and must produce a fully + // non-overlapping rowset. + set_be_param("cloud_single_rowset_compaction_segment_group_size", + outputInfo.segments.toString()) + runCumulativeCompaction(backend.Host, backend.HttpPort, tabletId) + + def afterSecondCompaction = + showTablet(tableName, backend.Host, backend.HttpPort, tabletId) + def finalRowset = rowsetByVersion(afterSecondCompaction, 2) + def finalInfo = parseRowsetInfo(finalRowset) + assertEquals("NONOVERLAPPING", finalInfo.overlap) + if (expectMultipleOutputSegmentsPerGroup) { + assertTrue(finalInfo.segments > 1, finalRowset) + def finalRowsetMeta = getRowsetMeta(tabletId, 2) + assertEquals(finalInfo.segments, finalRowsetMeta.num_segments.toString().toInteger()) + assertFalse(finalRowsetMeta.segments_key_bounds_aggregated ?: false) + assertFalse(finalRowsetMeta.segments_key_bounds_truncated ?: false) + def segmentKeyBounds = finalRowsetMeta.segments_key_bounds + assertEquals(finalInfo.segments, segmentKeyBounds.size()) + segmentKeyBounds.eachWithIndex { keyBounds, int segmentId -> + assertTrue(compareEncodedKeys(keyBounds.min_key, keyBounds.max_key) <= 0, + "invalid key range at segment ${segmentId}") + if (segmentId > 0) { + def previousKeyBounds = segmentKeyBounds[segmentId - 1] + assertTrue(compareEncodedKeys( + previousKeyBounds.max_key, keyBounds.min_key) < 0, + "overlapping key ranges at segments ${segmentId - 1} and " + + "${segmentId}") + } + } + } + def finalCountResult = sql "SELECT COUNT(*) FROM ${tableName}" + assertEquals(expectedRows, finalCountResult[0][0]) + def finalPointRows = checkPointRows(tableName, expectedPointRows) + if (expectedPointRows == null && finalPointRows != pointRowsAfterCompaction) { + logger.warn("Point query result changed after repeated single rowset compaction" + + ", table=${tableName}, before=${pointRowsAfterCompaction}" + + ", after=${finalPointRows}") + } + } + + def sqlCacheOrigValue = sql("select @@enable_sql_cache")[0][0] + try { + sql "set enable_sql_cache=false" + GetDebugPoint().clearDebugPointsForAllBEs() + GetDebugPoint().enableDebugPointForAllBEs("MemTable.need_flush") + + // ==================================================== + logCaseSection("Standard cases: compact DUP, AGG, MOW, and MOR rowsets twice " + + "with 2 or 4 input segments per group") + [2, 4].each { int inputSegmentsPerGroup -> + checkSingleRowsetGroupedCompaction( + "test_cloud_single_rowset_grouped_compaction_g${inputSegmentsPerGroup}_dup", + "DUPLICATE KEY", "v INT", "", inputSegmentsPerGroup, 8192, 8192 * 2, + [["100", "100"], ["100", "101"]], false) + checkSingleRowsetGroupedCompaction( + "test_cloud_single_rowset_grouped_compaction_g${inputSegmentsPerGroup}_agg", + "AGGREGATE KEY", "v INT SUM", "", inputSegmentsPerGroup, 8192, 8192, + [["100", "201"]], false) + checkSingleRowsetGroupedCompaction( + "test_cloud_single_rowset_grouped_compaction_g${inputSegmentsPerGroup}_mow", + "UNIQUE KEY", "v INT", + ", \"enable_unique_key_merge_on_write\" = \"true\"", + inputSegmentsPerGroup, 8192, 8192, [["100", "101"]], false) + checkSingleRowsetGroupedCompaction( + "test_cloud_single_rowset_grouped_compaction_g${inputSegmentsPerGroup}_mor", + "UNIQUE KEY", "v INT", + ", \"enable_unique_key_merge_on_write\" = \"false\"", + inputSegmentsPerGroup, 8192, 8192, null, false) + } + + // ==================================================== + logCaseSection("MOW delete-bitmap case: compact each of three overlapping rowsets " + + "twice and verify cross-rowset delete bitmap conversion") + sql "DROP TABLE IF EXISTS test_cloud_single_rowset_compaction_mow_delete_bitmap" + sql """ + CREATE TABLE test_cloud_single_rowset_compaction_mow_delete_bitmap ( + k INT, + v INT + ) + UNIQUE KEY(k) + DISTRIBUTED BY HASH(k) BUCKETS 1 + PROPERTIES ( + "replication_num" = "1", + "disable_auto_compaction" = "true", + "enable_unique_key_merge_on_write" = "true" + ) + """ + + def loadMowRows = { + int startKey, int endKey, int valueBase, int duplicateRounds -> + StringBuilder content = new StringBuilder() + (startKey..endKey).each { int key -> + content.append("${key},${valueBase + key}\n") + } + // Append an inner key range to create duplicates within this rowset. Keep the + // offsets deterministic so a failure can be reproduced exactly. + int duplicateStartKey = startKey + 2048 + int duplicateEndKey = endKey - 3072 + for (int round = 0; round < duplicateRounds; ++round) { + (duplicateStartKey..duplicateEndKey).each { int key -> + content.append("${key},${valueBase + key + 1}\n") + } + } + streamLoad { + table "test_cloud_single_rowset_compaction_mow_delete_bitmap" + set "column_separator", "," + inputStream new ByteArrayInputStream(content.toString().getBytes()) + time 30000 + check { result, exception, startTime, endTime -> + if (exception != null) { + throw exception + } + def json = parseJson(result) + assertEquals("success", json.Status.toLowerCase()) + int originalRows = endKey - startKey + 1 + int duplicateRows = duplicateEndKey - duplicateStartKey + 1 + assertEquals(originalRows + duplicateRows * duplicateRounds, + json.NumberTotalRows) + assertEquals(0, json.NumberFilteredRows) + } + } + sql "sync" + } + + loadMowRows(1, 16384, 100000, 1) + loadMowRows(8193, 24576, 200000, 2) + loadMowRows(12289, 28672, 300000, 3) + + def readMowRows = { + def rows = sql """ + SELECT k, v + FROM test_cloud_single_rowset_compaction_mow_delete_bitmap + WHERE k IN (1, 8192, 8193, 12288, 12289, 16384, 16385, 24576, 24577, 28672) + ORDER BY k + """ + return rows.collect { row -> row.collect { column -> column.toString() } } + } + def expectedMowRows = [ + ["1", "100001"], + ["8192", "108193"], + ["8193", "208193"], + ["12288", "212289"], + ["12289", "312289"], + ["16384", "316385"], + ["16385", "316386"], + ["24576", "324577"], + ["24577", "324578"], + ["28672", "328672"] + ] + assertEquals(expectedMowRows, readMowRows()) + def mowCountBeforeCompaction = + sql "SELECT COUNT(*) FROM test_cloud_single_rowset_compaction_mow_delete_bitmap" + assertEquals(28672, mowCountBeforeCompaction[0][0]) + + def mowTablets = sql_return_maparray( + "SHOW TABLETS FROM test_cloud_single_rowset_compaction_mow_delete_bitmap") + assertEquals(1, mowTablets.size()) + def mowTabletId = mowTablets[0].TabletId + def mowBackendId = mowTablets[0].BackendId + def mowBackends = sql_return_maparray "SHOW BACKENDS" + def mowBackend = mowBackends.find { it.BackendId == mowBackendId } + assertNotNull(mowBackend) + + def checkMowQueryResult = { + def countResult = + sql "SELECT COUNT(*) FROM test_cloud_single_rowset_compaction_mow_delete_bitmap" + assertEquals(mowCountBeforeCompaction, countResult) + assertEquals(expectedMowRows, readMowRows()) + } + + def mowRowsetVersions = [2, 3, 4] + def mowMaxSegmentSizeByVersion = [2: 32768, 3: 8192, 4: 2048] + def finalMowSegmentCounts = [:] + def compactMowRowsetTwice = { int targetVersion -> + logger.info("Compact MOW rowset version ${targetVersion} twice") + def beforeFirstCompaction = showTablet( + "test_cloud_single_rowset_compaction_mow_delete_bitmap", + mowBackend.Host, mowBackend.HttpPort, mowTabletId) + def inputRowset = rowsetByVersion(beforeFirstCompaction, targetVersion) + def inputInfo = parseRowsetInfo(inputRowset) + assertEquals("OVERLAPPING", inputInfo.overlap) + assertTrue(inputInfo.segments > initialInputSegmentsPerGroup, inputRowset) + def untouchedVersions = mowRowsetVersions.findAll { it != targetVersion } + def untouchedRowsets = untouchedVersions.collectEntries { int version -> + [(version): rowsetByVersion(beforeFirstCompaction, version)] + } + + set_be_param("vertical_compaction_max_segment_size", + mowMaxSegmentSizeByVersion[targetVersion].toString()) + set_be_param("cloud_single_rowset_compaction_segment_group_size", + initialInputSegmentsPerGroup.toString()) + runCumulativeCompaction(mowBackend.Host, mowBackend.HttpPort, mowTabletId) + + def afterFirstCompaction = showTablet( + "test_cloud_single_rowset_compaction_mow_delete_bitmap", + mowBackend.Host, mowBackend.HttpPort, mowTabletId) + def groupedRowset = rowsetByVersion(afterFirstCompaction, targetVersion) + def groupedInfo = parseRowsetInfo(groupedRowset) + assertEquals("NONOVERLAPPING_WITHIN_GROUP", groupedInfo.overlap) + assertNotEquals(inputRowset, groupedRowset) + untouchedRowsets.each { int version, String rowset -> + assertEquals(rowset, rowsetByVersion(afterFirstCompaction, version)) + } + checkMowQueryResult() + + set_be_param("cloud_single_rowset_compaction_segment_group_size", + groupedInfo.segments.toString()) + runCumulativeCompaction(mowBackend.Host, mowBackend.HttpPort, mowTabletId) + + def afterSecondCompaction = showTablet( + "test_cloud_single_rowset_compaction_mow_delete_bitmap", + mowBackend.Host, mowBackend.HttpPort, mowTabletId) + def nonoverlappingRowset = rowsetByVersion(afterSecondCompaction, targetVersion) + def nonoverlappingInfo = parseRowsetInfo(nonoverlappingRowset) + assertEquals("NONOVERLAPPING", nonoverlappingInfo.overlap) + assertNotEquals(groupedRowset, nonoverlappingRowset) + finalMowSegmentCounts[targetVersion] = nonoverlappingInfo.segments + untouchedRowsets.each { int version, String rowset -> + assertEquals(rowset, rowsetByVersion(afterSecondCompaction, version)) + } + checkMowQueryResult() + } + + set_be_param("enable_rowid_conversion_correctness_check", "true") + set_be_param("compaction_batch_size", "512") + mowRowsetVersions.each { int version -> compactMowRowsetTwice(version) } + assertEquals(mowRowsetVersions.size(), finalMowSegmentCounts.values().toSet().size(), + "expected different final segment counts: ${finalMowSegmentCounts}") + + // ==================================================== + logCaseSection("Multi-segment case: verify disjoint key ranges after the second " + + "cumulative compaction") + set_be_param("cloud_single_rowset_compaction_segment_group_size", + initialInputSegmentsPerGroup.toString()) + set_be_param("vertical_compaction_max_segment_size", "2048") + set_be_param("compaction_batch_size", "512") + checkSingleRowsetGroupedCompaction( + "test_cloud_single_rowset_grouped_compact_multi_seg_dup", + "DUPLICATE KEY", "v INT", "", initialInputSegmentsPerGroup, 32768, 32768 * 2, + [["100", "100"], ["100", "101"]], true) + + // ==================================================== + logCaseSection("Schema-change case: verify CREATE INDEX clears the grouped rowset " + + "layout") + sql "DROP TABLE IF EXISTS test_cloud_grouped_compaction_schema_change" + sql """ + CREATE TABLE test_cloud_grouped_compaction_schema_change ( + k INT, + v INT + ) + DUPLICATE KEY(k) + DISTRIBUTED BY HASH(k) BUCKETS 1 + PROPERTIES ( + "replication_num" = "1", + "disable_auto_compaction" = "true" + ) + """ + + StringBuilder schemaChangeContent = new StringBuilder() + for (int i = 0; i < 2; i++) { + (1..8192).each { + schemaChangeContent.append("${it},${it + i}\n") + } + } + streamLoad { + table "test_cloud_grouped_compaction_schema_change" + set "column_separator", "," + inputStream new ByteArrayInputStream(schemaChangeContent.toString().getBytes()) + time 30000 + check { result, exception, startTime, endTime -> + if (exception != null) { + throw exception + } + def json = parseJson(result) + assertEquals("success", json.Status.toLowerCase()) + assertEquals(8192 * 2, json.NumberTotalRows) + assertEquals(0, json.NumberFilteredRows) + } + } + sql "sync" + + def schemaChangeTablets = + sql_return_maparray "SHOW TABLETS FROM test_cloud_grouped_compaction_schema_change" + assertEquals(1, schemaChangeTablets.size()) + def oldTabletId = schemaChangeTablets[0].TabletId + def oldBackendId = schemaChangeTablets[0].BackendId + def schemaChangeBackends = sql_return_maparray "SHOW BACKENDS" + def oldBackend = schemaChangeBackends.find { it.BackendId == oldBackendId } + assertNotNull(oldBackend) + + set_be_param("cloud_single_rowset_compaction_segment_group_size", + initialInputSegmentsPerGroup.toString()) + runCumulativeCompaction(oldBackend.Host, oldBackend.HttpPort, oldTabletId) + + def groupedTablet = showTablet("test_cloud_grouped_compaction_schema_change", + oldBackend.Host, oldBackend.HttpPort, oldTabletId) + def groupedRowset = rowsetByVersion(groupedTablet, 2) + def groupedInfo = parseRowsetInfo(groupedRowset) + assertEquals("NONOVERLAPPING_WITHIN_GROUP", groupedInfo.overlap) + + sql """ + CREATE INDEX idx_v ON test_cloud_grouped_compaction_schema_change(v) + USING INVERTED + """ + waitForSchemaChangeDone { + sql """ + SHOW ALTER TABLE COLUMN + WHERE TableName='test_cloud_grouped_compaction_schema_change' + ORDER BY CreateTime DESC LIMIT 1 + """ + time 90 + } + + schemaChangeTablets = + sql_return_maparray "SHOW TABLETS FROM test_cloud_grouped_compaction_schema_change" + assertEquals(1, schemaChangeTablets.size()) + def newTabletId = schemaChangeTablets[0].TabletId + assertNotEquals(oldTabletId, newTabletId) + def newBackendId = schemaChangeTablets[0].BackendId + schemaChangeBackends = sql_return_maparray "SHOW BACKENDS" + def newBackend = schemaChangeBackends.find { it.BackendId == newBackendId } + assertNotNull(newBackend) + + def rewrittenTablet = showTablet("test_cloud_grouped_compaction_schema_change", + newBackend.Host, newBackend.HttpPort, newTabletId) + def rewrittenRowset = rowsetByVersion(rewrittenTablet, 2) + def rewrittenInfo = parseRowsetInfo(rewrittenRowset) + assertNotEquals("NONOVERLAPPING_WITHIN_GROUP", rewrittenInfo.overlap) + def rewrittenRowsetMeta = getRowsetMeta(newTabletId, 2) + assertTrue((rewrittenRowsetMeta.segment_group_sizes ?: []).isEmpty()) + def rewrittenCount = + sql "SELECT COUNT(*) FROM test_cloud_grouped_compaction_schema_change" + assertEquals(8192 * 2, rewrittenCount[0][0]) + def rewrittenPointRows = + readPointRows("test_cloud_grouped_compaction_schema_change") + assertEquals([["100", "100"], ["100", "101"]], rewrittenPointRows) + } finally { + sql "set enable_sql_cache=${sqlCacheOrigValue}" + GetDebugPoint().clearDebugPointsForAllBEs() + } + } +}