This is an automated email from the ASF dual-hosted git repository.
mymeiyi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new 90f75bf67cc [improvement](compaction) Support key ranges in rowset
merger (#67360)
90f75bf67cc is described below
commit 90f75bf67cca9d3cd7a77937e141f860807a8e45
Author: meiyi <[email protected]>
AuthorDate: Wed Sep 16 15:20:05 2026 +0800
[improvement](compaction) Support key ranges in rowset merger (#67360)
Later, we will support parallel compaction. Parallel compaction will
split one compaction job into multiple subtasks, with each subtask
processing only the rows within a specific key range of the input
rowsets.
This pr add optional key-range and RuntimeState parameters to the rowset
merger, propagate them through TabletReader and rowset readers. It
supports unbounded lower and upper endpoints.
related pr: https://github.com/apache/doris/pull/67333
---
be/src/storage/iterator/vertical_block_reader.cpp | 5 +
be/src/storage/merger.cpp | 85 ++++++++++--
be/src/storage/merger.h | 21 ++-
be/src/storage/rowset/beta_rowset_reader.cpp | 38 +++--
be/src/storage/rowset/rowset_reader_context.h | 7 +-
be/src/storage/segment/segment_iterator.cpp | 7 +-
be/src/storage/tablet/tablet_reader.cpp | 87 +++++++-----
be/src/storage/tablet/tablet_reader.h | 10 +-
.../compaction/vertical_compaction_test.cpp | 154 +++++++++++++++++++++
be/test/storage/tablet_reader_test.cpp | 22 +++
10 files changed, 364 insertions(+), 72 deletions(-)
diff --git a/be/src/storage/iterator/vertical_block_reader.cpp
b/be/src/storage/iterator/vertical_block_reader.cpp
index 7abc56bb145..e3d1e38de40 100644
--- a/be/src/storage/iterator/vertical_block_reader.cpp
+++ b/be/src/storage/iterator/vertical_block_reader.cpp
@@ -30,6 +30,7 @@
#include "core/column/column_vector.h"
#include "core/data_type/data_type_number.h"
#include "exprs/aggregate/aggregate_function_reader.h"
+#include "runtime/runtime_state.h"
#include "storage/compaction/compaction.h"
#include "storage/iterator/vertical_merge_iterator.h"
#include "storage/iterators.h"
@@ -112,7 +113,11 @@ Status VerticalBlockReader::_get_segment_iterators(const
ReaderParams& read_para
<< ", version:" << read_params.version;
return res;
}
+ RuntimeState* runtime_state = read_params.runtime_state;
for (const auto& rs_split : read_params.rs_splits) {
+ if (runtime_state != nullptr) {
+ RETURN_IF_CANCELLED(runtime_state);
+ }
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
diff --git a/be/src/storage/merger.cpp b/be/src/storage/merger.cpp
index 97f8e12090e..6999cb48693 100644
--- a/be/src/storage/merger.cpp
+++ b/be/src/storage/merger.cpp
@@ -38,12 +38,14 @@
#include "common/logging.h"
#include "common/status.h"
#include "core/block/block.h"
+#include "runtime/runtime_state.h"
#include "storage/iterator/block_reader.h"
#include "storage/iterator/vertical_block_reader.h"
#include "storage/iterator/vertical_merge_iterator.h"
#include "storage/iterators.h"
#include "storage/olap_common.h"
#include "storage/olap_define.h"
+#include "storage/olap_tuple.h"
#include "storage/rowid_conversion.h"
#include "storage/rowset/beta_rowset.h"
#include "storage/rowset/rowset.h"
@@ -64,11 +66,44 @@
#include "util/slice.h"
namespace doris {
+
+namespace {
+
+Status set_key_range(TabletReader::ReaderParams* reader_params, const
TabletSchema& tablet_schema,
+ const Merger::KeyRange& range) {
+ DORIS_CHECK(range.lower_key.has_value() || range.upper_key.has_value());
+ if (tablet_schema.sort_type() == SortType::ZORDER ||
+ !tablet_schema.cluster_key_uids().empty()) {
+ return Status::Error<ErrorCode::INVALID_ARGUMENT>(
+ "Key range is not supported for Z-order or cluster-key
tablets");
+ }
+ auto validate_bound = [&](const std::optional<OlapTuple>& bound) -> Status
{
+ if (bound.has_value() &&
+ (bound->size() == 0 || bound->size() >
tablet_schema.num_key_columns())) {
+ return Status::Error<ErrorCode::INVALID_ARGUMENT>(
+ "Key range bound must be a non-empty leading key prefix:
bound columns={}, "
+ "key columns={}",
+ bound->size(), tablet_schema.num_key_columns());
+ }
+ return Status::OK();
+ };
+ RETURN_IF_ERROR(validate_bound(range.lower_key));
+ RETURN_IF_ERROR(validate_bound(range.upper_key));
+ reader_params->start_key.push_back(range.lower_key);
+ reader_params->end_key.push_back(range.upper_key);
+ reader_params->start_key_include = range.lower_inclusive;
+ reader_params->end_key_include = range.upper_inclusive;
+ return Status::OK();
+}
+
+} // namespace
+
Status Merger::vmerge_rowsets(BaseTabletSPtr tablet, ReaderType reader_type,
const TabletSchema& cur_tablet_schema,
const std::vector<RowsetReaderSharedPtr>&
src_rowset_readers,
RowsetWriter* dst_rowset_writer, Statistics*
stats_output,
- std::optional<std::pair<int64_t, int64_t>>
segment_range) {
+ std::optional<std::pair<int64_t, int64_t>>
segment_range,
+ std::optional<KeyRange> key_range, RuntimeState*
runtime_state) {
if (!cur_tablet_schema.cluster_key_uids().empty()) {
return Status::InternalError(
"mow table with cluster keys does not support non vertical
compaction");
@@ -78,6 +113,10 @@ Status Merger::vmerge_rowsets(BaseTabletSPtr tablet,
ReaderType reader_type,
reader_params.tablet = tablet;
reader_params.reader_type = reader_type;
reader_params.read_row_binlog = tablet->is_row_binlog_tablet();
+ reader_params.runtime_state = runtime_state;
+ if (key_range.has_value()) {
+ RETURN_IF_ERROR(set_key_range(&reader_params, cur_tablet_schema,
*key_range));
+ }
if (reader_params.read_row_binlog) {
// Row-binlog horizontal (non-vertical) compaction must produce a
globally
// (key, TSO)-ordered output.
@@ -121,6 +160,9 @@ Status Merger::vmerge_rowsets(BaseTabletSPtr tablet,
ReaderType reader_type,
size_t output_rows = 0;
bool eof = false;
while (!eof && !ExecEnv::GetInstance()->storage_engine().stopped()) {
+ if (runtime_state != nullptr) {
+ RETURN_IF_CANCELLED(runtime_state);
+ }
auto tablet_state = tablet->tablet_state();
if (tablet_state != TABLET_RUNNING && tablet_state != TABLET_NOTREADY)
{
tablet->clear_cache();
@@ -258,7 +300,8 @@ Status Merger::vertical_compact_one_group(
RowsetWriter* dst_rowset_writer, uint32_t max_rows_per_segment,
Statistics* stats_output,
std::vector<uint32_t> key_group_cluster_key_idxes, int64_t batch_size,
CompactionSampleInfo* sample_info, VerticalCompactionContextStats*
context_stats,
- bool enable_sparse_optimization, std::optional<std::pair<int64_t,
int64_t>> segment_range) {
+ bool enable_sparse_optimization, std::optional<std::pair<int64_t,
int64_t>> segment_range,
+ std::optional<KeyRange> key_range, RuntimeState* runtime_state) {
// build tablet reader
VLOG_NOTICE << "vertical compact one group, max_rows_per_segment=" <<
max_rows_per_segment;
VerticalBlockReader reader(row_source_buf, context_stats);
@@ -268,7 +311,11 @@ Status Merger::vertical_compact_one_group(
reader_params.tablet = tablet;
reader_params.reader_type = reader_type;
reader_params.read_row_binlog = tablet->is_row_binlog_tablet();
+ reader_params.runtime_state = runtime_state;
reader_params.enable_sparse_optimization = enable_sparse_optimization;
+ if (key_range.has_value()) {
+ RETURN_IF_ERROR(set_key_range(&reader_params, tablet_schema,
*key_range));
+ }
TabletReadSource read_source;
read_source.rs_splits.reserve(src_rowset_readers.size());
@@ -310,6 +357,9 @@ Status Merger::vertical_compact_one_group(
size_t output_rows = 0;
bool eof = false;
while (!eof && !ExecEnv::GetInstance()->storage_engine().stopped()) {
+ if (runtime_state != nullptr) {
+ RETURN_IF_CANCELLED(runtime_state);
+ }
auto tablet_state = tablet->tablet_state();
if (tablet_state != TABLET_RUNNING && tablet_state != TABLET_NOTREADY)
{
tablet->clear_cache();
@@ -494,14 +544,18 @@ int64_t estimate_batch_size(int group_index,
BaseTabletSPtr tablet, int64_t way_
// 2. compact groups one by one, generate a row_source_buf when compact key
group
// and use this row_source_buf to compact value column groups
// 3. build output rowset
-Status Merger::vertical_merge_rowsets(BaseTabletSPtr tablet, ReaderType
reader_type,
- const TabletSchema& tablet_schema,
- const
std::vector<RowsetReaderSharedPtr>& src_rowset_readers,
- RowsetWriter* dst_rowset_writer,
- uint32_t max_rows_per_segment, int64_t
merge_way_num,
- Statistics* stats_output,
- VerticalCompactionProgressCallback
progress_cb,
- std::optional<std::pair<int64_t,
int64_t>> segment_range) {
+Status Merger::vertical_merge_rowsets(
+ BaseTabletSPtr tablet, ReaderType reader_type, const TabletSchema&
tablet_schema,
+ const std::vector<RowsetReaderSharedPtr>& src_rowset_readers,
+ RowsetWriter* dst_rowset_writer, uint32_t max_rows_per_segment,
int64_t merge_way_num,
+ Statistics* stats_output, VerticalCompactionProgressCallback
progress_cb,
+ std::optional<std::pair<int64_t, int64_t>> segment_range,
std::optional<KeyRange> key_range,
+ RuntimeState* runtime_state) {
+ if (key_range.has_value() && reader_type !=
ReaderType::READER_BASE_COMPACTION) {
+ return Status::Error<ErrorCode::INVALID_ARGUMENT>(
+ "Key range is only supported for base compaction,
reader_type={}",
+ static_cast<int>(reader_type));
+ }
LOG(INFO) << "Start to do vertical compaction, tablet_id: " <<
tablet->tablet_id();
VerticalCompactionContextStats context_stats;
Defer log_context_stats {[&] {
@@ -552,7 +606,7 @@ Status Merger::vertical_merge_rowsets(BaseTabletSPtr
tablet, ReaderType reader_t
// When density <= threshold, enable sparse optimization
// threshold = 0 means disable, 1 means always enable (default)
bool enable_sparse_optimization = false;
- if (!segment_range.has_value()) {
+ if (!segment_range.has_value() && !key_range.has_value()) {
for (const auto& rs_reader : src_rowset_readers) {
total_rows += rs_reader->rowset()->rowset_meta()->num_rows();
}
@@ -605,7 +659,7 @@ Status Merger::vertical_merge_rowsets(BaseTabletSPtr
tablet, ReaderType reader_t
}
}
}
- if (!segment_range.has_value() && need_footer_collection) {
+ if (!segment_range.has_value() && !key_range.has_value() &&
need_footer_collection) {
for (const auto& rs_reader : src_rowset_readers) {
auto beta_rowset =
std::dynamic_pointer_cast<BetaRowset>(rs_reader->rowset());
if (!beta_rowset) {
@@ -705,6 +759,9 @@ Status Merger::vertical_merge_rowsets(BaseTabletSPtr
tablet, ReaderType reader_t
// compact group one by one
for (auto i = 0; i < column_groups.size(); ++i) {
+ if (runtime_state != nullptr) {
+ RETURN_IF_CANCELLED(runtime_state);
+ }
VLOG_NOTICE << "row source size: " << row_sources_buf.total_size();
bool is_key = (i == 0);
int64_t batch_size = config::compaction_batch_size != -1
@@ -720,7 +777,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, segment_range);
+ enable_sparse_optimization, segment_range, key_range,
runtime_state);
{
std::unique_lock<std::mutex> lock(sample_info_lock);
sample_infos[i] = sample_info;
@@ -751,7 +808,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()) {
+ if (!segment_range.has_value() && !key_range.has_value()) {
std::unique_lock<std::mutex> 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 7cfadd49764..3b88a938925 100644
--- a/be/src/storage/merger.h
+++ b/be/src/storage/merger.h
@@ -25,6 +25,7 @@
#include "common/status.h"
#include "io/io_common.h"
#include "storage/iterators.h"
+#include "storage/olap_tuple.h"
#include "storage/rowset/rowset_fwd.h"
#include "storage/simple_rowid_conversion.h"
#include "storage/tablet/tablet_fwd.h"
@@ -33,6 +34,7 @@ namespace doris {
class KeyBoundsPB;
class RowIdConversion;
class RowsetWriter;
+class RuntimeState;
namespace segment_v2 {
class VerticalSegmentWriter;
@@ -48,6 +50,13 @@ using VerticalCompactionProgressCallback =
class Merger {
public:
+ struct KeyRange {
+ std::optional<OlapTuple> lower_key;
+ std::optional<OlapTuple> upper_key;
+ bool lower_inclusive = false;
+ bool upper_inclusive = false;
+ };
+
struct Statistics {
int64_t cloud_local_read_time = 0;
int64_t cloud_remote_read_time = 0;
@@ -70,13 +79,17 @@ public:
BaseTabletSPtr tablet, ReaderType reader_type, const TabletSchema&
cur_tablet_schema,
const std::vector<RowsetReaderSharedPtr>& src_rowset_readers,
RowsetWriter* dst_rowset_writer, Statistics* stats_output,
- std::optional<std::pair<int64_t, int64_t>> segment_range =
std::nullopt);
+ std::optional<std::pair<int64_t, int64_t>> segment_range =
std::nullopt,
+ std::optional<KeyRange> key_range = std::nullopt,
+ RuntimeState* runtime_state = nullptr);
static Status vertical_merge_rowsets(
BaseTabletSPtr tablet, ReaderType reader_type, const TabletSchema&
tablet_schema,
const std::vector<RowsetReaderSharedPtr>& 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,
- std::optional<std::pair<int64_t, int64_t>> segment_range =
std::nullopt);
+ std::optional<std::pair<int64_t, int64_t>> segment_range =
std::nullopt,
+ std::optional<KeyRange> key_range = std::nullopt,
+ RuntimeState* runtime_state = nullptr);
// for vertical compaction
static void vertical_split_columns(const TabletSchema& tablet_schema,
@@ -92,7 +105,9 @@ public:
Statistics* stats_output, std::vector<uint32_t>
key_group_cluster_key_idxes,
int64_t batch_size, CompactionSampleInfo* sample_info,
VerticalCompactionContextStats* context_stats, bool
enable_sparse_optimization = false,
- std::optional<std::pair<int64_t, int64_t>> segment_range =
std::nullopt);
+ std::optional<std::pair<int64_t, int64_t>> segment_range =
std::nullopt,
+ std::optional<KeyRange> key_range = std::nullopt,
+ RuntimeState* runtime_state = nullptr);
// for segcompaction
static Status vertical_compact_one_group(
diff --git a/be/src/storage/rowset/beta_rowset_reader.cpp
b/be/src/storage/rowset/beta_rowset_reader.cpp
index 2262c6c05a4..209ca64841a 100644
--- a/be/src/storage/rowset/beta_rowset_reader.cpp
+++ b/be/src/storage/rowset/beta_rowset_reader.cpp
@@ -115,11 +115,17 @@ Status
BetaRowsetReader::get_segment_iterators(RowsetReaderContext* read_context
_read_options.tablet_id = _rowset->rowset_meta()->tablet_id();
_read_options.read_limit = _topn_limit;
if (_read_context->lower_bound_keys != nullptr) {
+ DORIS_CHECK(_read_context->upper_bound_keys != nullptr);
+ DORIS_CHECK_EQ(_read_context->lower_bound_keys->size(),
+ _read_context->upper_bound_keys->size());
for (int i = 0; i < _read_context->lower_bound_keys->size(); ++i) {
-
_read_options.key_ranges.emplace_back(&_read_context->lower_bound_keys->at(i),
-
_read_context->is_lower_keys_included->at(i),
-
&_read_context->upper_bound_keys->at(i),
-
_read_context->is_upper_keys_included->at(i));
+ const auto& lower_bound = _read_context->lower_bound_keys->at(i);
+ const auto& upper_bound = _read_context->upper_bound_keys->at(i);
+ const auto* lower_key = lower_bound.has_value() ? &*lower_bound :
nullptr;
+ const auto* upper_key = upper_bound.has_value() ? &*upper_bound :
nullptr;
+ _read_options.key_ranges.emplace_back(
+ lower_key, _read_context->is_lower_keys_included->at(i),
upper_key,
+ _read_context->is_upper_keys_included->at(i));
}
}
@@ -212,18 +218,20 @@ Status
BetaRowsetReader::get_segment_iterators(RowsetReaderContext* read_context
_read_options.target_cast_type_for_variants =
_read_context->target_cast_type_for_variants;
if (_read_context->runtime_state != nullptr) {
_read_options.io_ctx.query_id =
&_read_context->runtime_state->query_id();
- _read_options.io_ctx.read_file_cache =
-
_read_context->runtime_state->query_options().enable_file_cache;
- _read_options.io_ctx.is_disposable =
-
_read_context->runtime_state->query_options().disable_file_cache;
- auto* query_ctx = _read_context->runtime_state->get_query_ctx();
- if (_read_context->reader_type == ReaderType::READER_QUERY &&
query_ctx != nullptr) {
- _read_options.io_ctx.remote_scan_cache_write_limiter =
- query_ctx->remote_scan_cache_write_limiter();
+ if (_read_context->reader_type == ReaderType::READER_QUERY) {
+ _read_options.io_ctx.read_file_cache =
+
_read_context->runtime_state->query_options().enable_file_cache;
+ _read_options.io_ctx.is_disposable =
+
_read_context->runtime_state->query_options().disable_file_cache;
+ if (auto* query_ctx =
_read_context->runtime_state->get_query_ctx();
+ query_ctx != nullptr) {
+ _read_options.io_ctx.remote_scan_cache_write_limiter =
+ query_ctx->remote_scan_cache_write_limiter();
+ }
+ _read_options.io_ctx.inverted_index_snii_read_no_write_file_cache =
+ _read_context->runtime_state->query_options()
+ .inverted_index_snii_read_no_write_file_cache;
}
- _read_options.io_ctx.inverted_index_snii_read_no_write_file_cache =
- _read_context->runtime_state->query_options()
- .inverted_index_snii_read_no_write_file_cache;
}
if (_read_context->condition_cache_digest) {
diff --git a/be/src/storage/rowset/rowset_reader_context.h
b/be/src/storage/rowset/rowset_reader_context.h
index 5f8e536a655..cc6da6aa2f3 100644
--- a/be/src/storage/rowset/rowset_reader_context.h
+++ b/be/src/storage/rowset/rowset_reader_context.h
@@ -18,6 +18,7 @@
#ifndef DORIS_BE_SRC_OLAP_ROWSET_ROWSET_READER_CONTEXT_H
#define DORIS_BE_SRC_OLAP_ROWSET_ROWSET_READER_CONTEXT_H
+#include <optional>
#include <set>
#include <vector>
@@ -29,12 +30,12 @@
#include "storage/index/ann/ann_topn_runtime.h"
#include "storage/olap_common.h"
#include "storage/predicate/column_predicate.h"
+#include "storage/row_cursor.h"
#include "storage/rowid_conversion.h"
#include "storage/schema.h"
namespace doris {
-class RowCursor;
class DeleteBitmap;
class DeleteHandler;
class TabletSchema;
@@ -63,9 +64,9 @@ struct RowsetReaderContext {
const std::vector<std::shared_ptr<ColumnPredicate>>* predicates = nullptr;
// value column predicate in UNIQUE table
const std::vector<std::shared_ptr<ColumnPredicate>>* value_predicates =
nullptr;
- const std::vector<RowCursor>* lower_bound_keys = nullptr;
+ const std::vector<std::optional<RowCursor>>* lower_bound_keys = nullptr;
const std::vector<bool>* is_lower_keys_included = nullptr;
- const std::vector<RowCursor>* upper_bound_keys = nullptr;
+ const std::vector<std::optional<RowCursor>>* upper_bound_keys = nullptr;
const std::vector<bool>* is_upper_keys_included = nullptr;
const DeleteHandler* delete_handler = nullptr;
OlapReaderStatistics* stats = nullptr;
diff --git a/be/src/storage/segment/segment_iterator.cpp
b/be/src/storage/segment/segment_iterator.cpp
index d53663d9d2f..99d53fd7f86 100644
--- a/be/src/storage/segment/segment_iterator.cpp
+++ b/be/src/storage/segment/segment_iterator.cpp
@@ -709,8 +709,11 @@ Status SegmentIterator::_get_row_ranges_by_keys() {
return Status::OK();
}
- // Read & seek key columns is a waste of time when no key column in _schema
- if (std::none_of(_schema->columns().begin(), _schema->columns().end(),
+ // Read & seek key columns is a waste of time when no key column in
_schema.
+ // However, parallel base compaction uses key ranges to split tasks, so
value-only vertical
+ // groups must still apply the key range to read the same physical rows as
the key group.
+ if (_opts.io_ctx.reader_type != ReaderType::READER_BASE_COMPACTION &&
+ std::none_of(_schema->columns().begin(), _schema->columns().end(),
[&](const TabletColumnPtr& col) {
return col &&
_opts.tablet_schema->column_by_uid(col->unique_id()).is_key();
diff --git a/be/src/storage/tablet/tablet_reader.cpp
b/be/src/storage/tablet/tablet_reader.cpp
index 38eae53632f..77224b957f4 100644
--- a/be/src/storage/tablet/tablet_reader.cpp
+++ b/be/src/storage/tablet/tablet_reader.cpp
@@ -96,26 +96,27 @@ Status TabletReader::_capture_rs_readers(const
ReaderParams& read_params) {
bool is_lower_key_included = _keys_param.start_key_include;
bool is_upper_key_included = _keys_param.end_key_include;
+ DORIS_CHECK_EQ(_keys_param.start_keys.size(), _keys_param.end_keys.size());
for (int i = 0; i < _keys_param.start_keys.size(); ++i) {
- // lower bound
- RowCursor& start_key = _keys_param.start_keys[i];
- RowCursor& end_key = _keys_param.end_keys[i];
-
- if (!is_lower_key_included) {
- if (compare_row_key(start_key, end_key) >= 0) {
- VLOG_NOTICE << "return EOF when lower key not include"
- << ", start_key=" << start_key.to_string()
- << ", end_key=" << end_key.to_string();
- eof = true;
- break;
- }
- } else {
- if (compare_row_key(start_key, end_key) > 0) {
- VLOG_NOTICE << "return EOF when lower key include="
- << ", start_key=" << start_key.to_string()
- << ", end_key=" << end_key.to_string();
- eof = true;
- break;
+ const auto& start_key = _keys_param.start_keys[i];
+ const auto& end_key = _keys_param.end_keys[i];
+ if (start_key.has_value() && end_key.has_value()) {
+ if (!is_lower_key_included) {
+ if (compare_row_key(*start_key, *end_key) >= 0) {
+ VLOG_NOTICE << "return EOF when lower key not include"
+ << ", start_key=" << start_key->to_string()
+ << ", end_key=" << end_key->to_string();
+ eof = true;
+ break;
+ }
+ } else {
+ if (compare_row_key(*start_key, *end_key) > 0) {
+ VLOG_NOTICE << "return EOF when lower key include="
+ << ", start_key=" << start_key->to_string()
+ << ", end_key=" << end_key->to_string();
+ eof = true;
+ break;
+ }
}
}
@@ -250,6 +251,11 @@ Status TabletReader::_init_params(const ReaderParams&
read_params) {
Status TabletReader::_init_keys_param(const ReaderParams& read_params) {
SCOPED_RAW_TIMER(&_stats.tablet_reader_init_keys_param_timer_ns);
+ if (read_params.start_key.size() != read_params.end_key.size()) {
+ return Status::Error<INVALID_ARGUMENT>(
+ "The number of start keys does not equal the number of end
keys: {} vs {}",
+ read_params.start_key.size(), read_params.end_key.size());
+ }
if (read_params.start_key.empty()) {
return Status::OK();
}
@@ -258,10 +264,16 @@ Status TabletReader::_init_keys_param(const ReaderParams&
read_params) {
_keys_param.end_key_include = read_params.end_key_include;
size_t start_key_size = read_params.start_key.size();
- //_keys_param.start_keys.resize(start_key_size);
- std::vector<RowCursor>(start_key_size).swap(_keys_param.start_keys);
-
- size_t scan_key_size = read_params.start_key.front().size();
+
std::vector<std::optional<RowCursor>>(start_key_size).swap(_keys_param.start_keys);
+
+ const OlapTuple* first_key =
+ read_params.start_key.front().has_value() ?
&*read_params.start_key.front()
+ : read_params.end_key.front().has_value() ?
&*read_params.end_key.front()
+ : nullptr;
+ if (first_key == nullptr) {
+ return Status::Error<INVALID_ARGUMENT>("A key range must have at least
one bound");
+ }
+ size_t scan_key_size = first_key->size();
if (scan_key_size > _tablet_schema->num_columns()) {
return Status::Error<INVALID_ARGUMENT>(
"Input param are invalid. Column count is bigger than
num_columns of schema. "
@@ -270,13 +282,17 @@ Status TabletReader::_init_keys_param(const ReaderParams&
read_params) {
}
for (size_t i = 0; i < start_key_size; ++i) {
- if (read_params.start_key[i].size() != scan_key_size) {
+ if (!read_params.start_key[i].has_value()) {
+ continue;
+ }
+ if (read_params.start_key[i]->size() != scan_key_size) {
return Status::Error<INVALID_ARGUMENT>(
"The start_key.at({}).size={}, not equals the
scan_key_size={}", i,
- read_params.start_key[i].size(), scan_key_size);
+ read_params.start_key[i]->size(), scan_key_size);
}
- Status res = _keys_param.start_keys[i].init(_tablet_schema,
read_params.start_key[i]);
+ auto& start_key = _keys_param.start_keys[i].emplace();
+ Status res = start_key.init(_tablet_schema, *read_params.start_key[i]);
if (!res.ok()) {
LOG(WARNING) << "fail to init row cursor. res = " << res;
return res;
@@ -284,22 +300,31 @@ Status TabletReader::_init_keys_param(const ReaderParams&
read_params) {
}
size_t end_key_size = read_params.end_key.size();
- //_keys_param.end_keys.resize(end_key_size);
- std::vector<RowCursor>(end_key_size).swap(_keys_param.end_keys);
+
std::vector<std::optional<RowCursor>>(end_key_size).swap(_keys_param.end_keys);
for (size_t i = 0; i < end_key_size; ++i) {
- if (read_params.end_key[i].size() != scan_key_size) {
+ if (!read_params.end_key[i].has_value()) {
+ continue;
+ }
+ if (read_params.end_key[i]->size() != scan_key_size) {
return Status::Error<INVALID_ARGUMENT>(
"The end_key.at({}).size={}, not equals the
scan_key_size={}", i,
- read_params.end_key[i].size(), scan_key_size);
+ read_params.end_key[i]->size(), scan_key_size);
}
- Status res = _keys_param.end_keys[i].init(_tablet_schema,
read_params.end_key[i]);
+ auto& end_key = _keys_param.end_keys[i].emplace();
+ Status res = end_key.init(_tablet_schema, *read_params.end_key[i]);
if (!res.ok()) {
LOG(WARNING) << "fail to init row cursor. res = " << res;
return res;
}
}
+ for (size_t i = 0; i < start_key_size; ++i) {
+ if (!_keys_param.start_keys[i].has_value() &&
!_keys_param.end_keys[i].has_value()) {
+ return Status::Error<INVALID_ARGUMENT>("Key range {} has no
bounds", i);
+ }
+ }
+
//TODO:check the valid of start_key and end_key.(eg. start_key <= end_key)
return Status::OK();
diff --git a/be/src/storage/tablet/tablet_reader.h
b/be/src/storage/tablet/tablet_reader.h
index 03a15982d54..db8e557ef54 100644
--- a/be/src/storage/tablet/tablet_reader.h
+++ b/be/src/storage/tablet/tablet_reader.h
@@ -24,6 +24,7 @@
#include <stdint.h>
#include <memory>
+#include <optional>
#include <set>
#include <string>
#include <unordered_set>
@@ -87,8 +88,8 @@ inline int compare_row_key(const RowCursor& lhs, const
RowCursor& rhs) {
class TabletReader {
struct KeysParam {
- std::vector<RowCursor> start_keys;
- std::vector<RowCursor> end_keys;
+ std::vector<std::optional<RowCursor>> start_keys;
+ std::vector<std::optional<RowCursor>> end_keys;
bool start_key_include = false;
bool end_key_include = false;
};
@@ -135,8 +136,9 @@ public:
bool use_page_cache = false;
Version version = Version(-1, 0);
- std::vector<OlapTuple> start_key;
- std::vector<OlapTuple> end_key;
+ // The vectors are range-aligned; nullopt represents an unbounded
endpoint.
+ std::vector<std::optional<OlapTuple>> start_key;
+ std::vector<std::optional<OlapTuple>> end_key;
bool start_key_include = false;
bool end_key_include = false;
diff --git a/be/test/storage/compaction/vertical_compaction_test.cpp
b/be/test/storage/compaction/vertical_compaction_test.cpp
index b5627662170..108ded0a809 100644
--- a/be/test/storage/compaction/vertical_compaction_test.cpp
+++ b/be/test/storage/compaction/vertical_compaction_test.cpp
@@ -31,6 +31,7 @@
#include <iostream>
#include <memory>
+#include <optional>
#include <string>
#include <tuple>
#include <unordered_map>
@@ -45,20 +46,24 @@
#include "core/column/column.h"
#include "core/column/column_nullable.h"
#include "core/data_type/data_type.h"
+#include "core/field.h"
#include "gtest/gtest_pred_impl.h"
#include "io/cache/block_file_cache_factory.h"
#include "io/fs/local_file_system.h"
#include "io/io_common.h"
#include "json2pb/json_to_pb.h"
#include "runtime/exec_env.h"
+#include "runtime/runtime_state.h"
#include "runtime/thread_context.h"
#include "storage/delete/delete_handler.h"
#include "storage/iterator/vertical_merge_iterator.h"
#include "storage/merger.h"
#include "storage/olap_common.h"
+#include "storage/olap_tuple.h"
#include "storage/options.h"
#include "storage/rowid_conversion.h"
#include "storage/rowset/beta_rowset.h"
+#include "storage/rowset/beta_rowset_reader.h"
#include "storage/rowset/rowset.h"
#include "storage/rowset/rowset_factory.h"
#include "storage/rowset/rowset_meta.h"
@@ -681,6 +686,155 @@ TEST_F(VerticalCompactionTest, TestDupKeyVerticalMerge) {
}
}
+TEST_F(VerticalCompactionTest, MergeHonorsRuntimeStateCancellation) {
+ constexpr int num_segments = 1;
+ std::vector<std::vector<std::vector<std::tuple<int64_t, int64_t>>>>
input_data;
+ generate_input_data(1, num_segments, 10, NONOVERLAPPING, input_data);
+ auto tablet_schema = create_schema();
+ auto input_rowset = create_rowset(tablet_schema, NONOVERLAPPING,
input_data.front(), 0);
+ auto tablet = create_tablet(*tablet_schema, false);
+
+ for (const bool is_vertical : {false, true}) {
+ RowsetReaderSharedPtr input_reader;
+ ASSERT_TRUE(input_rowset->create_reader(&input_reader).ok());
+ std::vector<RowsetReaderSharedPtr> input_readers =
{std::move(input_reader)};
+ auto writer_context =
+ create_rowset_writer_context(tablet_schema, NONOVERLAPPING,
3456, {0, 0});
+ auto writer_result =
+ RowsetFactory::create_rowset_writer(*engine_ref,
writer_context, is_vertical);
+ ASSERT_TRUE(writer_result.has_value()) << writer_result.error();
+ auto output_writer = std::move(writer_result).value();
+
+ RuntimeState runtime_state;
+ runtime_state.cancel(Status::Cancelled("injected compaction
cancellation"));
+ Merger::Statistics stats;
+ Status status;
+ if (is_vertical) {
+ status = Merger::vertical_merge_rowsets(
+ tablet, ReaderType::READER_BASE_COMPACTION,
*tablet_schema, input_readers,
+ output_writer.get(), 100, num_segments, &stats, nullptr,
std::nullopt,
+ std::nullopt, &runtime_state);
+ } else {
+ status = Merger::vmerge_rowsets(tablet,
ReaderType::READER_BASE_COMPACTION,
+ *tablet_schema, input_readers,
output_writer.get(),
+ &stats, std::nullopt,
std::nullopt, &runtime_state);
+ }
+ EXPECT_TRUE(status.is<ErrorCode::CANCELLED>()) << status;
+ }
+}
+
+TEST_F(VerticalCompactionTest, MergeHonorsKeyRanges) {
+ constexpr int num_rows = 10;
+ std::vector<std::vector<std::vector<std::tuple<int64_t, int64_t>>>>
input_data;
+ generate_input_data(1, 1, num_rows, NONOVERLAPPING, input_data);
+ auto tablet_schema = create_schema();
+ auto input_rowset = create_rowset(tablet_schema, NONOVERLAPPING,
input_data.front(), 0);
+ auto tablet = create_tablet(*tablet_schema, false);
+
+ auto make_key = [](int32_t value) {
+ OlapTuple key;
+ key.add_field(Field::create_field<TYPE_INT>(value));
+ return key;
+ };
+ auto run_case = [&](bool is_vertical, Merger::KeyRange key_range, int32_t
expected_begin,
+ int32_t expected_end) {
+ RowsetReaderSharedPtr input_reader;
+ ASSERT_TRUE(input_rowset->create_reader(&input_reader).ok());
+ std::vector<RowsetReaderSharedPtr> input_readers =
{std::move(input_reader)};
+ auto writer_context =
+ create_rowset_writer_context(tablet_schema, NONOVERLAPPING,
100, {0, 0});
+ auto writer_result =
+ RowsetFactory::create_rowset_writer(*engine_ref,
writer_context, is_vertical);
+ ASSERT_TRUE(writer_result.has_value()) << writer_result.error();
+ auto output_writer = std::move(writer_result).value();
+
+ Merger::Statistics stats;
+ RuntimeState runtime_state;
+ TQueryOptions query_options;
+ query_options.__set_enable_file_cache(false);
+ query_options.__set_disable_file_cache(false);
+ runtime_state.set_query_options(query_options);
+ Status status;
+ if (is_vertical) {
+ status = Merger::vertical_merge_rowsets(
+ tablet, ReaderType::READER_BASE_COMPACTION,
*tablet_schema, input_readers,
+ output_writer.get(), 100, 1, &stats, nullptr, std::nullopt,
+ std::move(key_range), &runtime_state);
+ } else {
+ status = Merger::vmerge_rowsets(tablet,
ReaderType::READER_BASE_COMPACTION,
+ *tablet_schema, input_readers,
output_writer.get(),
+ &stats, std::nullopt,
std::move(key_range),
+ &runtime_state);
+ }
+ ASSERT_TRUE(status.ok()) << status;
+ const auto& read_options =
+ assert_cast<const
BetaRowsetReader&>(*input_readers.front())._read_options;
+ EXPECT_TRUE(read_options.io_ctx.read_file_cache);
+ EXPECT_TRUE(read_options.io_ctx.is_disposable);
+
+ RowsetSharedPtr output_rowset;
+ ASSERT_EQ(Status::OK(), output_writer->build(output_rowset));
+ ASSERT_TRUE(output_rowset);
+ ASSERT_EQ(expected_end - expected_begin, output_rowset->num_rows());
+
+ RowsetReaderContext reader_context;
+ reader_context.tablet_schema = tablet_schema;
+ reader_context.need_ordered_result = false;
+ auto read_schema = std::make_shared<ReadSchema>(
+ project_columns_by_ordinal(tablet_schema->columns(),
std::vector<ColumnId> {0, 1}));
+ reader_context.read_schema = read_schema;
+ RowsetReaderSharedPtr output_reader;
+ create_and_init_rowset_reader(output_rowset.get(), reader_context,
&output_reader);
+
+ int32_t expected_key = expected_begin;
+ do {
+ auto block = read_schema->create_read_block();
+ status = output_reader->next_batch(&block);
+ const auto& columns = block.get_columns_with_type_and_name();
+ ASSERT_EQ(2, columns.size());
+ for (size_t row = 0; row < block.rows(); ++row) {
+ EXPECT_EQ(expected_key, columns[0].column->get_int(row));
+ EXPECT_EQ(expected_key + 1, columns[1].column->get_int(row));
+ ++expected_key;
+ }
+ } while (status.ok());
+ EXPECT_TRUE(status.is<END_OF_FILE>()) << status;
+ EXPECT_EQ(expected_end, expected_key);
+ };
+
+ for (bool is_vertical : {false, true}) {
+ run_case(is_vertical, {make_key(2), make_key(7), true, false}, 2, 7);
+ run_case(is_vertical, {std::nullopt, make_key(3), false, true}, 0, 4);
+ run_case(is_vertical, {make_key(6), std::nullopt, false, false}, 7,
num_rows);
+ }
+
+ RowsetReaderSharedPtr input_reader;
+ ASSERT_TRUE(input_rowset->create_reader(&input_reader).ok());
+ std::vector<RowsetReaderSharedPtr> input_readers =
{std::move(input_reader)};
+ auto writer_context = create_rowset_writer_context(tablet_schema,
NONOVERLAPPING, 100, {0, 0});
+ auto writer_result = RowsetFactory::create_rowset_writer(*engine_ref,
writer_context, false);
+ ASSERT_TRUE(writer_result.has_value()) << writer_result.error();
+ auto output_writer = std::move(writer_result).value();
+ Merger::Statistics stats;
+ auto expect_invalid = [&](const TabletSchema& schema, Merger::KeyRange
key_range) {
+ auto status = Merger::vmerge_rowsets(tablet,
ReaderType::READER_BASE_COMPACTION, schema,
+ input_readers,
output_writer.get(), &stats,
+ std::nullopt,
std::move(key_range));
+ EXPECT_TRUE(status.is<ErrorCode::INVALID_ARGUMENT>()) << status;
+ };
+
+ auto wide_key = make_key(2);
+ wide_key.add_field(Field::create_field<TYPE_INT>(3));
+ expect_invalid(*tablet_schema, {std::move(wide_key), std::nullopt, true,
false});
+
+ TabletSchemaPB zorder_schema_pb;
+ tablet_schema->to_schema_pb(&zorder_schema_pb);
+ zorder_schema_pb.set_sort_type(SortType::ZORDER);
+ TabletSchema zorder_schema;
+ zorder_schema.init_from_pb(zorder_schema_pb);
+ expect_invalid(zorder_schema, {make_key(2), std::nullopt, true, false});
+}
+
TEST_F(VerticalCompactionTest, TestDupWithoutKeyVerticalMerge) {
auto num_input_rowset = 2;
auto num_segments = 2;
diff --git a/be/test/storage/tablet_reader_test.cpp
b/be/test/storage/tablet_reader_test.cpp
index 0978308c458..5d6ffeea57d 100644
--- a/be/test/storage/tablet_reader_test.cpp
+++ b/be/test/storage/tablet_reader_test.cpp
@@ -135,4 +135,26 @@ TEST_F(TabletReaderTest,
remove_delete_columns_keeps_unrelated_paths) {
EXPECT_EQ(size_t(2), access_paths.size());
}
+
+TEST_F(TabletReaderTest, initializes_unbounded_key_ranges) {
+ auto schema = create_schema({{"k1", 10}});
+ OlapTuple upper_key;
+ upper_key.add_field(Field::create_field<TYPE_INT>(10));
+ OlapTuple lower_key;
+ lower_key.add_field(Field::create_field<TYPE_INT>(20));
+
+ TabletReader::ReaderParams params;
+ params.start_key.emplace_back(std::nullopt);
+ params.end_key.emplace_back(std::move(upper_key));
+ params.start_key.emplace_back(std::move(lower_key));
+ params.end_key.emplace_back(std::nullopt);
+
+ TabletReader reader;
+ reader._tablet_schema = std::move(schema);
+ ASSERT_TRUE(reader._init_keys_param(params).ok());
+ EXPECT_FALSE(reader._keys_param.start_keys[0].has_value());
+ EXPECT_TRUE(reader._keys_param.end_keys[0].has_value());
+ EXPECT_TRUE(reader._keys_param.start_keys[1].has_value());
+ EXPECT_FALSE(reader._keys_param.end_keys[1].has_value());
+}
} // namespace doris
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]