This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new 01d6797c2fd branch-4.1: [fix](compaction) Release exhausted segment
contexts during vertical compaction #65963 (#67599)
01d6797c2fd is described below
commit 01d6797c2fd199c43d1c50219cdb24007f0c3d07
Author: Jamie <[email protected]>
AuthorDate: Thu Sep 10 09:24:14 2026 +0800
branch-4.1: [fix](compaction) Release exhausted segment contexts during
vertical compaction #65963 (#67599)
### What problem does this PR solve?
Issue Number: None
Related PR: #65963
Problem Summary: Backport #65963 to `branch-4.1`.
Vertical compaction kept exhausted segment iterator contexts and
dictionary/block memory alive until all value-column groups finished.
This change tracks final row-source consumption and releases exhausted
contexts immediately, while adding active/peak context diagnostics and
regression coverage.
The backport adapts the storage/iterator versus olap/vec source layout
and preserves the target branch APIs. Branches without
`use_insert_order_when_same` retain their existing ordering behavior.
### Release note
Release exhausted segment contexts during vertical compaction to bound
retained memory.
### Check List (For Author)
- Test: Static validation
- Clang Format 16 `--dry-run --Werror`: passed for all changed C++ files
- `git diff --check`: passed
- Changed-file scope, latest base, and cherry-pick provenance verified
- Focused local compilation was attempted on branch-4.1, but the
available generated Thrift/Protobuf headers did not match this branch;
hosted CI is pending
- Behavior changed: Yes
- Releases a segment context after its final row-source reference is
consumed
- Adds active and peak segment-context diagnostics
- Does this need documentation: No
---
be/src/storage/iterator/vertical_block_reader.cpp | 9 +-
be/src/storage/iterator/vertical_block_reader.h | 7 +-
.../storage/iterator/vertical_merge_iterator.cpp | 109 ++++++--
be/src/storage/iterator/vertical_merge_iterator.h | 53 +++-
be/src/storage/merger.cpp | 16 +-
be/src/storage/merger.h | 3 +-
.../compaction/vertical_compaction_test.cpp | 285 +++++++++++++++++++++
7 files changed, 444 insertions(+), 38 deletions(-)
diff --git a/be/src/storage/iterator/vertical_block_reader.cpp
b/be/src/storage/iterator/vertical_block_reader.cpp
index dbc4d7f55ba..1fe1274dac5 100644
--- a/be/src/storage/iterator/vertical_block_reader.cpp
+++ b/be/src/storage/iterator/vertical_block_reader.cpp
@@ -150,16 +150,17 @@ Status VerticalBlockReader::_init_collect_iter(const
ReaderParams& read_params,
_vcollect_iter = new_vertical_fifo_merge_iterator(
std::move(*segment_iters_ptr), iterator_init_flag,
rowset_ids,
ori_return_col_size, read_params.tablet->keys_type(),
seq_col_idx,
- _row_sources_buffer);
+ _row_sources_buffer, _context_stats);
} else {
_vcollect_iter = new_vertical_heap_merge_iterator(
std::move(*segment_iters_ptr), iterator_init_flag,
rowset_ids,
ori_return_col_size, read_params.tablet->keys_type(),
seq_col_idx,
- _row_sources_buffer,
read_params.key_group_cluster_key_idxes);
+ _row_sources_buffer, _context_stats,
read_params.key_group_cluster_key_idxes);
}
} else {
- _vcollect_iter =
new_vertical_mask_merge_iterator(std::move(*segment_iters_ptr),
- ori_return_col_size,
_row_sources_buffer);
+ _vcollect_iter =
+
new_vertical_mask_merge_iterator(std::move(*segment_iters_ptr),
ori_return_col_size,
+ _row_sources_buffer,
_context_stats);
}
// init collect iterator
StorageReadOptions opts;
diff --git a/be/src/storage/iterator/vertical_block_reader.h
b/be/src/storage/iterator/vertical_block_reader.h
index 80e36434f08..baba2a215f8 100644
--- a/be/src/storage/iterator/vertical_block_reader.h
+++ b/be/src/storage/iterator/vertical_block_reader.h
@@ -44,11 +44,13 @@ struct RowsetId;
class RowSourcesBuffer;
struct RowBatch;
+struct VerticalCompactionContextStats;
class VerticalBlockReader final : public TabletReader {
public:
- VerticalBlockReader(RowSourcesBuffer* row_sources_buffer)
- : _row_sources_buffer(row_sources_buffer) {
+ VerticalBlockReader(RowSourcesBuffer* row_sources_buffer,
+ VerticalCompactionContextStats* context_stats =
nullptr)
+ : _row_sources_buffer(row_sources_buffer),
_context_stats(context_stats) {
_id = nextId++;
}
@@ -113,6 +115,7 @@ private:
Status (VerticalBlockReader::*_next_block_func)(Block* block, bool* eof) =
nullptr;
RowSourcesBuffer* _row_sources_buffer;
+ VerticalCompactionContextStats* _context_stats;
ColumnPtr _delete_filter_column;
// for agg mode
diff --git a/be/src/storage/iterator/vertical_merge_iterator.cpp
b/be/src/storage/iterator/vertical_merge_iterator.cpp
index 1df2b6c4bc4..856bd78226e 100644
--- a/be/src/storage/iterator/vertical_merge_iterator.cpp
+++ b/be/src/storage/iterator/vertical_merge_iterator.cpp
@@ -17,6 +17,7 @@
#include "storage/iterator/vertical_merge_iterator.h"
+#include <bvar/bvar.h>
#include <fcntl.h>
#include <gen_cpp/olap_file.pb.h>
#include <stdlib.h>
@@ -41,6 +42,13 @@ namespace doris {
#include "common/compile_check_begin.h"
using namespace ErrorCode;
+namespace {
+
+bvar::Adder<int64_t> g_vertical_compaction_active_segment_contexts(
+ "vertical_compaction_active_segment_contexts");
+
+} // namespace
+
// -------------- row source ---------------//
RowSource::RowSource(uint16_t source_num, bool agg_flag) {
_data = (source_num & SOURCE_FLAG) | (source_num & AGG_FLAG);
@@ -85,8 +93,14 @@ Status RowSourcesBuffer::append(const
std::vector<RowSource>& row_sources) {
_reset_buffer();
}
}
+ uint64_t source_position = _total_size;
for (const auto& source : row_sources) {
_buffer.push_back(source.data());
+ auto source_num = source.get_source_num();
+ if (source_num >= _last_source_positions.size()) {
+ _last_source_positions.resize(source_num + 1);
+ }
+ _last_source_positions[source_num] = source_position++;
}
_total_size += row_sources.size();
return Status::OK();
@@ -94,6 +108,7 @@ Status RowSourcesBuffer::append(const
std::vector<RowSource>& row_sources) {
Status RowSourcesBuffer::seek_to_begin() {
_buf_idx = 0;
+ _read_index = 0;
if (_fd > 0) {
auto offset = lseek(_fd, 0, SEEK_SET);
if (offset != 0) {
@@ -274,6 +289,45 @@ Status RowSourcesBuffer::_deserialize() {
}
// ---------- vertical merge iterator context ----------//
+VerticalMergeIteratorContext::~VerticalMergeIteratorContext() {
+ release_resources();
+}
+
+void VerticalMergeIteratorContext::release_resources() {
+ _valid = false;
+ _iter.reset();
+ _block.reset();
+ // Unlike the physical EOF fallback, source exhaustion proves that this
context will not
+ // access these blocks again. External IteratorRowRef/RowBatch owners keep
them alive through
+ // shared_ptr, so dropping the context-owned references here is
intentional.
+ _block_list.clear();
+ _mark_inactive();
+}
+
+void VerticalMergeIteratorContext::_mark_active() {
+ DCHECK(!_is_active_context_counted);
+ _is_active_context_counted = true;
+ g_vertical_compaction_active_segment_contexts << 1;
+ if (_context_stats != nullptr) {
+ ++_context_stats->active_segment_contexts;
+ _context_stats->active_segment_contexts_peak =
+ std::max(_context_stats->active_segment_contexts_peak,
+ _context_stats->active_segment_contexts);
+ }
+}
+
+void VerticalMergeIteratorContext::_mark_inactive() {
+ if (!_is_active_context_counted) {
+ return;
+ }
+ _is_active_context_counted = false;
+ g_vertical_compaction_active_segment_contexts << -1;
+ if (_context_stats != nullptr) {
+ DCHECK_GT(_context_stats->active_segment_contexts, 0);
+ --_context_stats->active_segment_contexts;
+ }
+}
+
Status VerticalMergeIteratorContext::block_reset(const std::shared_ptr<Block>&
block) {
if (!block->columns()) {
const Schema& schema = _iter->schema();
@@ -376,6 +430,7 @@ Status VerticalMergeIteratorContext::init(const
StorageReadOptions& opts,
sample_info->rows += rows();
}
if (valid()) {
+ _mark_active();
RETURN_IF_ERROR(advance());
}
_inited = true;
@@ -460,6 +515,9 @@ Status VerticalMergeIteratorContext::_load_next_block() {
// the column iterator in the segment iterator will hold the
dictionary.
// Release the segment iterator to free up the dictionary.
_iter.reset();
+ if (_block_list.empty()) {
+ _mark_inactive();
+ }
return Status::OK();
} else {
return st;
@@ -582,7 +640,7 @@ Status VerticalHeapMergeIterator::init(const
StorageReadOptions& opts,
auto& iter = _origin_iters[seg_order];
auto ctx = std::make_unique<VerticalMergeIteratorContext>(
std::move(iter), _rowset_ids[seg_order], _ori_return_cols,
seg_order, _seq_col_idx,
- _key_group_cluster_key_idxes);
+ _context_stats, _key_group_cluster_key_idxes);
_ori_iter_ctx.push_back(std::move(ctx));
}
_origin_iters.clear();
@@ -648,9 +706,9 @@ Status VerticalFifoMergeIterator::next_batch(Block* block) {
cur_order++) {
auto& next_iter = _origin_iters[cur_order];
std::unique_ptr<VerticalMergeIteratorContext> next_ctx(
- new VerticalMergeIteratorContext(std::move(next_iter),
-
_rowset_ids[cur_order], _ori_return_cols,
- cur_order,
_seq_col_idx));
+ new VerticalMergeIteratorContext(
+ std::move(next_iter), _rowset_ids[cur_order],
_ori_return_cols,
+ cur_order, _seq_col_idx, _context_stats,
std::vector<uint32_t> {}));
RETURN_IF_ERROR(next_ctx->init(_opts));
if (next_ctx->valid()) {
_cur_iter_ctx.swap(next_ctx);
@@ -688,9 +746,9 @@ Status VerticalFifoMergeIterator::init(const
StorageReadOptions& opts,
// will not be pushed into heap, we should init next one util we find a
valid iter
// so this rowset can work in heap
for (auto& iter : _origin_iters) {
- std::unique_ptr<VerticalMergeIteratorContext> ctx(
- new VerticalMergeIteratorContext(std::move(iter),
_rowset_ids[seg_order],
- _ori_return_cols, seg_order,
_seq_col_idx));
+ std::unique_ptr<VerticalMergeIteratorContext> ctx(new
VerticalMergeIteratorContext(
+ std::move(iter), _rowset_ids[seg_order], _ori_return_cols,
seg_order, _seq_col_idx,
+ _context_stats));
RETURN_IF_ERROR(ctx->init(opts, sample_info));
if (!ctx->valid()) {
++seg_order;
@@ -718,6 +776,14 @@ Status
VerticalMaskMergeIterator::check_all_iter_finished() {
}
return Status::OK();
}
+
+void VerticalMaskMergeIterator::consume_row_sources(uint16_t order, size_t
count) {
+ _row_sources_buf->advance(count);
+ if (_row_sources_buf->is_source_exhausted(order)) {
+ _origin_iter_ctx[order]->release_resources();
+ }
+}
+
Status VerticalMaskMergeIterator::next_row(IteratorRowRef* ref) {
DCHECK(_row_sources_buf);
auto st = _row_sources_buf->has_remaining();
@@ -745,7 +811,7 @@ Status VerticalMaskMergeIterator::next_row(IteratorRowRef*
ref) {
}
ctx->set_is_first_row(false);
- _row_sources_buf->advance();
+ consume_row_sources(order);
return Status::OK();
}
RETURN_IF_ERROR(ctx->advance());
@@ -755,7 +821,7 @@ Status VerticalMaskMergeIterator::next_row(IteratorRowRef*
ref) {
_filtered_rows++;
}
- _row_sources_buf->advance();
+ consume_row_sources(order);
return Status::OK();
}
@@ -778,15 +844,16 @@ Status
VerticalMaskMergeIterator::unique_key_next_row(IteratorRowRef* ref) {
// Except first row, we call advance first and than get cur row
ctx->set_cur_row_ref(ref);
ctx->set_is_first_row(false);
- _row_sources_buf->advance();
+ consume_row_sources(order);
return Status::OK();
}
RETURN_IF_ERROR(ctx->advance());
- _row_sources_buf->advance();
if (!row_source.agg_flag()) {
ctx->set_cur_row_ref(ref);
+ consume_row_sources(order);
return Status::OK();
}
+ consume_row_sources(order);
_filtered_rows++;
st = _row_sources_buf->has_remaining();
}
@@ -843,7 +910,7 @@ Status
VerticalMaskMergeIterator::unique_key_next_batch(std::vector<RowBatch>* b
// If current row has agg_flag=true, skip it (single row)
if (row_source.agg_flag()) {
- _row_sources_buf->advance();
+ consume_row_sources(order);
_filtered_rows++;
continue;
}
@@ -864,7 +931,7 @@ Status
VerticalMaskMergeIterator::unique_key_next_batch(std::vector<RowBatch>* b
RETURN_IF_ERROR(ctx->advance_by(run_count - 1));
}
- _row_sources_buf->advance(run_count);
+ consume_row_sources(order, run_count);
// Try to merge into current batch or create new batch
if (current_block == block && start_row == batch_start + batch_count) {
@@ -941,8 +1008,8 @@ Status VerticalMaskMergeIterator::init(const
StorageReadOptions& opts,
RowsetId rs_id;
for (auto& iter : _origin_iters) {
- auto ctx =
std::make_unique<VerticalMergeIteratorContext>(std::move(iter), rs_id,
-
_ori_return_cols, -1, -1);
+ auto ctx = std::make_unique<VerticalMergeIteratorContext>(
+ std::move(iter), rs_id, _ori_return_cols, -1, -1,
_context_stats);
_origin_iter_ctx.push_back(std::move(ctx));
}
_origin_iters.clear();
@@ -957,26 +1024,28 @@ std::shared_ptr<RowwiseIterator>
new_vertical_heap_merge_iterator(
std::vector<RowwiseIteratorUPtr>&& inputs, const std::vector<bool>&
iterator_init_flag,
const std::vector<RowsetId>& rowset_ids, size_t ori_return_cols,
KeysType keys_type,
uint32_t seq_col_idx, RowSourcesBuffer* row_sources,
+ VerticalCompactionContextStats* context_stats,
std::vector<uint32_t> key_group_cluster_key_idxes) {
return std::make_shared<VerticalHeapMergeIterator>(
std::move(inputs), iterator_init_flag, rowset_ids,
ori_return_cols, keys_type,
- seq_col_idx, row_sources, key_group_cluster_key_idxes);
+ seq_col_idx, row_sources, context_stats,
key_group_cluster_key_idxes);
}
std::shared_ptr<RowwiseIterator> new_vertical_fifo_merge_iterator(
std::vector<RowwiseIteratorUPtr>&& inputs, const std::vector<bool>&
iterator_init_flag,
const std::vector<RowsetId>& rowset_ids, size_t ori_return_cols,
KeysType keys_type,
- uint32_t seq_col_idx, RowSourcesBuffer* row_sources) {
+ uint32_t seq_col_idx, RowSourcesBuffer* row_sources,
+ VerticalCompactionContextStats* context_stats) {
return std::make_shared<VerticalFifoMergeIterator>(std::move(inputs),
iterator_init_flag,
rowset_ids,
ori_return_cols, keys_type,
- seq_col_idx,
row_sources);
+ seq_col_idx,
row_sources, context_stats);
}
std::shared_ptr<RowwiseIterator> new_vertical_mask_merge_iterator(
std::vector<RowwiseIteratorUPtr>&& inputs, size_t ori_return_cols,
- RowSourcesBuffer* row_sources) {
+ RowSourcesBuffer* row_sources, VerticalCompactionContextStats*
context_stats) {
return std::make_shared<VerticalMaskMergeIterator>(std::move(inputs),
ori_return_cols,
- row_sources);
+ row_sources,
context_stats);
}
#include "common/compile_check_end.h"
diff --git a/be/src/storage/iterator/vertical_merge_iterator.h
b/be/src/storage/iterator/vertical_merge_iterator.h
index 865399c4747..dbb01c63678 100644
--- a/be/src/storage/iterator/vertical_merge_iterator.h
+++ b/be/src/storage/iterator/vertical_merge_iterator.h
@@ -42,6 +42,12 @@ namespace doris {
#include "common/compile_check_begin.h"
enum KeysType : int;
+// Owned by one vertical_merge_rowsets invocation and updated synchronously by
its iterators.
+struct VerticalCompactionContextStats {
+ int64_t active_segment_contexts = 0;
+ int64_t active_segment_contexts_peak = 0;
+};
+
// Row source represent row location in multi-segments
// use a uint16_t to store info
// the lower 15 bits means segment_id in segment pool, and the higher 1 bits
means agg flag.
@@ -102,11 +108,17 @@ public:
void advance(int64_t step = 1) {
DCHECK(_buf_idx + step <= _buffer.size());
_buf_idx += step;
+ _read_index += step;
}
uint64_t buf_idx() const { return _buf_idx; }
uint64_t total_size() const { return _total_size; }
uint64_t buffered_size() { return _buffer.size(); }
+ bool is_source_exhausted(uint16_t source) const {
+ DCHECK(source < _last_source_positions.size());
+ return source < _last_source_positions.size() &&
+ _read_index > _last_source_positions[source];
+ }
void set_agg_flag(uint64_t index, bool agg);
bool get_agg_flag(uint64_t index);
@@ -142,6 +154,8 @@ private:
int _fd = -1;
PaddedPODArray<UInt16> _buffer;
uint64_t _total_size = 0;
+ uint64_t _read_index = 0;
+ std::vector<uint64_t> _last_source_positions;
};
// --------------- VerticalMergeIteratorContext ------------- //
@@ -150,6 +164,7 @@ class VerticalMergeIteratorContext {
public:
VerticalMergeIteratorContext(RowwiseIteratorUPtr&& iter, RowsetId
rowset_id,
size_t ori_return_cols, uint32_t order,
uint32_t seq_col_idx,
+ VerticalCompactionContextStats* context_stats,
std::vector<uint32_t>
key_group_cluster_key_idxes = {})
: _iter(std::move(iter)),
_rowset_id(rowset_id),
@@ -157,14 +172,15 @@ public:
_order(order),
_seq_col_idx(seq_col_idx),
_num_key_columns(_iter->schema().num_key_columns()),
-
_key_group_cluster_key_idxes(std::move(key_group_cluster_key_idxes)) {}
+
_key_group_cluster_key_idxes(std::move(key_group_cluster_key_idxes)),
+ _context_stats(context_stats) {}
VerticalMergeIteratorContext(const VerticalMergeIteratorContext&) = delete;
VerticalMergeIteratorContext(VerticalMergeIteratorContext&&) = delete;
VerticalMergeIteratorContext& operator=(const
VerticalMergeIteratorContext&) = delete;
VerticalMergeIteratorContext& operator=(VerticalMergeIteratorContext&&) =
delete;
- ~VerticalMergeIteratorContext() = default;
+ ~VerticalMergeIteratorContext();
Status block_reset(const std::shared_ptr<Block>& block);
Status init(const StorageReadOptions& opts, CompactionSampleInfo*
sample_info = nullptr);
bool compare(const VerticalMergeIteratorContext& rhs) const;
@@ -230,9 +246,15 @@ public:
const std::shared_ptr<Block>& block_ptr() const { return _block; }
+ // No later row source references this context. The returned
IteratorRowRef/RowBatch keeps
+ // its own shared_ptr<Block>, so the segment reader and context-owned
blocks can be released.
+ void release_resources();
+
private:
// Load next block into _block
Status _load_next_block();
+ void _mark_active();
+ void _mark_inactive();
RowwiseIteratorUPtr _iter;
RowsetId _rowset_id;
@@ -250,6 +272,7 @@ private:
size_t _block_row_max = 0;
int64_t _num_key_columns;
const std::vector<uint32_t> _key_group_cluster_key_idxes;
+ VerticalCompactionContextStats* _context_stats;
size_t _cur_batch_num = 0;
// used to store data load from iterator->next_batch(Block*)
@@ -258,6 +281,7 @@ private:
std::list<std::shared_ptr<Block>> _block_list;
// use to identify whether it's first block load from RowwiseIterator
bool _is_first_row = true;
+ bool _is_active_context_counted = false;
bool _record_rowids = false;
std::vector<RowLocation> _block_row_locations;
};
@@ -271,6 +295,7 @@ public:
std::vector<RowsetId> rowset_ids, size_t
ori_return_cols,
KeysType keys_type, int32_t seq_col_idx,
RowSourcesBuffer* row_sources_buf,
+ VerticalCompactionContextStats* context_stats,
std::vector<uint32_t>
key_group_cluster_key_idxes)
: _origin_iters(std::move(iters)),
_iterator_init_flags(std::move(iterator_init_flags)),
@@ -279,6 +304,7 @@ public:
_keys_type(keys_type),
_seq_col_idx(seq_col_idx),
_row_sources_buf(row_sources_buf),
+ _context_stats(context_stats),
_key_group_cluster_key_idxes(std::move(key_group_cluster_key_idxes)) {}
~VerticalHeapMergeIterator() override = default;
@@ -323,6 +349,7 @@ private:
KeysType _keys_type;
int32_t _seq_col_idx = -1;
RowSourcesBuffer* _row_sources_buf;
+ VerticalCompactionContextStats* _context_stats;
uint32_t _merged_rows = 0;
StorageReadOptions _opts;
bool _record_rowids = false;
@@ -338,14 +365,16 @@ public:
std::vector<bool> iterator_init_flags,
std::vector<RowsetId> rowset_ids, size_t
ori_return_cols,
KeysType keys_type, int32_t seq_col_idx,
- RowSourcesBuffer* row_sources_buf)
+ RowSourcesBuffer* row_sources_buf,
+ VerticalCompactionContextStats* context_stats)
: _origin_iters(std::move(iters)),
_iterator_init_flags(std::move(iterator_init_flags)),
_rowset_ids(std::move(rowset_ids)),
_ori_return_cols(ori_return_cols),
_keys_type(keys_type),
_seq_col_idx(seq_col_idx),
- _row_sources_buf(row_sources_buf) {}
+ _row_sources_buf(row_sources_buf),
+ _context_stats(context_stats) {}
~VerticalFifoMergeIterator() override = default;
VerticalFifoMergeIterator(const VerticalFifoMergeIterator&) = delete;
@@ -377,6 +406,7 @@ private:
KeysType _keys_type;
int32_t _seq_col_idx = -1;
RowSourcesBuffer* _row_sources_buf;
+ VerticalCompactionContextStats* _context_stats;
uint32_t _merged_rows = 0;
StorageReadOptions _opts;
bool _record_rowids = false;
@@ -399,10 +429,12 @@ class VerticalMaskMergeIterator : public RowwiseIterator {
public:
// VerticalMaskMergeIterator takes the ownership of input iterators
VerticalMaskMergeIterator(std::vector<RowwiseIteratorUPtr>&& iters, size_t
ori_return_cols,
- RowSourcesBuffer* row_sources_buf)
+ RowSourcesBuffer* row_sources_buf,
+ VerticalCompactionContextStats* context_stats)
: _origin_iters(std::move(iters)),
_ori_return_cols(ori_return_cols),
- _row_sources_buf(row_sources_buf) {}
+ _row_sources_buf(row_sources_buf),
+ _context_stats(context_stats) {}
~VerticalMaskMergeIterator() override = default;
VerticalMaskMergeIterator(const VerticalMaskMergeIterator&) = delete;
@@ -429,6 +461,8 @@ private:
int64_t _get_size(Block* block) { return block->rows(); }
Status check_all_iter_finished();
+ // Advance the row-source cursor and release its context after the final
reference.
+ void consume_row_sources(uint16_t order, size_t count = 1);
// released after build ctx
std::vector<RowwiseIteratorUPtr> _origin_iters;
@@ -441,6 +475,7 @@ private:
int _block_row_max = 0;
size_t _filtered_rows = 0;
RowSourcesBuffer* _row_sources_buf;
+ VerticalCompactionContextStats* _context_stats;
StorageReadOptions _opts;
CompactionSampleInfo* _sample_info = nullptr;
};
@@ -450,16 +485,18 @@ std::shared_ptr<RowwiseIterator>
new_vertical_heap_merge_iterator(
std::vector<RowwiseIteratorUPtr>&& inputs, const std::vector<bool>&
iterator_init_flag,
const std::vector<RowsetId>& rowset_ids, size_t _ori_return_cols,
KeysType key_type,
uint32_t seq_col_idx, RowSourcesBuffer* row_sources_buf,
+ VerticalCompactionContextStats* context_stats,
std::vector<uint32_t> key_group_cluster_key_idxes);
std::shared_ptr<RowwiseIterator> new_vertical_fifo_merge_iterator(
std::vector<RowwiseIteratorUPtr>&& inputs, const std::vector<bool>&
iterator_init_flag,
const std::vector<RowsetId>& rowset_ids, size_t _ori_return_cols,
KeysType key_type,
- uint32_t seq_col_idx, RowSourcesBuffer* row_sources_buf);
+ uint32_t seq_col_idx, RowSourcesBuffer* row_sources_buf,
+ VerticalCompactionContextStats* context_stats);
std::shared_ptr<RowwiseIterator> new_vertical_mask_merge_iterator(
std::vector<RowwiseIteratorUPtr>&& inputs, size_t ori_return_cols,
- RowSourcesBuffer* row_sources_buf);
+ RowSourcesBuffer* row_sources_buf, VerticalCompactionContextStats*
context_stats);
#include "common/compile_check_end.h"
} // namespace doris
diff --git a/be/src/storage/merger.cpp b/be/src/storage/merger.cpp
index f574f63f056..f5364d48787 100644
--- a/be/src/storage/merger.cpp
+++ b/be/src/storage/merger.cpp
@@ -60,6 +60,7 @@
#include "storage/tablet/tablet_reader.h"
#include "storage/types.h"
#include "storage/utils.h"
+#include "util/defer_op.h"
#include "util/slice.h"
namespace doris {
@@ -250,10 +251,11 @@ Status Merger::vertical_compact_one_group(
const std::vector<RowsetReaderSharedPtr>& src_rowset_readers,
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, bool enable_sparse_optimization) {
+ CompactionSampleInfo* sample_info, VerticalCompactionContextStats*
context_stats,
+ bool enable_sparse_optimization) {
// build tablet reader
VLOG_NOTICE << "vertical compact one group, max_rows_per_segment=" <<
max_rows_per_segment;
- VerticalBlockReader reader(row_source_buf);
+ VerticalBlockReader reader(row_source_buf, context_stats);
TabletReader::ReaderParams reader_params;
reader_params.is_key_column_group = is_key;
reader_params.key_group_cluster_key_idxes = key_group_cluster_key_idxes;
@@ -494,6 +496,13 @@ Status Merger::vertical_merge_rowsets(BaseTabletSPtr
tablet, ReaderType reader_t
Statistics* stats_output,
VerticalCompactionProgressCallback
progress_cb) {
LOG(INFO) << "Start to do vertical compaction, tablet_id: " <<
tablet->tablet_id();
+ VerticalCompactionContextStats context_stats;
+ Defer log_context_stats {[&] {
+ DCHECK_EQ(context_stats.active_segment_contexts, 0);
+ LOG(INFO) << "Vertical compaction segment context statistics,
tablet_id: "
+ << tablet->tablet_id() << ",
vertical_compaction_active_segment_contexts_peak: "
+ << context_stats.active_segment_contexts_peak;
+ }};
std::vector<std::vector<uint32_t>> column_groups;
std::vector<uint32_t> key_group_cluster_key_idxes;
// If BE config vertical_compaction_num_columns_per_group has been
modified from
@@ -698,7 +707,8 @@ Status Merger::vertical_merge_rowsets(BaseTabletSPtr
tablet, ReaderType reader_t
Status st = vertical_compact_one_group(
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,
enable_sparse_optimization);
+ key_group_cluster_key_idxes, batch_size, &sample_info,
&context_stats,
+ enable_sparse_optimization);
{
std::unique_lock<std::mutex> lock(sample_info_lock);
sample_infos[i] = sample_info;
diff --git a/be/src/storage/merger.h b/be/src/storage/merger.h
index 4dcf4787c20..5185c97e6b0 100644
--- a/be/src/storage/merger.h
+++ b/be/src/storage/merger.h
@@ -38,6 +38,7 @@ class SegmentWriter;
class RowSourcesBuffer;
class VerticalBlockReader;
+struct VerticalCompactionContextStats;
using VerticalCompactionProgressCallback =
std::function<void(int64_t total_groups, int64_t completed_groups)>;
@@ -85,7 +86,7 @@ public:
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,
- bool enable_sparse_optimization = false);
+ VerticalCompactionContextStats* context_stats, bool
enable_sparse_optimization = false);
// for segcompaction
static Status vertical_compact_one_group(
diff --git a/be/test/storage/compaction/vertical_compaction_test.cpp
b/be/test/storage/compaction/vertical_compaction_test.cpp
index 079bbb3f9af..1d9f8abe14c 100644
--- a/be/test/storage/compaction/vertical_compaction_test.cpp
+++ b/be/test/storage/compaction/vertical_compaction_test.cpp
@@ -16,6 +16,7 @@
// specific language governing permissions and limitations
// under the License.
+#include <bvar/variable.h>
#include <gen_cpp/AgentService_types.h>
#include <gen_cpp/Descriptors_types.h>
#include <gen_cpp/PaloInternalService_types.h>
@@ -50,6 +51,7 @@
#include "io/io_common.h"
#include "json2pb/json_to_pb.h"
#include "runtime/exec_env.h"
+#include "runtime/thread_context.h"
#include "storage/compaction/compaction.h"
#include "storage/delete/delete_handler.h"
#include "storage/iterator/vertical_merge_iterator.h"
@@ -73,6 +75,7 @@
#include "storage/tablet/tablet_schema.h"
#include "storage/txn/txn_manager.h"
#include "storage/utils.h"
+#include "util/defer_op.h"
#include "util/uid_util.h"
namespace doris {
@@ -584,6 +587,9 @@ TEST_F(VerticalCompactionTest, TestRowSourcesBuffer) {
size_t limit = 10;
static_cast<void>(buffer.flush());
static_cast<void>(buffer.seek_to_begin());
+ EXPECT_FALSE(buffer.is_source_exhausted(0));
+ EXPECT_FALSE(buffer.is_source_exhausted(1));
+ EXPECT_FALSE(buffer.is_source_exhausted(2));
int idx = -1;
while (buffer.has_remaining().ok()) {
@@ -594,7 +600,12 @@ TEST_F(VerticalCompactionTest, TestRowSourcesBuffer) {
auto same = buffer.same_source_count(cur, limit);
EXPECT_EQ(same, 2);
buffer.advance(same);
+ EXPECT_TRUE(buffer.is_source_exhausted(cur));
}
+ static_cast<void>(buffer.seek_to_begin());
+ EXPECT_FALSE(buffer.is_source_exhausted(0));
+ EXPECT_FALSE(buffer.is_source_exhausted(1));
+ EXPECT_FALSE(buffer.is_source_exhausted(2));
RowSourcesBuffer buffer1(101, absolute_dir,
ReaderType::READER_CUMULATIVE_COMPACTION);
EXPECT_TRUE(buffer1.append(tmp_row_source).ok());
@@ -695,6 +706,9 @@ TEST_F(VerticalCompactionTest,
TestRowSourcesBufferSpillThreshold) {
++read_back;
}
EXPECT_EQ(read_back, expected_total);
+ for (uint16_t source = 0; source < 8; ++source) {
+ EXPECT_TRUE(buffer.is_source_exhausted(source));
+ }
}
TEST_F(VerticalCompactionTest, TestDupKeyVerticalMerge) {
@@ -1169,6 +1183,275 @@ TEST_F(VerticalCompactionTest,
<< "the higher sequence row should stay visible after compaction";
}
+TEST_F(VerticalCompactionTest,
TestUniqueKeyNonOverlappingSegmentContextRetention) {
+ constexpr uint32_t num_input_rowsets = 3;
+ constexpr uint32_t num_segments_per_rowset = 4;
+ constexpr uint32_t rows_per_segment = 64;
+ constexpr int64_t total_segments = num_input_rowsets *
num_segments_per_rowset;
+ constexpr int64_t total_rows = total_segments * rows_per_segment;
+
+ auto old_compaction_batch_size = config::compaction_batch_size;
+ auto old_sparse_threshold =
config::sparse_column_compaction_threshold_percent;
+ Defer restore_config {[&] {
+ config::compaction_batch_size = old_compaction_batch_size;
+ config::sparse_column_compaction_threshold_percent =
old_sparse_threshold;
+ }};
+ config::compaction_batch_size = 32;
+ config::sparse_column_compaction_threshold_percent = 0;
+
+ std::vector<std::vector<std::vector<std::tuple<int64_t, int64_t>>>>
input_data;
+ for (uint32_t rowset_id = 0; rowset_id < num_input_rowsets; ++rowset_id) {
+ std::vector<std::vector<std::tuple<int64_t, int64_t>>> rowset_data;
+ for (uint32_t segment_id = 0; segment_id < num_segments_per_rowset;
++segment_id) {
+ std::vector<std::tuple<int64_t, int64_t>> segment_data;
+ for (uint32_t row_id = 0; row_id < rows_per_segment; ++row_id) {
+ int64_t logical_row = segment_id * rows_per_segment + row_id;
+ int64_t key = logical_row * num_input_rowsets + rowset_id;
+ segment_data.emplace_back(key, key + 1);
+ }
+ rowset_data.emplace_back(std::move(segment_data));
+ }
+ input_data.emplace_back(std::move(rowset_data));
+ }
+
+ TabletSchemaSPtr tablet_schema = create_schema(UNIQUE_KEYS);
+ std::vector<RowsetSharedPtr> input_rowsets;
+ for (uint32_t rowset_id = 0; rowset_id < num_input_rowsets; ++rowset_id) {
+ auto rowset =
+ create_rowset(tablet_schema, NONOVERLAPPING,
input_data[rowset_id], rowset_id);
+ ASSERT_FALSE(rowset->rowset_meta()->is_segments_overlapping());
+ ASSERT_EQ(num_segments_per_rowset, rowset->num_segments());
+ input_rowsets.push_back(rowset);
+ }
+
+ TabletSharedPtr tablet = create_tablet(*tablet_schema, false);
+ auto run_case = [&](double sparse_threshold) {
+ config::sparse_column_compaction_threshold_percent = sparse_threshold;
+ tablet->compaction_density.store(1.0);
+
+ std::vector<RowsetReaderSharedPtr> input_rs_readers;
+ for (const auto& rowset : input_rowsets) {
+ RowsetReaderSharedPtr rs_reader;
+ ASSERT_TRUE(rowset->create_reader(&rs_reader).ok());
+ input_rs_readers.push_back(std::move(rs_reader));
+ }
+
+ auto writer_context =
+ create_rowset_writer_context(tablet_schema, NONOVERLAPPING,
UINT32_MAX,
+ {0,
input_rowsets.back()->end_version()});
+ auto res = RowsetFactory::create_rowset_writer(*engine_ref,
writer_context, true);
+ ASSERT_TRUE(res.has_value()) << res.error();
+ auto output_rs_writer = std::move(res).value();
+
+ ASSERT_EQ("0",
+
bvar::Variable::describe_exposed("vertical_compaction_active_segment_contexts"));
+
+ Merger::Statistics stats;
+ auto st = Merger::vertical_merge_rowsets(
+ tablet, ReaderType::READER_BASE_COMPACTION, *tablet_schema,
input_rs_readers,
+ output_rs_writer.get(), UINT32_MAX, total_segments, &stats);
+ ASSERT_TRUE(st.ok()) << st;
+
+ EXPECT_EQ("0",
+
bvar::Variable::describe_exposed("vertical_compaction_active_segment_contexts"));
+ EXPECT_EQ(total_rows, stats.output_rows);
+ EXPECT_EQ(0, stats.merged_rows);
+ EXPECT_EQ(0, stats.filtered_rows);
+
+ RowsetSharedPtr output_rowset;
+ ASSERT_EQ(Status::OK(), output_rs_writer->build(output_rowset));
+ ASSERT_TRUE(output_rowset);
+ EXPECT_EQ(total_rows, output_rowset->num_rows());
+ };
+
+ run_case(0);
+ run_case(1.0);
+}
+
+TEST_F(VerticalCompactionTest, TestUniqueKeySegmentContextMemoryAmplification)
{
+ constexpr uint32_t num_input_rowsets = 10;
+ constexpr uint32_t batch_size = 32;
+ constexpr uint32_t payload_size = 8 * 1024;
+
+ auto old_compaction_batch_size = config::compaction_batch_size;
+ auto old_sparse_threshold =
config::sparse_column_compaction_threshold_percent;
+ Defer restore_config {[&] {
+ config::compaction_batch_size = old_compaction_batch_size;
+ config::sparse_column_compaction_threshold_percent =
old_sparse_threshold;
+ }};
+ config::compaction_batch_size = batch_size;
+ config::sparse_column_compaction_threshold_percent = 0;
+
+ TabletSchemaSPtr tablet_schema = std::make_shared<TabletSchema>();
+ TabletSchemaPB tablet_schema_pb;
+ tablet_schema_pb.set_keys_type(UNIQUE_KEYS);
+ tablet_schema_pb.set_num_short_key_columns(1);
+ tablet_schema_pb.set_num_rows_per_row_block(1024);
+ tablet_schema_pb.set_compress_kind(COMPRESS_NONE);
+ tablet_schema_pb.set_next_column_unique_id(4);
+
+ ColumnPB* key_column = tablet_schema_pb.add_column();
+ key_column->set_unique_id(1);
+ key_column->set_name("key");
+ key_column->set_type("INT");
+ key_column->set_is_key(true);
+ key_column->set_length(4);
+ key_column->set_index_length(4);
+ key_column->set_is_nullable(false);
+ key_column->set_is_bf_column(false);
+
+ ColumnPB* value_column = tablet_schema_pb.add_column();
+ value_column->set_unique_id(2);
+ value_column->set_name("value");
+ value_column->set_type("VARCHAR");
+ value_column->set_is_key(false);
+ value_column->set_length(payload_size);
+ value_column->set_index_length(20);
+ value_column->set_is_nullable(false);
+ value_column->set_is_bf_column(false);
+
+ ColumnPB* delete_sign_column = tablet_schema_pb.add_column();
+ delete_sign_column->set_unique_id(3);
+ delete_sign_column->set_name(DELETE_SIGN);
+ delete_sign_column->set_type("TINYINT");
+ delete_sign_column->set_is_key(false);
+ delete_sign_column->set_length(1);
+ delete_sign_column->set_index_length(1);
+ delete_sign_column->set_is_nullable(false);
+ delete_sign_column->set_is_bf_column(false);
+
+ tablet_schema->init_from_pb(tablet_schema_pb);
+ TabletSharedPtr tablet = create_tablet(*tablet_schema, false);
+
+ struct Observation {
+ int64_t memory_peak;
+ };
+ std::vector<Observation> observations;
+
+ auto run_case = [&](uint32_t num_segments_per_rowset, uint32_t
rows_per_segment,
+ int64_t base_version) {
+ int64_t total_segments = num_input_rowsets * num_segments_per_rowset;
+ int64_t total_rows = total_segments * rows_per_segment;
+ std::string payload(payload_size, 'x');
+ std::vector<RowsetSharedPtr> input_rowsets;
+ std::vector<RowsetReaderSharedPtr> input_rs_readers;
+
+ for (uint32_t rowset_id = 0; rowset_id < num_input_rowsets;
++rowset_id) {
+ auto writer_context = create_rowset_writer_context(
+ tablet_schema, NONOVERLAPPING, UINT32_MAX,
+ {base_version + rowset_id, base_version + rowset_id});
+ auto res = RowsetFactory::create_rowset_writer(*engine_ref,
writer_context, true);
+ ASSERT_TRUE(res.has_value()) << res.error();
+ auto rowset_writer = std::move(res).value();
+
+ for (uint32_t segment_id = 0; segment_id <
num_segments_per_rowset; ++segment_id) {
+ Block block = tablet_schema->create_block();
+ auto columns = std::move(block).mutate_columns();
+ for (uint32_t row_id = 0; row_id < rows_per_segment; ++row_id)
{
+ int32_t logical_row = segment_id * rows_per_segment +
row_id;
+ int32_t key = logical_row * num_input_rowsets + rowset_id;
+ uint8_t delete_sign = 0;
+ columns[0]->insert_data(reinterpret_cast<const
char*>(&key), sizeof(key));
+ columns[1]->insert_data(payload.data(), payload.size());
+ columns[2]->insert_data(reinterpret_cast<const
char*>(&delete_sign),
+ sizeof(delete_sign));
+ }
+ ASSERT_TRUE(add_block_with_columns(rowset_writer.get(),
&block, &columns).ok());
+ ASSERT_TRUE(rowset_writer->flush().ok());
+ }
+
+ RowsetSharedPtr rowset;
+ ASSERT_EQ(Status::OK(), rowset_writer->build(rowset));
+ ASSERT_FALSE(rowset->rowset_meta()->is_segments_overlapping());
+ ASSERT_EQ(num_segments_per_rowset, rowset->num_segments());
+ ASSERT_EQ(num_segments_per_rowset * rows_per_segment,
rowset->num_rows());
+ input_rowsets.push_back(rowset);
+
+ RowsetReaderSharedPtr rs_reader;
+ ASSERT_TRUE(rowset->create_reader(&rs_reader).ok());
+ input_rs_readers.push_back(std::move(rs_reader));
+ }
+
+ auto output_writer_context =
+ create_rowset_writer_context(tablet_schema, NONOVERLAPPING,
UINT32_MAX,
+ {base_version, base_version +
num_input_rowsets - 1});
+ auto output_res =
+ RowsetFactory::create_rowset_writer(*engine_ref,
output_writer_context, true);
+ ASSERT_TRUE(output_res.has_value()) << output_res.error();
+ auto output_rs_writer = std::move(output_res).value();
+
+ ASSERT_EQ("0",
+
bvar::Variable::describe_exposed("vertical_compaction_active_segment_contexts"));
+
+ Merger::Statistics stats;
+ int64_t memory_peak = 0;
+ {
+ SCOPED_PEAK_MEM(&memory_peak);
+ auto st = Merger::vertical_merge_rowsets(
+ tablet, ReaderType::READER_BASE_COMPACTION,
*tablet_schema, input_rs_readers,
+ output_rs_writer.get(), UINT32_MAX, total_segments,
&stats);
+ ASSERT_TRUE(st.ok()) << st;
+ }
+
+ ASSERT_EQ("0",
+
bvar::Variable::describe_exposed("vertical_compaction_active_segment_contexts"));
+ ASSERT_EQ(total_rows, stats.output_rows);
+ ASSERT_EQ(0, stats.merged_rows);
+ ASSERT_EQ(0, stats.filtered_rows);
+
+ RowsetSharedPtr output_rowset;
+ ASSERT_EQ(Status::OK(), output_rs_writer->build(output_rowset));
+ ASSERT_TRUE(output_rowset);
+ ASSERT_EQ(total_rows, output_rowset->num_rows());
+
+ RowsetReaderContext reader_context;
+ reader_context.tablet_schema = tablet_schema;
+ reader_context.need_ordered_result = false;
+ std::vector<uint32_t> return_columns = {0, 1, 2};
+ reader_context.return_columns = &return_columns;
+ RowsetReaderSharedPtr output_rs_reader;
+ create_and_init_rowset_reader(output_rowset.get(), reader_context,
&output_rs_reader);
+
+ int64_t expected_key = 0;
+ Status read_status;
+ do {
+ Block output_block = tablet_schema->create_block();
+ read_status = output_rs_reader->next_batch(&output_block);
+ const auto& output_columns =
output_block.get_columns_with_type_and_name();
+ ASSERT_EQ(3, output_columns.size());
+ for (size_t row = 0; row < output_block.rows(); ++row) {
+ ASSERT_EQ(expected_key,
output_columns[0].column->get_int(row));
+ ASSERT_EQ(payload,
output_columns[1].column->get_data_at(row).to_string());
+ ASSERT_EQ(0, output_columns[2].column->get_int(row));
+ ++expected_key;
+ }
+ } while (read_status.ok());
+ ASSERT_TRUE(read_status.is<END_OF_FILE>()) << read_status;
+ ASSERT_EQ(total_rows, expected_key);
+
+ observations.push_back({memory_peak});
+ };
+
+ // Both cases contain exactly 32000 rows and the same 8 KiB value payload
per row.
+ // Only the segment distribution differs.
+ run_case(10, 320, 1000);
+ ASSERT_EQ(1, observations.size());
+ run_case(50, 64, 2000);
+ ASSERT_EQ(2, observations.size());
+
+ const auto& low_segment_case = observations[0];
+ const auto& high_segment_case = observations[1];
+ LOG(INFO) << "equal-data vertical compaction observation:
low_segments_memory_peak="
+ << low_segment_case.memory_peak
+ << ", high_segments_memory_peak=" <<
high_segment_case.memory_peak;
+ EXPECT_GT(low_segment_case.memory_peak, 0);
+ EXPECT_GT(high_segment_case.memory_peak, 0);
+ auto memory_peak_delta = low_segment_case.memory_peak >
high_segment_case.memory_peak
+ ? low_segment_case.memory_peak -
high_segment_case.memory_peak
+ : high_segment_case.memory_peak -
low_segment_case.memory_peak;
+ EXPECT_LT(memory_peak_delta, 2 * 1024 * 1024);
+}
+
TEST_F(VerticalCompactionTest, TestDupKeyVerticalMergeWithDelete) {
auto num_input_rowset = 2;
auto num_segments = 2;
@@ -1422,10 +1705,12 @@ TEST_F(VerticalCompactionTest, TestAggKeyVerticalMerge)
{
Merger::Statistics stats;
RowIdConversion rowid_conversion;
stats.rowid_conversion = &rowid_conversion;
+ ASSERT_EQ("0",
bvar::Variable::describe_exposed("vertical_compaction_active_segment_contexts"));
auto s = Merger::vertical_merge_rowsets(tablet,
ReaderType::READER_BASE_COMPACTION,
*tablet_schema, input_rs_readers,
output_rs_writer.get(), 100,
num_segments, &stats);
EXPECT_TRUE(s.ok());
+ EXPECT_EQ("0",
bvar::Variable::describe_exposed("vertical_compaction_active_segment_contexts"));
RowsetSharedPtr out_rowset;
EXPECT_EQ(Status::OK(), output_rs_writer->build(out_rowset));
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]