| |
| // Licensed to the Apache Software Foundation (ASF) under one |
| // or more contributor license agreements. See the NOTICE file |
| // distributed with this work for additional information |
| // regarding copyright ownership. The ASF licenses this file |
| // to you under the Apache License, Version 2.0 (the |
| // "License"); you may not use this file except in compliance |
| // with the License. You may obtain a copy of the License at |
| // |
| // http://www.apache.org/licenses/LICENSE-2.0 |
| // |
| // Unless required by applicable law or agreed to in writing, |
| // software distributed under the License is distributed on an |
| // "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY |
| // KIND, either express or implied. See the License for the |
| // specific language governing permissions and limitations |
| // under the License. |
| |
| #include <bvar/variable.h> |
| #include <gen_cpp/AgentService_types.h> |
| #include <gen_cpp/Descriptors_types.h> |
| #include <gen_cpp/PaloInternalService_types.h> |
| #include <gen_cpp/Types_types.h> |
| #include <gen_cpp/olap_common.pb.h> |
| #include <gen_cpp/olap_file.pb.h> |
| #include <glog/logging.h> |
| #include <gtest/gtest-message.h> |
| #include <gtest/gtest-test-part.h> |
| #include <stdint.h> |
| #include <unistd.h> |
| |
| #include <iostream> |
| #include <memory> |
| #include <string> |
| #include <tuple> |
| #include <unordered_map> |
| #include <utility> |
| #include <vector> |
| |
| #include "common/status.h" |
| #include "core/assert_cast.h" |
| #include "core/block/block.h" |
| #include "core/block/column_with_type_and_name.h" |
| #include "core/column/column.h" |
| #include "core/column/column_nullable.h" |
| #include "core/data_type/data_type.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/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/options.h" |
| #include "storage/rowid_conversion.h" |
| #include "storage/rowset/beta_rowset.h" |
| #include "storage/rowset/rowset.h" |
| #include "storage/rowset/rowset_factory.h" |
| #include "storage/rowset/rowset_meta.h" |
| #include "storage/rowset/rowset_reader.h" |
| #include "storage/rowset/rowset_reader_context.h" |
| #include "storage/rowset/rowset_writer.h" |
| #include "storage/rowset/rowset_writer_context.h" |
| #include "storage/schema.h" |
| #include "storage/segment/segment.h" |
| #include "storage/storage_engine.h" |
| #include "storage/tablet/tablet.h" |
| #include "storage/tablet/tablet_meta.h" |
| #include "storage/tablet/tablet_schema.h" |
| #include "storage/utils.h" |
| #include "util/defer_op.h" |
| #include "util/uid_util.h" |
| |
| namespace doris { |
| using namespace ErrorCode; |
| |
| static const uint32_t MAX_PATH_LEN = 1024; |
| static StorageEngine* engine_ref = nullptr; |
| |
| class VerticalCompactionTest : public ::testing::Test { |
| protected: |
| void SetUp() override { |
| config::vertical_compaction_max_row_source_memory_mb = 1; |
| char buffer[MAX_PATH_LEN]; |
| EXPECT_NE(getcwd(buffer, MAX_PATH_LEN), nullptr); |
| absolute_dir = std::string(buffer) + kTestDir; |
| auto st = io::global_local_filesystem()->delete_directory(absolute_dir); |
| ASSERT_TRUE(st.ok()) << st; |
| st = io::global_local_filesystem()->create_directory(absolute_dir); |
| ASSERT_TRUE(st.ok()) << st; |
| EXPECT_TRUE(io::global_local_filesystem() |
| ->create_directory(absolute_dir + "/tablet_path") |
| .ok()); |
| |
| doris::EngineOptions options; |
| auto engine = std::make_unique<StorageEngine>(options); |
| engine_ref = engine.get(); |
| ExecEnv::GetInstance()->set_storage_engine(std::move(engine)); |
| |
| _data_dir = new DataDir(*engine_ref, absolute_dir, 100000000); |
| static_cast<void>(_data_dir->init()); |
| } |
| void TearDown() override { |
| SAFE_DELETE(_data_dir); |
| EXPECT_TRUE(io::global_local_filesystem()->delete_directory(absolute_dir).ok()); |
| engine_ref = nullptr; |
| ExecEnv::GetInstance()->set_storage_engine(nullptr); |
| } |
| |
| Status add_block_with_columns(RowsetWriter* rowset_writer, Block* block, |
| MutableColumns* columns) { |
| block->set_columns(std::move(*columns)); |
| auto st = rowset_writer->add_block(block); |
| *columns = std::move(*block).mutate_columns(); |
| return st; |
| } |
| |
| TabletSchemaSPtr create_schema(KeysType keys_type = DUP_KEYS, bool without_key = false) { |
| TabletSchemaSPtr tablet_schema = std::make_shared<TabletSchema>(); |
| TabletSchemaPB tablet_schema_pb; |
| tablet_schema_pb.set_keys_type(keys_type); |
| tablet_schema_pb.set_num_short_key_columns(without_key ? 0 : 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* column_1 = tablet_schema_pb.add_column(); |
| column_1->set_unique_id(1); |
| column_1->set_name("c1"); |
| column_1->set_type("INT"); |
| column_1->set_is_key(!without_key); |
| column_1->set_length(4); |
| column_1->set_index_length(4); |
| column_1->set_is_nullable(false); |
| column_1->set_is_bf_column(false); |
| |
| ColumnPB* column_2 = tablet_schema_pb.add_column(); |
| column_2->set_unique_id(2); |
| column_2->set_name("c2"); |
| column_2->set_type("INT"); |
| column_2->set_length(4); |
| column_2->set_index_length(4); |
| column_2->set_is_nullable(true); |
| column_2->set_is_key(false); |
| column_2->set_is_nullable(false); |
| column_2->set_is_bf_column(false); |
| |
| // unique table must contains the DELETE_SIGN column |
| if (keys_type == UNIQUE_KEYS) { |
| ColumnPB* column_3 = tablet_schema_pb.add_column(); |
| column_3->set_unique_id(3); |
| column_3->set_name(DELETE_SIGN); |
| column_3->set_type("TINYINT"); |
| column_3->set_length(1); |
| column_3->set_index_length(1); |
| column_3->set_is_nullable(false); |
| column_3->set_is_key(false); |
| column_3->set_is_nullable(false); |
| column_3->set_is_bf_column(false); |
| } |
| |
| tablet_schema->init_from_pb(tablet_schema_pb); |
| return tablet_schema; |
| } |
| |
| TabletSchemaSPtr create_agg_schema() { |
| TabletSchemaSPtr tablet_schema = std::make_shared<TabletSchema>(); |
| TabletSchemaPB tablet_schema_pb; |
| tablet_schema_pb.set_keys_type(KeysType::AGG_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* column_1 = tablet_schema_pb.add_column(); |
| column_1->set_unique_id(1); |
| column_1->set_name("c1"); |
| column_1->set_type("INT"); |
| column_1->set_is_key(true); |
| column_1->set_length(4); |
| column_1->set_index_length(4); |
| column_1->set_is_nullable(false); |
| column_1->set_is_bf_column(false); |
| |
| ColumnPB* column_2 = tablet_schema_pb.add_column(); |
| column_2->set_unique_id(2); |
| column_2->set_name("c2"); |
| column_2->set_type("INT"); |
| column_2->set_length(4); |
| column_2->set_index_length(4); |
| column_2->set_is_nullable(true); |
| column_2->set_is_key(false); |
| column_2->set_is_nullable(false); |
| column_2->set_is_bf_column(false); |
| column_2->set_aggregation("SUM"); |
| |
| tablet_schema->init_from_pb(tablet_schema_pb); |
| return tablet_schema; |
| } |
| |
| RowsetWriterContext create_rowset_writer_context(TabletSchemaSPtr tablet_schema, |
| const SegmentsOverlapPB& overlap, |
| uint32_t max_rows_per_segment, |
| Version version) { |
| static int64_t inc_id = 1000; |
| RowsetWriterContext rowset_writer_context; |
| RowsetId rowset_id; |
| rowset_id.init(inc_id); |
| rowset_writer_context.rowset_id = rowset_id; |
| rowset_writer_context.rowset_type = BETA_ROWSET; |
| rowset_writer_context.rowset_state = VISIBLE; |
| rowset_writer_context.tablet_schema = tablet_schema; |
| rowset_writer_context.tablet_path = absolute_dir + "/tablet_path"; |
| rowset_writer_context.version = version; |
| rowset_writer_context.segments_overlap = overlap; |
| rowset_writer_context.max_rows_per_segment = max_rows_per_segment; |
| inc_id++; |
| return rowset_writer_context; |
| } |
| |
| void create_and_init_rowset_reader(Rowset* rowset, RowsetReaderContext& context, |
| RowsetReaderSharedPtr* result) { |
| auto s = rowset->create_reader(result); |
| EXPECT_TRUE(s.ok()); |
| EXPECT_TRUE(*result != nullptr); |
| |
| s = (*result)->init(&context); |
| EXPECT_TRUE(s.ok()); |
| } |
| |
| RowsetSharedPtr create_rowset( |
| TabletSchemaSPtr tablet_schema, const SegmentsOverlapPB& overlap, |
| std::vector<std::vector<std::tuple<int64_t, int64_t>>> rowset_data, int64_t version) { |
| if (overlap == NONOVERLAPPING) { |
| for (auto i = 1; i < rowset_data.size(); i++) { |
| auto& last_seg_data = rowset_data[i - 1]; |
| auto& cur_seg_data = rowset_data[i]; |
| int64_t last_seg_max = std::get<0>(last_seg_data[last_seg_data.size() - 1]); |
| int64_t cur_seg_min = std::get<0>(cur_seg_data[0]); |
| EXPECT_LT(last_seg_max, cur_seg_min); |
| } |
| } |
| auto writer_context = create_rowset_writer_context(tablet_schema, overlap, UINT32_MAX, |
| {version, version}); |
| |
| auto res = RowsetFactory::create_rowset_writer(*engine_ref, writer_context, true); |
| EXPECT_TRUE(res.has_value()) << res.error(); |
| auto rowset_writer = std::move(res).value(); |
| |
| uint32_t num_rows = 0; |
| for (int i = 0; i < rowset_data.size(); ++i) { |
| Block block = tablet_schema->create_block(); |
| auto columns = std::move(block).mutate_columns(); |
| for (int rid = 0; rid < rowset_data[i].size(); ++rid) { |
| int32_t c1 = std::get<0>(rowset_data[i][rid]); |
| int32_t c2 = std::get<1>(rowset_data[i][rid]); |
| columns[0]->insert_data((const char*)&c1, sizeof(c1)); |
| columns[1]->insert_data((const char*)&c2, sizeof(c2)); |
| |
| if (tablet_schema->keys_type() == UNIQUE_KEYS) { |
| uint8_t num = 0; |
| columns[2]->insert_data((const char*)&num, sizeof(num)); |
| } |
| num_rows++; |
| } |
| auto s = add_block_with_columns(rowset_writer.get(), &block, &columns); |
| EXPECT_TRUE(s.ok()); |
| s = rowset_writer->flush(); |
| EXPECT_TRUE(s.ok()); |
| } |
| |
| RowsetSharedPtr rowset; |
| EXPECT_EQ(Status::OK(), rowset_writer->build(rowset)); |
| EXPECT_EQ(rowset_data.size(), rowset->rowset_meta()->num_segments()); |
| EXPECT_EQ(num_rows, rowset->rowset_meta()->num_rows()); |
| return rowset; |
| } |
| |
| void init_rs_meta(RowsetMetaSharedPtr& rs_meta, int64_t start, int64_t end) { |
| std::string json_rowset_meta = R"({ |
| "rowset_id": 540081, |
| "tablet_id": 15673, |
| "partition_id": 10000, |
| "tablet_schema_hash": 567997577, |
| "rowset_type": "BETA_ROWSET", |
| "rowset_state": "VISIBLE", |
| "empty": false |
| })"; |
| RowsetMetaPB rowset_meta_pb; |
| json2pb::JsonToProtoMessage(json_rowset_meta, &rowset_meta_pb); |
| rowset_meta_pb.set_start_version(start); |
| rowset_meta_pb.set_end_version(end); |
| rowset_meta_pb.set_creation_time(10000); |
| rs_meta->init_from_pb(rowset_meta_pb); |
| } |
| |
| RowsetSharedPtr create_delete_predicate(const TabletSchemaSPtr& schema, |
| DeletePredicatePB del_pred, int64_t version) { |
| RowsetMetaSharedPtr rsm(new RowsetMeta()); |
| init_rs_meta(rsm, version, version); |
| RowsetId id; |
| id.init(version * 1000); |
| rsm->set_rowset_id(id); |
| rsm->set_delete_predicate(std::move(del_pred)); |
| rsm->set_tablet_schema(schema); |
| return std::make_shared<BetaRowset>(schema, rsm, ""); |
| } |
| |
| TabletSharedPtr create_tablet(const TabletSchema& tablet_schema, |
| bool enable_unique_key_merge_on_write) { |
| std::vector<TColumn> cols; |
| std::unordered_map<uint32_t, uint32_t> col_ordinal_to_unique_id; |
| for (auto i = 0; i < tablet_schema.num_columns(); i++) { |
| const TabletColumn& column = tablet_schema.column(i); |
| TColumn col; |
| col.column_type.type = TPrimitiveType::INT; |
| col.__set_column_name(column.name()); |
| col.__set_is_key(column.is_key()); |
| cols.push_back(col); |
| col_ordinal_to_unique_id[i] = column.unique_id(); |
| } |
| |
| TTabletSchema t_tablet_schema; |
| t_tablet_schema.__set_short_key_column_count(tablet_schema.num_short_key_columns()); |
| t_tablet_schema.__set_schema_hash(3333); |
| if (tablet_schema.keys_type() == UNIQUE_KEYS) { |
| t_tablet_schema.__set_keys_type(TKeysType::UNIQUE_KEYS); |
| } else if (tablet_schema.keys_type() == DUP_KEYS) { |
| t_tablet_schema.__set_keys_type(TKeysType::DUP_KEYS); |
| } else if (tablet_schema.keys_type() == AGG_KEYS) { |
| t_tablet_schema.__set_keys_type(TKeysType::AGG_KEYS); |
| } |
| t_tablet_schema.__set_storage_type(TStorageType::COLUMN); |
| t_tablet_schema.__set_columns(cols); |
| TabletMetaSharedPtr tablet_meta( |
| new TabletMeta(2, 2, 2, 2, 2, 2, t_tablet_schema, 2, col_ordinal_to_unique_id, |
| UniqueId(1, 2), TTabletType::TABLET_TYPE_DISK, |
| TCompressionType::LZ4F, 0, enable_unique_key_merge_on_write)); |
| |
| TabletSharedPtr tablet(new Tablet(*engine_ref, tablet_meta, _data_dir)); |
| static_cast<void>(tablet->init()); |
| bool exists = false; |
| auto res = io::global_local_filesystem()->exists(tablet->tablet_path(), &exists); |
| EXPECT_TRUE(res.ok() && !exists); |
| res = io::global_local_filesystem()->create_directory(tablet->tablet_path()); |
| EXPECT_TRUE(res.ok()); |
| return tablet; |
| } |
| |
| // all rowset's data are same |
| void generate_input_data( |
| uint32_t num_input_rowset, uint32_t num_segments, uint32_t rows_per_segment, |
| const SegmentsOverlapPB& overlap, |
| std::vector<std::vector<std::vector<std::tuple<int64_t, int64_t>>>>& input_data) { |
| for (auto i = 0; i < num_input_rowset; i++) { |
| std::vector<std::vector<std::tuple<int64_t, int64_t>>> rowset_data; |
| for (auto j = 0; j < num_segments; j++) { |
| std::vector<std::tuple<int64_t, int64_t>> segment_data; |
| for (auto n = 0; n < rows_per_segment; n++) { |
| int64_t c1 = j * rows_per_segment + n; |
| int64_t c2 = c1 + 1; |
| segment_data.emplace_back(c1, c2); |
| } |
| rowset_data.emplace_back(segment_data); |
| } |
| input_data.emplace_back(rowset_data); |
| } |
| } |
| |
| void block_create(TabletSchemaSPtr tablet_schema, Block* block) { |
| block->clear(); |
| size_t num_columns = tablet_schema->num_columns(); |
| if (num_columns > 0 && tablet_schema->columns().back()->name() == BeConsts::ROW_STORE_COL) { |
| --num_columns; |
| } |
| std::vector<ColumnId> schema_column_ids(num_columns); |
| for (uint32_t cid = 0; cid < num_columns; ++cid) { |
| schema_column_ids[cid] = cid; |
| } |
| Schema schema(tablet_schema->columns(), schema_column_ids); |
| const auto& column_ids = schema.column_ids(); |
| for (size_t i = 0; i < schema.num_column_ids(); ++i) { |
| auto column_desc = schema.column(column_ids[i]); |
| auto data_type = Schema::get_data_type_ptr(*column_desc); |
| EXPECT_TRUE(data_type != nullptr); |
| auto column = data_type->create_column(); |
| block->insert(ColumnWithTypeAndName(std::move(column), data_type, column_desc->name())); |
| } |
| } |
| |
| private: |
| const std::string kTestDir = "/ut_dir/vertical_compaction_test"; |
| std::string absolute_dir; |
| DataDir* _data_dir = nullptr; |
| }; |
| |
| TEST_F(VerticalCompactionTest, TestRowSourcesBuffer) { |
| RowSourcesBuffer buffer(100, absolute_dir, ReaderType::READER_CUMULATIVE_COMPACTION); |
| RowSource s1(0, 0); |
| RowSource s2(0, 0); |
| RowSource s3(1, 1); |
| RowSource s4(1, 0); |
| RowSource s5(2, 0); |
| RowSource s6(2, 0); |
| std::vector<RowSource> tmp_row_source; |
| tmp_row_source.emplace_back(s1); |
| tmp_row_source.emplace_back(s2); |
| tmp_row_source.emplace_back(s3); |
| tmp_row_source.emplace_back(s4); |
| tmp_row_source.emplace_back(s5); |
| tmp_row_source.emplace_back(s6); |
| |
| EXPECT_TRUE(buffer.append(tmp_row_source).ok()); |
| EXPECT_EQ(buffer.total_size(), 6); |
| 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()) { |
| if (++idx == 1) { |
| EXPECT_TRUE(buffer.current().agg_flag()); |
| } |
| auto cur = buffer.current().get_source_num(); |
| 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()); |
| EXPECT_TRUE(buffer1.append(tmp_row_source).ok()); |
| buffer1.set_agg_flag(2, false); |
| buffer1.set_agg_flag(4, true); |
| static_cast<void>(buffer1.flush()); |
| static_cast<void>(buffer1.seek_to_begin()); |
| EXPECT_EQ(buffer1.total_size(), 12); |
| idx = -1; |
| while (buffer1.has_remaining().ok()) { |
| if (++idx == 1) { |
| EXPECT_FALSE(buffer1.current().agg_flag()); |
| } |
| if (++idx == 0) { |
| EXPECT_TRUE(buffer1.current().agg_flag()); |
| } |
| std::cout << buffer1.buf_idx() << std::endl; |
| auto cur = buffer1.current().get_source_num(); |
| auto same = buffer1.same_source_count(cur, limit); |
| EXPECT_EQ(same, 2); |
| buffer1.advance(same); |
| } |
| } |
| |
| // Regression test for RowSourcesBuffer::append spill threshold. |
| // |
| // Background: |
| // PaddedPODArray::allocated_bytes() returns the total allocated memory which |
| // INCLUDES pad_left and pad_right. These padding bytes are NOT usable for |
| // storing elements. Earlier, append() used `allocated_bytes() - size*sizeof` |
| // as "available room" to decide whether to skip spilling. This over-estimates |
| // the truly usable space (by pad_left + pad_right bytes), so when the buffer |
| // has already crossed the configured memory limit, append() may incorrectly |
| // decide that the upcoming push_back will fit without reallocation, skip the |
| // spill, and then push_back triggers a reallocation that doubles the buffer, |
| // exceeding the configured `vertical_compaction_max_row_source_memory_mb`. |
| // |
| // This test simulates the case by setting a very small memory limit (1 MB) and |
| // repeatedly appending row sources. After the first time the buffer crosses |
| // the limit, the next append must trigger a spill (file write + reset) instead |
| // of silently growing the in-memory buffer beyond the limit. |
| TEST_F(VerticalCompactionTest, TestRowSourcesBufferSpillThreshold) { |
| // 1 MB limit (set in SetUp as well, but make it explicit here). |
| config::vertical_compaction_max_row_source_memory_mb = 1; |
| const size_t mem_limit_bytes = |
| static_cast<size_t>(config::vertical_compaction_max_row_source_memory_mb) * 1024 * 1024; |
| |
| RowSourcesBuffer buffer(200, absolute_dir, ReaderType::READER_CUMULATIVE_COMPACTION); |
| |
| // Build a batch of row sources. Use a moderate batch size so that the |
| // buffer's allocated_bytes() can become very close to the limit before |
| // a single append crosses it. |
| constexpr size_t kBatchSize = 4096; |
| std::vector<RowSource> batch; |
| batch.reserve(kBatchSize); |
| for (size_t i = 0; i < kBatchSize; ++i) { |
| batch.emplace_back(static_cast<uint16_t>(i % 8), false); |
| } |
| |
| // Total elements that fit in the memory limit (a safe upper bound). |
| // Each element is 2 bytes (UInt16), so ~512K elements per MB. |
| const size_t total_appends = (mem_limit_bytes / sizeof(uint16_t)) * 4 / kBatchSize + 8; |
| |
| size_t expected_total = 0; |
| for (size_t i = 0; i < total_appends; ++i) { |
| ASSERT_TRUE(buffer.append(batch).ok()); |
| expected_total += kBatchSize; |
| |
| // Invariant: in-memory buffered_size() must never exceed what the |
| // memory limit allows (in elements). Otherwise the spill logic is |
| // broken (the bug described above). |
| // Allow a small slack equal to one batch because the spill check is |
| // performed BEFORE the push_back that crosses the threshold. |
| size_t buffered_elems = buffer.buffered_size(); |
| size_t buffered_bytes = buffered_elems * sizeof(uint16_t); |
| // After each append, buffered_bytes should be <= mem_limit + one batch size. |
| // It must NOT grow unboundedly (e.g., 2x of the limit due to PODArray |
| // reallocation that the buggy version would allow). |
| EXPECT_LE(buffered_bytes, mem_limit_bytes + kBatchSize * sizeof(uint16_t)) |
| << "RowSourcesBuffer in-memory size exceeded the configured limit, " |
| << "spill threshold logic is broken. iter=" << i |
| << ", buffered_elems=" << buffered_elems; |
| } |
| |
| EXPECT_EQ(buffer.total_size(), expected_total); |
| |
| // Make sure data is persisted and can be read back correctly. |
| ASSERT_TRUE(buffer.flush().ok()); |
| ASSERT_TRUE(buffer.seek_to_begin().ok()); |
| |
| size_t read_back = 0; |
| while (buffer.has_remaining().ok()) { |
| // Verify that the source num matches the pattern we wrote. |
| auto cur = buffer.current().get_source_num(); |
| EXPECT_EQ(cur, (read_back % kBatchSize) % 8); |
| buffer.advance(1); |
| ++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) { |
| auto num_input_rowset = 2; |
| auto num_segments = 2; |
| auto rows_per_segment = 2 * 100 * 1024; |
| SegmentsOverlapPB overlap = NONOVERLAPPING; |
| std::vector<std::vector<std::vector<std::tuple<int64_t, int64_t>>>> input_data; |
| generate_input_data(num_input_rowset, num_segments, rows_per_segment, overlap, input_data); |
| for (auto rs_id = 0; rs_id < input_data.size(); rs_id++) { |
| for (auto s_id = 0; s_id < input_data[rs_id].size(); s_id++) { |
| for (auto row_id = 0; row_id < input_data[rs_id][s_id].size(); row_id++) { |
| LOG(INFO) << "input data: " << std::get<0>(input_data[rs_id][s_id][row_id]) << " " |
| << std::get<1>(input_data[rs_id][s_id][row_id]); |
| } |
| } |
| } |
| |
| TabletSchemaSPtr tablet_schema = create_schema(); |
| // create input rowset |
| std::vector<RowsetSharedPtr> input_rowsets; |
| SegmentsOverlapPB new_overlap = overlap; |
| for (auto i = 0; i < num_input_rowset; i++) { |
| if (overlap == OVERLAP_UNKNOWN) { |
| if (i == 0) { |
| new_overlap = NONOVERLAPPING; |
| } else { |
| new_overlap = OVERLAPPING; |
| } |
| } |
| RowsetSharedPtr rowset = create_rowset(tablet_schema, new_overlap, input_data[i], i); |
| input_rowsets.push_back(rowset); |
| } |
| // create input rowset reader |
| std::vector<RowsetReaderSharedPtr> input_rs_readers; |
| for (auto& rowset : input_rowsets) { |
| RowsetReaderSharedPtr rs_reader; |
| ASSERT_TRUE(rowset->create_reader(&rs_reader).ok()); |
| input_rs_readers.push_back(std::move(rs_reader)); |
| } |
| |
| // create output rowset writer |
| auto writer_context = create_rowset_writer_context(tablet_schema, NONOVERLAPPING, 3456, |
| {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(); |
| |
| // merge input rowset |
| TabletSharedPtr tablet = create_tablet(*tablet_schema, false); |
| Merger::Statistics stats; |
| RowIdConversion rowid_conversion; |
| stats.rowid_conversion = &rowid_conversion; |
| auto s = Merger::vertical_merge_rowsets(tablet, ReaderType::READER_BASE_COMPACTION, |
| *tablet_schema, input_rs_readers, |
| output_rs_writer.get(), 100, num_segments, &stats); |
| ASSERT_TRUE(s.ok()) << s; |
| RowsetSharedPtr out_rowset; |
| EXPECT_EQ(Status::OK(), output_rs_writer->build(out_rowset)); |
| ASSERT_TRUE(out_rowset); |
| |
| // create output rowset reader |
| RowsetReaderContext reader_context; |
| reader_context.tablet_schema = tablet_schema; |
| reader_context.need_ordered_result = false; |
| std::vector<uint32_t> return_columns = {0, 1}; |
| reader_context.return_columns = &return_columns; |
| RowsetReaderSharedPtr output_rs_reader; |
| LOG(INFO) << "create rowset reader in test"; |
| create_and_init_rowset_reader(out_rowset.get(), reader_context, &output_rs_reader); |
| |
| // read output rowset data |
| Block output_block; |
| std::vector<std::tuple<int64_t, int64_t>> output_data; |
| do { |
| block_create(tablet_schema, &output_block); |
| s = output_rs_reader->next_batch(&output_block); |
| auto columns = output_block.get_columns_with_type_and_name(); |
| EXPECT_EQ(columns.size(), 2); |
| for (auto i = 0; i < output_block.rows(); i++) { |
| output_data.emplace_back(columns[0].column->get_int(i), columns[1].column->get_int(i)); |
| } |
| } while (s == Status::OK()); |
| EXPECT_EQ(Status::Error<END_OF_FILE>(""), s); |
| EXPECT_EQ(out_rowset->rowset_meta()->num_rows(), output_data.size()); |
| EXPECT_EQ(output_data.size(), num_input_rowset * num_segments * rows_per_segment); |
| // check vertical compaction result |
| for (auto id = 0; id < output_data.size(); id++) { |
| LOG(INFO) << "output data: " << std::get<0>(output_data[id]) << " " |
| << std::get<1>(output_data[id]); |
| } |
| int dst_id = 0; |
| for (auto rs_id = 0; rs_id < input_data.size(); rs_id++) { |
| dst_id = 0; |
| for (auto s_id = 0; s_id < input_data[rs_id].size(); s_id++) { |
| for (auto row_id = 0; row_id < input_data[rs_id][s_id].size(); row_id++) { |
| LOG(INFO) << "input data: " << std::get<0>(input_data[rs_id][s_id][row_id]) << " " |
| << std::get<1>(input_data[rs_id][s_id][row_id]); |
| EXPECT_EQ(std::get<0>(input_data[rs_id][s_id][row_id]), |
| std::get<0>(output_data[dst_id])); |
| EXPECT_EQ(std::get<1>(input_data[rs_id][s_id][row_id]), |
| std::get<1>(output_data[dst_id])); |
| dst_id += 2; |
| } |
| } |
| } |
| } |
| |
| TEST_F(VerticalCompactionTest, TestDupWithoutKeyVerticalMerge) { |
| auto num_input_rowset = 2; |
| auto num_segments = 2; |
| auto rows_per_segment = 2 * 100 * 1024; |
| SegmentsOverlapPB overlap = NONOVERLAPPING; |
| std::vector<std::vector<std::vector<std::tuple<int64_t, int64_t>>>> input_data; |
| generate_input_data(num_input_rowset, num_segments, rows_per_segment, overlap, input_data); |
| for (auto rs_id = 0; rs_id < input_data.size(); rs_id++) { |
| for (auto s_id = 0; s_id < input_data[rs_id].size(); s_id++) { |
| for (auto row_id = 0; row_id < input_data[rs_id][s_id].size(); row_id++) { |
| LOG(INFO) << "input data: " << std::get<0>(input_data[rs_id][s_id][row_id]) << " " |
| << std::get<1>(input_data[rs_id][s_id][row_id]); |
| } |
| } |
| } |
| |
| TabletSchemaSPtr tablet_schema = create_schema(DUP_KEYS, true); |
| // create input rowset |
| std::vector<RowsetSharedPtr> input_rowsets; |
| SegmentsOverlapPB new_overlap = overlap; |
| for (auto i = 0; i < num_input_rowset; i++) { |
| if (overlap == OVERLAP_UNKNOWN) { |
| if (i == 0) { |
| new_overlap = NONOVERLAPPING; |
| } else { |
| new_overlap = OVERLAPPING; |
| } |
| } |
| RowsetSharedPtr rowset = create_rowset(tablet_schema, new_overlap, input_data[i], i); |
| input_rowsets.push_back(rowset); |
| } |
| // create input rowset reader |
| std::vector<RowsetReaderSharedPtr> input_rs_readers; |
| for (auto& rowset : input_rowsets) { |
| RowsetReaderSharedPtr rs_reader; |
| EXPECT_TRUE(rowset->create_reader(&rs_reader).ok()); |
| input_rs_readers.push_back(std::move(rs_reader)); |
| } |
| |
| // create output rowset writer |
| auto writer_context = create_rowset_writer_context(tablet_schema, NONOVERLAPPING, 3456, |
| {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(); |
| |
| // merge input rowset |
| TabletSharedPtr tablet = create_tablet(*tablet_schema, false); |
| Merger::Statistics stats; |
| RowIdConversion rowid_conversion; |
| stats.rowid_conversion = &rowid_conversion; |
| auto s = Merger::vertical_merge_rowsets(tablet, ReaderType::READER_BASE_COMPACTION, |
| *tablet_schema, input_rs_readers, |
| output_rs_writer.get(), 100, num_segments, &stats); |
| ASSERT_TRUE(s.ok()) << s; |
| RowsetSharedPtr out_rowset; |
| EXPECT_EQ(Status::OK(), output_rs_writer->build(out_rowset)); |
| |
| // create output rowset reader |
| RowsetReaderContext reader_context; |
| reader_context.tablet_schema = tablet_schema; |
| reader_context.need_ordered_result = false; |
| std::vector<uint32_t> return_columns = {0, 1}; |
| reader_context.return_columns = &return_columns; |
| RowsetReaderSharedPtr output_rs_reader; |
| LOG(INFO) << "create rowset reader in test"; |
| create_and_init_rowset_reader(out_rowset.get(), reader_context, &output_rs_reader); |
| |
| // read output rowset data |
| Block output_block; |
| std::vector<std::tuple<int64_t, int64_t>> output_data; |
| do { |
| block_create(tablet_schema, &output_block); |
| s = output_rs_reader->next_batch(&output_block); |
| auto columns = output_block.get_columns_with_type_and_name(); |
| EXPECT_EQ(columns.size(), 2); |
| for (auto i = 0; i < output_block.rows(); i++) { |
| output_data.emplace_back(columns[0].column->get_int(i), columns[1].column->get_int(i)); |
| } |
| } while (s == Status::OK()); |
| EXPECT_EQ(Status::Error<END_OF_FILE>(""), s); |
| EXPECT_EQ(out_rowset->rowset_meta()->num_rows(), output_data.size()); |
| EXPECT_EQ(output_data.size(), num_input_rowset * num_segments * rows_per_segment); |
| // check vertical compaction result |
| for (auto id = 0; id < output_data.size(); id++) { |
| LOG(INFO) << "output data: " << std::get<0>(output_data[id]) << " " |
| << std::get<1>(output_data[id]); |
| } |
| int dst_id = 0; |
| |
| for (auto rs_id = 0; rs_id < input_data.size(); rs_id++) { |
| dst_id = 0; |
| for (auto s_id = 0; s_id < input_data[rs_id].size(); s_id++) { |
| for (auto row_id = 0; row_id < input_data[rs_id][s_id].size(); row_id++) { |
| LOG(INFO) << "input data: " << std::get<0>(input_data[rs_id][s_id][row_id]) << " " |
| << std::get<1>(input_data[rs_id][s_id][row_id]); |
| EXPECT_EQ(std::get<0>(input_data[rs_id][s_id][row_id]), |
| std::get<0>(output_data[dst_id])); |
| EXPECT_EQ(std::get<1>(input_data[rs_id][s_id][row_id]), |
| std::get<1>(output_data[dst_id])); |
| dst_id += 1; |
| } |
| } |
| } |
| } |
| |
| TEST_F(VerticalCompactionTest, TestUniqueKeyVerticalMerge) { |
| auto num_input_rowset = 2; |
| auto num_segments = 2; |
| auto rows_per_segment = 2 * 100 * 1024; |
| SegmentsOverlapPB overlap = NONOVERLAPPING; |
| std::vector<std::vector<std::vector<std::tuple<int64_t, int64_t>>>> input_data; |
| generate_input_data(num_input_rowset, num_segments, rows_per_segment, overlap, input_data); |
| for (auto rs_id = 0; rs_id < input_data.size(); rs_id++) { |
| for (auto s_id = 0; s_id < input_data[rs_id].size(); s_id++) { |
| for (auto row_id = 0; row_id < input_data[rs_id][s_id].size(); row_id++) { |
| LOG(INFO) << "input data: " << std::get<0>(input_data[rs_id][s_id][row_id]) << " " |
| << std::get<1>(input_data[rs_id][s_id][row_id]); |
| } |
| } |
| } |
| |
| TabletSchemaSPtr tablet_schema = create_schema(UNIQUE_KEYS); |
| // create input rowset |
| std::vector<RowsetSharedPtr> input_rowsets; |
| SegmentsOverlapPB new_overlap = overlap; |
| for (auto i = 0; i < num_input_rowset; i++) { |
| if (overlap == OVERLAP_UNKNOWN) { |
| if (i == 0) { |
| new_overlap = NONOVERLAPPING; |
| } else { |
| new_overlap = OVERLAPPING; |
| } |
| } |
| RowsetSharedPtr rowset = create_rowset(tablet_schema, new_overlap, input_data[i], i); |
| input_rowsets.push_back(rowset); |
| } |
| // create input rowset reader |
| std::vector<RowsetReaderSharedPtr> input_rs_readers; |
| for (auto& rowset : input_rowsets) { |
| RowsetReaderSharedPtr rs_reader; |
| EXPECT_TRUE(rowset->create_reader(&rs_reader).ok()); |
| input_rs_readers.push_back(std::move(rs_reader)); |
| } |
| |
| // create output rowset writer |
| auto writer_context = create_rowset_writer_context(tablet_schema, NONOVERLAPPING, 3456, |
| {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(); |
| |
| // merge input rowset |
| TabletSharedPtr tablet = create_tablet(*tablet_schema, false); |
| Merger::Statistics stats; |
| RowIdConversion rowid_conversion; |
| stats.rowid_conversion = &rowid_conversion; |
| auto s = Merger::vertical_merge_rowsets(tablet, ReaderType::READER_BASE_COMPACTION, |
| *tablet_schema, input_rs_readers, |
| output_rs_writer.get(), 10000, num_segments, &stats); |
| EXPECT_TRUE(s.ok()); |
| RowsetSharedPtr out_rowset; |
| EXPECT_EQ(Status::OK(), output_rs_writer->build(out_rowset)); |
| |
| // create output rowset reader |
| RowsetReaderContext reader_context; |
| reader_context.tablet_schema = tablet_schema; |
| reader_context.need_ordered_result = false; |
| std::vector<uint32_t> return_columns = {0, 1}; |
| reader_context.return_columns = &return_columns; |
| RowsetReaderSharedPtr output_rs_reader; |
| LOG(INFO) << "create rowset reader in test"; |
| create_and_init_rowset_reader(out_rowset.get(), reader_context, &output_rs_reader); |
| |
| // read output rowset data |
| Block output_block; |
| std::vector<std::tuple<int64_t, int64_t>> output_data; |
| do { |
| block_create(tablet_schema, &output_block); |
| s = output_rs_reader->next_batch(&output_block); |
| auto columns = output_block.get_columns_with_type_and_name(); |
| EXPECT_EQ(columns.size(), 2); |
| for (auto i = 0; i < output_block.rows(); i++) { |
| output_data.emplace_back(columns[0].column->get_int(i), columns[1].column->get_int(i)); |
| } |
| } while (s == Status::OK()); |
| EXPECT_EQ(Status::Error<END_OF_FILE>(""), s); |
| EXPECT_EQ(out_rowset->rowset_meta()->num_rows(), output_data.size()); |
| EXPECT_EQ(output_data.size(), num_segments * rows_per_segment); |
| // check vertical compaction result |
| for (auto id = 0; id < output_data.size(); id++) { |
| LOG(INFO) << "output data: " << std::get<0>(output_data[id]) << " " |
| << std::get<1>(output_data[id]); |
| } |
| int dst_id = 0; |
| for (auto s_id = 0; s_id < input_data[0].size(); s_id++) { |
| for (auto row_id = 0; row_id < input_data[0][s_id].size(); row_id++) { |
| EXPECT_EQ(std::get<0>(input_data[0][s_id][row_id]), std::get<0>(output_data[dst_id])); |
| EXPECT_EQ(std::get<1>(input_data[0][s_id][row_id]), std::get<1>(output_data[dst_id])); |
| dst_id++; |
| } |
| } |
| } |
| |
| 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; |
| auto rows_per_segment = 2 * 100 * 1024; |
| SegmentsOverlapPB overlap = NONOVERLAPPING; |
| std::vector<std::vector<std::vector<std::tuple<int64_t, int64_t>>>> input_data; |
| generate_input_data(num_input_rowset, num_segments, rows_per_segment, overlap, input_data); |
| for (auto rs_id = 0; rs_id < input_data.size(); rs_id++) { |
| for (auto s_id = 0; s_id < input_data[rs_id].size(); s_id++) { |
| for (auto row_id = 0; row_id < input_data[rs_id][s_id].size(); row_id++) { |
| LOG(INFO) << "input data: " << std::get<0>(input_data[rs_id][s_id][row_id]) << " " |
| << std::get<1>(input_data[rs_id][s_id][row_id]); |
| } |
| } |
| } |
| |
| TabletSchemaSPtr tablet_schema = create_schema(DUP_KEYS); |
| TabletSharedPtr tablet = create_tablet(*tablet_schema, false); |
| // create input rowset |
| SegmentsOverlapPB new_overlap = overlap; |
| std::vector<RowsetSharedPtr> input_rowsets; |
| for (auto i = 0; i < num_input_rowset; i++) { |
| if (overlap == OVERLAP_UNKNOWN) { |
| new_overlap = (i == 0) ? NONOVERLAPPING : OVERLAPPING; |
| } |
| input_rowsets.push_back(create_rowset(tablet_schema, new_overlap, input_data[i], i)); |
| } |
| |
| // delete data with key < 100 |
| std::vector<TCondition> conditions; |
| TCondition condition; |
| condition.column_name = tablet->tablet_schema()->column(0).name(); |
| condition.condition_op = "<"; |
| condition.condition_values.clear(); |
| condition.condition_values.push_back("100"); |
| conditions.push_back(condition); |
| DeletePredicatePB del_pred; |
| auto st = DeleteHandler::generate_delete_predicate(*tablet->tablet_schema(), conditions, |
| &del_pred); |
| ASSERT_TRUE(st.ok()) << st; |
| input_rowsets.push_back(create_delete_predicate(tablet->tablet_schema(), std::move(del_pred), |
| num_input_rowset)); |
| |
| // create input rowset reader |
| std::vector<RowsetReaderSharedPtr> input_rs_readers; |
| for (auto& rowset : input_rowsets) { |
| RowsetReaderSharedPtr rs_reader; |
| ASSERT_TRUE(rowset->create_reader(&rs_reader).ok()); |
| input_rs_readers.push_back(std::move(rs_reader)); |
| } |
| |
| // create output rowset writer |
| auto writer_context = create_rowset_writer_context(tablet_schema, NONOVERLAPPING, 3456, |
| {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(); |
| // merge input rowset |
| Merger::Statistics stats; |
| RowIdConversion rowid_conversion; |
| stats.rowid_conversion = &rowid_conversion; |
| st = Merger::vertical_merge_rowsets(tablet, ReaderType::READER_BASE_COMPACTION, *tablet_schema, |
| input_rs_readers, output_rs_writer.get(), 100, num_segments, |
| &stats); |
| ASSERT_TRUE(st.ok()) << st; |
| RowsetSharedPtr out_rowset; |
| EXPECT_EQ(Status::OK(), output_rs_writer->build(out_rowset)); |
| ASSERT_TRUE(out_rowset); |
| |
| // create output rowset reader |
| RowsetReaderContext reader_context; |
| reader_context.tablet_schema = tablet_schema; |
| reader_context.need_ordered_result = false; |
| std::vector<uint32_t> return_columns = {0, 1}; |
| reader_context.return_columns = &return_columns; |
| RowsetReaderSharedPtr output_rs_reader; |
| LOG(INFO) << "create rowset reader in test"; |
| create_and_init_rowset_reader(out_rowset.get(), reader_context, &output_rs_reader); |
| |
| // read output rowset data |
| Block output_block; |
| std::vector<std::tuple<int64_t, int64_t>> output_data; |
| do { |
| block_create(tablet_schema, &output_block); |
| st = output_rs_reader->next_batch(&output_block); |
| auto columns = output_block.get_columns_with_type_and_name(); |
| EXPECT_EQ(columns.size(), 2); |
| for (auto i = 0; i < output_block.rows(); i++) { |
| output_data.emplace_back(columns[0].column->get_int(i), columns[1].column->get_int(i)); |
| } |
| } while (st.ok()); |
| EXPECT_TRUE(st.is<END_OF_FILE>()) << st; |
| EXPECT_EQ(out_rowset->rowset_meta()->num_rows(), output_data.size()); |
| EXPECT_EQ(output_data.size(), |
| num_input_rowset * num_segments * rows_per_segment - num_input_rowset * 100); |
| // All keys less than 1000 are deleted by delete handler |
| for (auto& item : output_data) { |
| ASSERT_GE(std::get<0>(item), 100); |
| } |
| } |
| |
| TEST_F(VerticalCompactionTest, TestDupWithoutKeyVerticalMergeWithDelete) { |
| auto num_input_rowset = 2; |
| auto num_segments = 2; |
| auto rows_per_segment = 2 * 100 * 1024; |
| SegmentsOverlapPB overlap = NONOVERLAPPING; |
| std::vector<std::vector<std::vector<std::tuple<int64_t, int64_t>>>> input_data; |
| generate_input_data(num_input_rowset, num_segments, rows_per_segment, overlap, input_data); |
| for (auto rs_id = 0; rs_id < input_data.size(); rs_id++) { |
| for (auto s_id = 0; s_id < input_data[rs_id].size(); s_id++) { |
| for (auto row_id = 0; row_id < input_data[rs_id][s_id].size(); row_id++) { |
| LOG(INFO) << "input data: " << std::get<0>(input_data[rs_id][s_id][row_id]) << " " |
| << std::get<1>(input_data[rs_id][s_id][row_id]); |
| } |
| } |
| } |
| |
| TabletSchemaSPtr tablet_schema = create_schema(DUP_KEYS, true); |
| TabletSharedPtr tablet = create_tablet(*tablet_schema, false); |
| // create input rowset |
| SegmentsOverlapPB new_overlap = overlap; |
| std::vector<RowsetSharedPtr> input_rowsets; |
| for (auto i = 0; i < num_input_rowset; i++) { |
| if (overlap == OVERLAP_UNKNOWN) { |
| new_overlap = (i == 0) ? NONOVERLAPPING : OVERLAPPING; |
| } |
| input_rowsets.push_back(create_rowset(tablet_schema, new_overlap, input_data[i], i)); |
| } |
| |
| // delete data with key < 100 |
| std::vector<TCondition> conditions; |
| TCondition condition; |
| condition.column_name = tablet->tablet_schema()->column(0).name(); |
| condition.condition_op = "<"; |
| condition.condition_values.clear(); |
| condition.condition_values.push_back("100"); |
| conditions.push_back(condition); |
| DeletePredicatePB del_pred; |
| auto st = DeleteHandler::generate_delete_predicate(*tablet->tablet_schema(), conditions, |
| &del_pred); |
| ASSERT_TRUE(st.ok()) << st; |
| input_rowsets.push_back(create_delete_predicate(tablet->tablet_schema(), std::move(del_pred), |
| num_input_rowset)); |
| |
| // create input rowset reader |
| std::vector<RowsetReaderSharedPtr> input_rs_readers; |
| for (auto& rowset : input_rowsets) { |
| RowsetReaderSharedPtr rs_reader; |
| ASSERT_TRUE(rowset->create_reader(&rs_reader).ok()); |
| input_rs_readers.push_back(std::move(rs_reader)); |
| } |
| |
| // create output rowset writer |
| auto writer_context = create_rowset_writer_context(tablet_schema, NONOVERLAPPING, 3456, |
| {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(); |
| // merge input rowset |
| Merger::Statistics stats; |
| RowIdConversion rowid_conversion; |
| stats.rowid_conversion = &rowid_conversion; |
| st = Merger::vertical_merge_rowsets(tablet, ReaderType::READER_BASE_COMPACTION, *tablet_schema, |
| input_rs_readers, output_rs_writer.get(), 100, num_segments, |
| &stats); |
| ASSERT_TRUE(st.ok()) << st; |
| RowsetSharedPtr out_rowset; |
| EXPECT_EQ(Status::OK(), output_rs_writer->build(out_rowset)); |
| ASSERT_TRUE(out_rowset); |
| |
| // create output rowset reader |
| RowsetReaderContext reader_context; |
| reader_context.tablet_schema = tablet_schema; |
| reader_context.need_ordered_result = false; |
| std::vector<uint32_t> return_columns = {0, 1}; |
| reader_context.return_columns = &return_columns; |
| RowsetReaderSharedPtr output_rs_reader; |
| LOG(INFO) << "create rowset reader in test"; |
| create_and_init_rowset_reader(out_rowset.get(), reader_context, &output_rs_reader); |
| |
| // read output rowset data |
| Block output_block; |
| std::vector<std::tuple<int64_t, int64_t>> output_data; |
| do { |
| block_create(tablet_schema, &output_block); |
| st = output_rs_reader->next_batch(&output_block); |
| auto columns = output_block.get_columns_with_type_and_name(); |
| EXPECT_EQ(columns.size(), 2); |
| for (auto i = 0; i < output_block.rows(); i++) { |
| output_data.emplace_back(columns[0].column->get_int(i), columns[1].column->get_int(i)); |
| } |
| } while (st.ok()); |
| EXPECT_TRUE(st.is<END_OF_FILE>()) << st; |
| EXPECT_EQ(out_rowset->rowset_meta()->num_rows(), output_data.size()); |
| EXPECT_EQ(output_data.size(), |
| num_input_rowset * num_segments * rows_per_segment - num_input_rowset * 100); |
| // All keys less than 1000 are deleted by delete handler |
| for (auto& item : output_data) { |
| ASSERT_GE(std::get<0>(item), 100); |
| } |
| } |
| |
| TEST_F(VerticalCompactionTest, TestAggKeyVerticalMerge) { |
| auto num_input_rowset = 2; |
| auto num_segments = 2; |
| auto rows_per_segment = 2 * 100 * 1024; |
| SegmentsOverlapPB overlap = NONOVERLAPPING; |
| std::vector<std::vector<std::vector<std::tuple<int64_t, int64_t>>>> input_data; |
| generate_input_data(num_input_rowset, num_segments, rows_per_segment, overlap, input_data); |
| for (auto rs_id = 0; rs_id < input_data.size(); rs_id++) { |
| for (auto s_id = 0; s_id < input_data[rs_id].size(); s_id++) { |
| for (auto row_id = 0; row_id < input_data[rs_id][s_id].size(); row_id++) { |
| LOG(INFO) << "input data: " << std::get<0>(input_data[rs_id][s_id][row_id]) << " " |
| << std::get<1>(input_data[rs_id][s_id][row_id]); |
| } |
| } |
| } |
| |
| TabletSchemaSPtr tablet_schema = create_agg_schema(); |
| // create input rowset |
| std::vector<RowsetSharedPtr> input_rowsets; |
| SegmentsOverlapPB new_overlap = overlap; |
| for (auto i = 0; i < num_input_rowset; i++) { |
| if (overlap == OVERLAP_UNKNOWN) { |
| if (i == 0) { |
| new_overlap = NONOVERLAPPING; |
| } else { |
| new_overlap = OVERLAPPING; |
| } |
| } |
| RowsetSharedPtr rowset = create_rowset(tablet_schema, new_overlap, input_data[i], i); |
| input_rowsets.push_back(rowset); |
| } |
| // create input rowset reader |
| std::vector<RowsetReaderSharedPtr> input_rs_readers; |
| for (auto& rowset : input_rowsets) { |
| RowsetReaderSharedPtr rs_reader; |
| EXPECT_TRUE(rowset->create_reader(&rs_reader).ok()); |
| input_rs_readers.push_back(std::move(rs_reader)); |
| } |
| |
| // create output rowset writer |
| auto writer_context = create_rowset_writer_context(tablet_schema, NONOVERLAPPING, 3456, |
| {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(); |
| |
| // merge input rowset |
| TabletSharedPtr tablet = create_tablet(*tablet_schema, false); |
| 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)); |
| |
| // create output rowset reader |
| RowsetReaderContext reader_context; |
| reader_context.tablet_schema = tablet_schema; |
| reader_context.need_ordered_result = false; |
| std::vector<uint32_t> return_columns = {0, 1}; |
| reader_context.return_columns = &return_columns; |
| RowsetReaderSharedPtr output_rs_reader; |
| LOG(INFO) << "create rowset reader in test"; |
| create_and_init_rowset_reader(out_rowset.get(), reader_context, &output_rs_reader); |
| |
| // read output rowset data |
| Block output_block; |
| std::vector<std::tuple<int64_t, int64_t>> output_data; |
| do { |
| block_create(tablet_schema, &output_block); |
| s = output_rs_reader->next_batch(&output_block); |
| auto columns = output_block.get_columns_with_type_and_name(); |
| EXPECT_EQ(columns.size(), 2); |
| for (auto i = 0; i < output_block.rows(); i++) { |
| output_data.emplace_back(columns[0].column->get_int(i), columns[1].column->get_int(i)); |
| } |
| } while (s == Status::OK()); |
| EXPECT_EQ(Status::Error<END_OF_FILE>(""), s); |
| EXPECT_EQ(out_rowset->rowset_meta()->num_rows(), output_data.size()); |
| EXPECT_EQ(output_data.size(), num_segments * rows_per_segment); |
| // check vertical compaction result |
| for (auto id = 0; id < output_data.size(); id++) { |
| LOG(INFO) << "output data: " << std::get<0>(output_data[id]) << " " |
| << std::get<1>(output_data[id]); |
| } |
| int dst_id = 0; |
| for (auto s_id = 0; s_id < input_data[0].size(); s_id++) { |
| for (auto row_id = 0; row_id < input_data[0][s_id].size(); row_id++) { |
| LOG(INFO) << "input data: " << std::get<0>(input_data[0][s_id][row_id]) << " " |
| << std::get<1>(input_data[0][s_id][row_id]); |
| EXPECT_EQ(std::get<0>(input_data[0][s_id][row_id]), std::get<0>(output_data[dst_id])); |
| EXPECT_EQ(std::get<1>(input_data[0][s_id][row_id]) * 2, |
| std::get<1>(output_data[dst_id])); |
| dst_id++; |
| } |
| } |
| } |
| |
| // Test sparse compaction when a value group starts with a reserve-only column. |
| TEST_F(VerticalCompactionTest, TestUniqueKeyVerticalMergeWithNullableSparseColumn) { |
| const auto original_threshold = config::sparse_column_compaction_threshold_percent; |
| const auto original_columns_per_group = config::vertical_compaction_num_columns_per_group; |
| Defer restore_config {[original_threshold, original_columns_per_group]() { |
| config::sparse_column_compaction_threshold_percent = original_threshold; |
| config::vertical_compaction_num_columns_per_group = original_columns_per_group; |
| }}; |
| config::sparse_column_compaction_threshold_percent = 1.0; |
| config::vertical_compaction_num_columns_per_group = 2; |
| |
| auto num_input_rowset = 2; |
| auto num_segments = 1; |
| auto rows_per_segment = 100; |
| |
| // The first value column only reserves capacity, while the nullable BIGINT column |
| // pre-fills actual_rows slots for in-place replacement. |
| 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(5); |
| |
| ColumnPB* column_1 = tablet_schema_pb.add_column(); |
| column_1->set_unique_id(1); |
| column_1->set_name("k1"); |
| column_1->set_type("INT"); |
| column_1->set_is_key(true); |
| column_1->set_length(4); |
| column_1->set_index_length(4); |
| column_1->set_is_nullable(false); |
| column_1->set_is_bf_column(false); |
| |
| ColumnPB* column_2 = tablet_schema_pb.add_column(); |
| column_2->set_unique_id(2); |
| column_2->set_name("v0"); |
| column_2->set_type("BOOLEAN"); |
| column_2->set_length(1); |
| column_2->set_index_length(1); |
| column_2->set_is_key(false); |
| column_2->set_is_nullable(false); |
| column_2->set_is_bf_column(false); |
| |
| ColumnPB* column_3 = tablet_schema_pb.add_column(); |
| column_3->set_unique_id(3); |
| column_3->set_name("v1"); |
| column_3->set_type("BIGINT"); |
| column_3->set_length(8); |
| column_3->set_index_length(8); |
| column_3->set_is_key(false); |
| column_3->set_is_nullable(true); |
| column_3->set_is_bf_column(false); |
| |
| // DELETE_SIGN column required for unique keys |
| ColumnPB* column_4 = tablet_schema_pb.add_column(); |
| column_4->set_unique_id(4); |
| column_4->set_name(DELETE_SIGN); |
| column_4->set_type("TINYINT"); |
| column_4->set_length(1); |
| column_4->set_index_length(1); |
| column_4->set_is_nullable(false); |
| column_4->set_is_key(false); |
| column_4->set_is_bf_column(false); |
| |
| tablet_schema->init_from_pb(tablet_schema_pb); |
| |
| // Create input rowsets with mixed NULL values in v1. |
| std::vector<RowsetSharedPtr> input_rowsets; |
| for (auto i = 0; i < num_input_rowset; i++) { |
| RowsetWriterContext rowset_writer_context; |
| static int64_t inc_id = 2000; |
| RowsetId rowset_id; |
| rowset_id.init(inc_id++); |
| rowset_writer_context.rowset_id = rowset_id; |
| rowset_writer_context.tablet_id = 12345; |
| rowset_writer_context.tablet_schema_hash = 1111; |
| rowset_writer_context.partition_id = 10; |
| rowset_writer_context.rowset_type = BETA_ROWSET; |
| rowset_writer_context.tablet_path = absolute_dir + "/tablet_path"; |
| rowset_writer_context.rowset_state = VISIBLE; |
| rowset_writer_context.tablet_schema = tablet_schema; |
| rowset_writer_context.version = Version(i * 10, i * 10); |
| rowset_writer_context.segments_overlap = NONOVERLAPPING; |
| |
| auto res = RowsetFactory::create_rowset_writer(*engine_ref, rowset_writer_context, true); |
| ASSERT_TRUE(res.has_value()) << res.error(); |
| auto rowset_writer = std::move(res).value(); |
| |
| Block block = tablet_schema->create_block(); |
| auto columns = std::move(block).mutate_columns(); |
| |
| for (int rid = 0; rid < rows_per_segment; ++rid) { |
| int32_t k1 = i * rows_per_segment + rid; |
| columns[0]->insert_data((const char*)&k1, sizeof(k1)); |
| |
| uint8_t v0 = rid % 2; |
| columns[1]->insert_data((const char*)&v0, sizeof(v0)); |
| |
| // The first row is non-NULL, so the first sparse batch must execute replace. |
| if (rid % 10 == 0) { |
| int64_t v1 = static_cast<int64_t>(k1) * 10; |
| columns[2]->insert_data((const char*)&v1, sizeof(v1)); |
| } else { |
| columns[2]->insert_default(); |
| } |
| |
| uint8_t delete_sign = 0; |
| columns[3]->insert_data((const char*)&delete_sign, sizeof(delete_sign)); |
| } |
| |
| auto s = add_block_with_columns(rowset_writer.get(), &block, &columns); |
| ASSERT_TRUE(s.ok()) << s; |
| s = rowset_writer->flush(); |
| ASSERT_TRUE(s.ok()) << s; |
| |
| RowsetSharedPtr rowset; |
| ASSERT_EQ(Status::OK(), rowset_writer->build(rowset)); |
| input_rowsets.push_back(rowset); |
| } |
| |
| // Create input rowset readers |
| std::vector<RowsetReaderSharedPtr> input_rs_readers; |
| for (auto& rowset : input_rowsets) { |
| RowsetReaderSharedPtr rs_reader; |
| ASSERT_TRUE(rowset->create_reader(&rs_reader).ok()); |
| input_rs_readers.push_back(std::move(rs_reader)); |
| } |
| |
| // Create output rowset writer |
| auto writer_context = create_rowset_writer_context(tablet_schema, NONOVERLAPPING, 3456, |
| {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(); |
| |
| // Create tablet and run vertical merge |
| TabletSharedPtr tablet = create_tablet(*tablet_schema, false); |
| Merger::Statistics stats; |
| RowIdConversion rowid_conversion; |
| stats.rowid_conversion = &rowid_conversion; |
| |
| auto s = Merger::vertical_merge_rowsets(tablet, ReaderType::READER_BASE_COMPACTION, |
| *tablet_schema, input_rs_readers, |
| output_rs_writer.get(), 10000, num_segments, &stats); |
| ASSERT_TRUE(s.ok()) << s; |
| |
| RowsetSharedPtr out_rowset; |
| ASSERT_EQ(Status::OK(), output_rs_writer->build(out_rowset)); |
| |
| 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(out_rowset.get(), reader_context, &output_rs_reader); |
| |
| Block output_block; |
| size_t output_rows = 0; |
| do { |
| block_create(tablet_schema, &output_block); |
| s = output_rs_reader->next_batch(&output_block); |
| auto columns = output_block.get_columns_with_type_and_name(); |
| ASSERT_EQ(columns.size(), 3); |
| const auto& nullable_v1 = assert_cast<const ColumnNullable&>(*columns[2].column); |
| for (size_t row = 0; row < output_block.rows(); ++row) { |
| int64_t k1 = columns[0].column->get_int(row); |
| EXPECT_EQ(k1, output_rows); |
| EXPECT_EQ(columns[1].column->get_bool(row), k1 % 2 != 0); |
| if (k1 % 10 == 0) { |
| EXPECT_FALSE(nullable_v1.is_null_at(row)); |
| EXPECT_EQ(nullable_v1.get_nested_column().get_int(row), k1 * 10); |
| } else { |
| EXPECT_TRUE(nullable_v1.is_null_at(row)); |
| } |
| ++output_rows; |
| } |
| } while (s.ok()); |
| EXPECT_TRUE(s.is<END_OF_FILE>()) << s; |
| EXPECT_EQ(output_rows, num_input_rowset * rows_per_segment); |
| EXPECT_EQ(out_rowset->rowset_meta()->num_rows(), output_rows); |
| } |
| |
| // Test that first-time compaction (no historical sampling) uses footer raw_data_bytes |
| // to estimate batch_size instead of hardcoded 992. |
| // This test verifies the footer-based estimation path is triggered and compaction succeeds. |
| TEST_F(VerticalCompactionTest, TestFirstCompactionUsesFooterEstimation) { |
| // Use small data to ensure compaction completes quickly |
| auto num_input_rowset = 2; |
| auto num_segments = 1; |
| auto rows_per_segment = 1024; |
| SegmentsOverlapPB overlap = NONOVERLAPPING; |
| std::vector<std::vector<std::vector<std::tuple<int64_t, int64_t>>>> input_data; |
| generate_input_data(num_input_rowset, num_segments, rows_per_segment, overlap, input_data); |
| |
| TabletSchemaSPtr tablet_schema = create_schema(); |
| |
| // Create input rowsets |
| std::vector<RowsetSharedPtr> input_rowsets; |
| for (auto i = 0; i < num_input_rowset; i++) { |
| RowsetSharedPtr rowset = create_rowset(tablet_schema, overlap, input_data[i], i); |
| input_rowsets.push_back(rowset); |
| } |
| |
| // Create input rowset readers |
| std::vector<RowsetReaderSharedPtr> input_rs_readers; |
| for (auto& rowset : input_rowsets) { |
| RowsetReaderSharedPtr rs_reader; |
| ASSERT_TRUE(rowset->create_reader(&rs_reader).ok()); |
| input_rs_readers.push_back(std::move(rs_reader)); |
| } |
| |
| // Create output rowset writer |
| auto writer_context = create_rowset_writer_context(tablet_schema, NONOVERLAPPING, 3456, |
| {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(); |
| |
| // Create tablet - fresh tablet has no historical sampling data, |
| // so estimate_batch_size will hit the else branch and use footer raw_data_bytes. |
| TabletSharedPtr tablet = create_tablet(*tablet_schema, false); |
| Merger::Statistics stats; |
| RowIdConversion rowid_conversion; |
| stats.rowid_conversion = &rowid_conversion; |
| |
| // Verify sample_infos are empty (no historical data) |
| auto& sample_infos = tablet->get_sample_infos(ReaderType::READER_BASE_COMPACTION); |
| EXPECT_TRUE(sample_infos.empty()); |
| |
| // Run vertical merge - this should use footer raw_data_bytes for batch size estimation |
| // since there is no historical sampling data. |
| // The log should contain "estimate batch size from footer" instead of the old hardcoded path. |
| auto s = Merger::vertical_merge_rowsets(tablet, ReaderType::READER_BASE_COMPACTION, |
| *tablet_schema, input_rs_readers, |
| output_rs_writer.get(), 100000, num_segments, &stats); |
| ASSERT_TRUE(s.ok()) << s; |
| |
| RowsetSharedPtr out_rowset; |
| ASSERT_EQ(Status::OK(), output_rs_writer->build(out_rowset)); |
| EXPECT_EQ(out_rowset->rowset_meta()->num_rows(), |
| num_input_rowset * num_segments * rows_per_segment); |
| |
| // After first compaction, sample_infos should be populated with historical data |
| // for subsequent compactions to use. |
| auto& updated_infos = tablet->get_sample_infos(ReaderType::READER_BASE_COMPACTION); |
| EXPECT_FALSE(updated_infos.empty()); |
| } |
| |
| // Test that raw_data_bytes in segment footer accurately reflects the original data size |
| // for different column types, which is the foundation of footer-based batch size estimation. |
| TEST_F(VerticalCompactionTest, TestFooterRawDataBytesAccuracy) { |
| // Create a schema with INT key + VARCHAR value to test both fixed and variable-length types |
| TabletSchemaSPtr tablet_schema = std::make_shared<TabletSchema>(); |
| TabletSchemaPB tablet_schema_pb; |
| tablet_schema_pb.set_keys_type(DUP_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(3); |
| |
| ColumnPB* col_int = tablet_schema_pb.add_column(); |
| col_int->set_unique_id(1); |
| col_int->set_name("c_int"); |
| col_int->set_type("INT"); |
| col_int->set_is_key(true); |
| col_int->set_length(4); |
| col_int->set_index_length(4); |
| col_int->set_is_nullable(false); |
| col_int->set_is_bf_column(false); |
| |
| ColumnPB* col_varchar = tablet_schema_pb.add_column(); |
| col_varchar->set_unique_id(2); |
| col_varchar->set_name("c_varchar"); |
| col_varchar->set_type("VARCHAR"); |
| col_varchar->set_is_key(false); |
| col_varchar->set_length(128); |
| col_varchar->set_index_length(20); |
| col_varchar->set_is_nullable(false); |
| col_varchar->set_is_bf_column(false); |
| |
| tablet_schema->init_from_pb(tablet_schema_pb); |
| |
| // Write 1000 rows: INT values + VARCHAR strings of exactly 20 bytes each |
| constexpr int kNumRows = 1000; |
| constexpr int kStringLen = 20; |
| std::string fixed_string(kStringLen, 'x'); |
| |
| auto writer_context = |
| create_rowset_writer_context(tablet_schema, NONOVERLAPPING, UINT32_MAX, {0, 0}); |
| 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(); |
| |
| Block block = tablet_schema->create_block(); |
| auto columns = std::move(block).mutate_columns(); |
| for (int i = 0; i < kNumRows; i++) { |
| int32_t int_val = i; |
| columns[0]->insert_data(reinterpret_cast<const char*>(&int_val), sizeof(int_val)); |
| columns[1]->insert_data(fixed_string.data(), fixed_string.size()); |
| } |
| 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_EQ(1, rowset->rowset_meta()->num_segments()); |
| ASSERT_EQ(kNumRows, rowset->rowset_meta()->num_rows()); |
| |
| // Load segments and read footer's raw_data_bytes |
| auto beta_rowset = std::dynamic_pointer_cast<BetaRowset>(rowset); |
| ASSERT_NE(beta_rowset, nullptr); |
| std::vector<segment_v2::SegmentSharedPtr> segments; |
| ASSERT_TRUE(beta_rowset->load_segments(&segments).ok()); |
| ASSERT_EQ(1, segments.size()); |
| ASSERT_EQ(kNumRows, segments[0]->num_rows()); |
| |
| // Collect raw_data_bytes per column from footer |
| std::unordered_map<int32_t, uint64_t> raw_bytes_by_uid; |
| auto st = segments[0]->traverse_column_meta_pbs([&](const segment_v2::ColumnMetaPB& meta) { |
| if (meta.unique_id() >= 0 && meta.has_raw_data_bytes()) { |
| raw_bytes_by_uid[meta.unique_id()] = meta.raw_data_bytes(); |
| } |
| }); |
| ASSERT_TRUE(st.ok()) << st; |
| |
| // Verify INT column (uid=1): raw_data_bytes should be exactly kNumRows * sizeof(int32_t). |
| // PageBuilder::get_raw_data_size() accumulates raw data bytes added via add(), |
| // for fixed-width types this is exactly N * sizeof(T). |
| ASSERT_TRUE(raw_bytes_by_uid.count(1) > 0) << "INT column raw_data_bytes not found in footer"; |
| EXPECT_EQ(raw_bytes_by_uid[1], kNumRows * sizeof(int32_t)) |
| << "INT column: expected " << kNumRows * sizeof(int32_t) |
| << " total raw_data_bytes, got " << raw_bytes_by_uid[1]; |
| |
| // Verify VARCHAR column (uid=2): raw_data_bytes should be exactly kNumRows * kStringLen. |
| // BinaryPlainPageBuilder/BinaryDictPageBuilder only accumulate src->size (the raw string |
| // payload), not offsets, varint length prefixes, or dictionary overhead. |
| ASSERT_TRUE(raw_bytes_by_uid.count(2) > 0) |
| << "VARCHAR column raw_data_bytes not found in footer"; |
| EXPECT_EQ(raw_bytes_by_uid[2], kNumRows * kStringLen) |
| << "VARCHAR column: expected " << kNumRows * kStringLen << " total raw_data_bytes, got " |
| << raw_bytes_by_uid[2]; |
| } |
| |
| // Verify that raw_data_bytes only counts non-null payload for nullable |
| // fixed-width columns. This is the premise that motivates the type_size |
| // lower bound in merger.cpp's footer-based per-row estimation: without it, |
| // a sparse nullable column would produce a per-row estimate far below the |
| // reader's actual memory footprint (which still allocates the full nested |
| // slot for null rows via ColumnNullable::insert_many_defaults). |
| TEST_F(VerticalCompactionTest, TestFooterRawDataBytesNullableSparse) { |
| TabletSchemaSPtr tablet_schema = std::make_shared<TabletSchema>(); |
| TabletSchemaPB tablet_schema_pb; |
| tablet_schema_pb.set_keys_type(DUP_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(3); |
| |
| ColumnPB* col_key = tablet_schema_pb.add_column(); |
| col_key->set_unique_id(1); |
| col_key->set_name("c_key"); |
| col_key->set_type("INT"); |
| col_key->set_is_key(true); |
| col_key->set_length(4); |
| col_key->set_index_length(4); |
| col_key->set_is_nullable(false); |
| col_key->set_is_bf_column(false); |
| |
| ColumnPB* col_val = tablet_schema_pb.add_column(); |
| col_val->set_unique_id(2); |
| col_val->set_name("c_val"); |
| col_val->set_type("INT"); |
| col_val->set_is_key(false); |
| col_val->set_length(4); |
| col_val->set_index_length(4); |
| col_val->set_is_nullable(true); |
| col_val->set_is_bf_column(false); |
| |
| tablet_schema->init_from_pb(tablet_schema_pb); |
| |
| constexpr int kNumRows = 1000; |
| constexpr int kNonNullCount = 100; // 10% non-null, 90% null |
| |
| auto writer_context = |
| create_rowset_writer_context(tablet_schema, NONOVERLAPPING, UINT32_MAX, {0, 0}); |
| 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(); |
| |
| Block block = tablet_schema->create_block(); |
| auto columns = std::move(block).mutate_columns(); |
| for (int i = 0; i < kNumRows; i++) { |
| int32_t key_val = i; |
| columns[0]->insert_data(reinterpret_cast<const char*>(&key_val), sizeof(key_val)); |
| if (i < kNonNullCount) { |
| int32_t val = i; |
| columns[1]->insert_data(reinterpret_cast<const char*>(&val), sizeof(val)); |
| } else { |
| columns[1]->insert_default(); // ColumnNullable default is null |
| } |
| } |
| 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_EQ(1, rowset->rowset_meta()->num_segments()); |
| |
| auto beta_rowset = std::dynamic_pointer_cast<BetaRowset>(rowset); |
| ASSERT_NE(beta_rowset, nullptr); |
| std::vector<segment_v2::SegmentSharedPtr> segments; |
| ASSERT_TRUE(beta_rowset->load_segments(&segments).ok()); |
| ASSERT_EQ(1, segments.size()); |
| |
| std::unordered_map<int32_t, uint64_t> raw_bytes_by_uid; |
| auto st = segments[0]->traverse_column_meta_pbs([&](const segment_v2::ColumnMetaPB& meta) { |
| if (meta.unique_id() >= 0 && meta.has_raw_data_bytes()) { |
| raw_bytes_by_uid[meta.unique_id()] = meta.raw_data_bytes(); |
| } |
| }); |
| ASSERT_TRUE(st.ok()) << st; |
| |
| // Key column (non-null INT): full coverage, raw == kNumRows * 4. |
| ASSERT_TRUE(raw_bytes_by_uid.count(1) > 0); |
| EXPECT_EQ(raw_bytes_by_uid[1], kNumRows * sizeof(int32_t)); |
| |
| // Value column (nullable INT, 90% null): raw_data_bytes only reflects |
| // the non-null payload because ScalarColumnWriter::append_nulls() does |
| // not advance the page builder. So raw == kNonNullCount * 4. |
| // |
| // If merger.cpp used `raw / total_rows` directly the per-row estimate |
| // would be ~0.4 bytes (+1 for null map = 1.4), but the reader actually |
| // allocates 4 bytes for every nested slot (regardless of null-ness) |
| // plus 1 byte of null map = 5 bytes/row. The fixed-width type_size |
| // lower bound is what closes that ~3.5x gap. |
| ASSERT_TRUE(raw_bytes_by_uid.count(2) > 0); |
| EXPECT_EQ(raw_bytes_by_uid[2], kNonNullCount * sizeof(int32_t)); |
| EXPECT_LT(raw_bytes_by_uid[2], kNumRows * sizeof(int32_t)); |
| } |
| |
| } // namespace doris |