github-actions[bot] commented on code in PR #65963:
URL: https://github.com/apache/doris/pull/65963#discussion_r3701202816
##########
be/test/storage/compaction/vertical_compaction_test.cpp:
##########
@@ -837,6 +851,275 @@ TEST_F(VerticalCompactionTest,
TestUniqueKeyVerticalMerge) {
}
}
+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",
Review Comment:
Please add a path-specific retention oracle here and for the AGG path. Both
runs only assert that the process-wide counter is zero after
`vertical_merge_rowsets()` returns; because `~VerticalMergeIteratorContext()`
always calls `release_resources()`, those assertions still pass if the new
calls from `unique_key_next_batch()` or `next_row()` are removed. The only
memory/peak oracle below forces sparse optimization off, so it covers fallback
only. A deterministic task-local peak assertion for sparse UNIQUE and AGG would
make these new lifetime paths regressible.
##########
be/test/storage/compaction/vertical_compaction_test.cpp:
##########
@@ -837,6 +851,275 @@ TEST_F(VerticalCompactionTest,
TestUniqueKeyVerticalMerge) {
}
}
+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);
Review Comment:
This makes the unit test unnecessarily expensive and the oracle
platform-sensitive. Each case writes, compacts, and rereads 32,000 x 8 KiB of
uncompressed values; together the two cases perform about 1.5 GiB of payload
work and leave roughly 1 GiB of files until fixture teardown. Yet pass/fail is
a fixed 2 MiB delta that also includes the overhead of 400 extra segment
iterators/contexts in the high-segment case. Please assert the per-invocation
`active_segment_contexts_peak` directly with small payloads and segment counts,
keeping at most a small memory smoke check if needed.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]