From 77830f9c927af62557262299b891b8e271fbcdcf Mon Sep 17 00:00:00 2001 From: ColinLee Date: Wed, 16 Sep 2026 09:38:46 +0800 Subject: [PATCH] perf(cpp): reduce tsfile close overhead --- cpp/src/common/tsfile_common.cc | 117 ++++++++++---------------- cpp/src/common/tsfile_common.h | 24 +++--- cpp/src/file/tsfile_io_writer.cc | 40 ++++----- cpp/src/file/tsfile_io_writer.h | 3 +- cpp/test/common/tsfile_common_test.cc | 85 ++++++++++++++++++- 5 files changed, 158 insertions(+), 111 deletions(-) diff --git a/cpp/src/common/tsfile_common.cc b/cpp/src/common/tsfile_common.cc index 3cac258ff..57bdfdaed 100644 --- a/cpp/src/common/tsfile_common.cc +++ b/cpp/src/common/tsfile_common.cc @@ -22,6 +22,7 @@ #include #include #include +#include #include "common/logger/elog.h" #include "common/schema.h" @@ -52,74 +53,56 @@ int TimeseriesIndex::add_chunk_meta(ChunkMeta* chunk_meta, /* ================ TSMIterator ================ */ int TSMIterator::init() { - // sort chunk_group_meta_list_ : {[measurementA, offsetA1], [measurementB, - // offsetB1], [measurementA, offsetA2], [measurementB, offsetB2]} -> - // {[measurementA, offsetA1], [measurementA, offsetA2], [measurementB, - // offsetB1], [measurementB, offsetB2]} + tsm_chunk_meta_info_.clear(); + const auto offset_less = [](const ChunkMeta* lhs, const ChunkMeta* rhs) { + return lhs->offset_of_chunk_header_ < rhs->offset_of_chunk_header_; + }; for (auto chunk_group_meta_iter = chunk_group_meta_list_.begin(); chunk_group_meta_iter != chunk_group_meta_list_.end(); chunk_group_meta_iter++) { - auto chunk_meta_list = chunk_group_meta_iter.get()->chunk_meta_list_; - // Use a map to group chunks by measurement_name_ - std::map> groups; - std::vector order; - for (auto it = chunk_meta_list.begin(); it != chunk_meta_list.end(); - it++) { - auto* chunk_meta = it.get(); - if (groups.find(chunk_meta->measurement_name_) == groups.end()) { - order.push_back(chunk_meta->measurement_name_); - } - groups[chunk_meta->measurement_name_].push_back(chunk_meta); - } - - // Sort each group of chunk metas by offset - for (auto it = groups.begin(); it != groups.end(); ++it) { - std::vector& group = it->second; - std::sort(group.begin(), group.end(), - [](ChunkMeta* a, ChunkMeta* b) { - return a->offset_of_chunk_header_ < - b->offset_of_chunk_header_; - }); + ChunkGroupMeta* chunk_group_meta = chunk_group_meta_iter.get(); + MeasurementChunkMetaMap chunk_group_map; + for (auto chunk_meta_iter = chunk_group_meta->chunk_meta_list_.begin(); + chunk_meta_iter != chunk_group_meta->chunk_meta_list_.end(); + chunk_meta_iter++) { + ChunkMeta* chunk_meta = chunk_meta_iter.get(); + chunk_group_map[chunk_meta->measurement_name_].push_back( + chunk_meta); } - // Clear and refill chunk_group_meta_list - chunk_group_meta_iter.get()->chunk_meta_list_.clear(); - for (const auto& measurement_name : order) { - for (auto chunk_meta : groups[measurement_name]) { - chunk_group_meta_iter.get()->chunk_meta_list_.push_back( - chunk_meta); + for (auto& measurement_entry : chunk_group_map) { + auto& chunk_metas = measurement_entry.second; + if (!std::is_sorted(chunk_metas.begin(), chunk_metas.end(), + offset_less)) { + std::sort(chunk_metas.begin(), chunk_metas.end(), offset_less); } } - } - // FIXME empty list - chunk_group_meta_iter_ = chunk_group_meta_list_.begin(); - while (chunk_group_meta_iter_ != chunk_group_meta_list_.end()) { - chunk_meta_iter_ = - chunk_group_meta_iter_.get()->chunk_meta_list_.begin(); - std::map> tmp; - while (chunk_meta_iter_ != - chunk_group_meta_iter_.get()->chunk_meta_list_.end()) { - tmp[chunk_meta_iter_.get()->measurement_name_].emplace_back( - chunk_meta_iter_.get()); - chunk_meta_iter_++; - } - if (!tmp.empty()) { - auto& merged = - tsm_chunk_meta_info_[chunk_group_meta_iter_.get()->device_id_]; - for (auto& m_entry : tmp) { - auto& vec = merged[m_entry.first]; - vec.insert(vec.end(), m_entry.second.begin(), - m_entry.second.end()); - } + auto device_pos = + tsm_chunk_meta_info_.lower_bound(chunk_group_meta->device_id_); + IDeviceIDComparator comparator; + const bool device_exists = + device_pos != tsm_chunk_meta_info_.end() && + !comparator(chunk_group_meta->device_id_, device_pos->first); + if (!device_exists) { + tsm_chunk_meta_info_.emplace_hint(device_pos, + chunk_group_meta->device_id_, + std::move(chunk_group_map)); + continue; } - chunk_group_meta_iter_++; - } - if (!tsm_chunk_meta_info_.empty() && - !tsm_chunk_meta_info_.begin()->second.empty()) { - tsm_measurement_iter_ = tsm_chunk_meta_info_.begin()->second.begin(); + auto& measurement_map = device_pos->second; + for (auto& measurement_entry : chunk_group_map) { + auto& chunk_metas = measurement_entry.second; + auto& merged_chunk_metas = measurement_map[measurement_entry.first]; + merged_chunk_metas.insert(merged_chunk_metas.end(), + chunk_metas.begin(), chunk_metas.end()); + } } + tsm_device_iter_ = tsm_chunk_meta_info_.begin(); + if (tsm_device_iter_ != tsm_chunk_meta_info_.end()) { + tsm_measurement_iter_ = tsm_device_iter_->second.begin(); + } return E_OK; } @@ -131,8 +114,6 @@ int TSMIterator::get_next(std::shared_ptr& ret_device_id, String& ret_measurement_name, TimeseriesIndex& ret_ts_index) { int ret = E_OK; - SimpleList chunk_meta_list_of_this_ts( - 1024, MOD_TIMESERIES_INDEX_OBJ); // FIXME if (tsm_measurement_iter_ == tsm_device_iter_->second.end()) { tsm_device_iter_++; if (!has_next()) { @@ -143,15 +124,13 @@ int TSMIterator::get_next(std::shared_ptr& ret_device_id, } ret_device_id = tsm_device_iter_->first; ret_measurement_name.shallow_copy_from(tsm_measurement_iter_->first); - for (auto meta : tsm_measurement_iter_->second) { - chunk_meta_list_of_this_ts.push_back(meta); - } - if (chunk_meta_list_of_this_ts.size() == 0) { + const std::vector& chunk_metas = tsm_measurement_iter_->second; + if (chunk_metas.empty()) { return E_TSFILE_WRITER_META_ERR; } - const bool multi_chunks = chunk_meta_list_of_this_ts.size() > 1; - ChunkMeta* first_chunk_meta = chunk_meta_list_of_this_ts.front(); + const bool multi_chunks = chunk_metas.size() > 1; + ChunkMeta* first_chunk_meta = chunk_metas.front(); const char meta_type = (multi_chunks ? 1 : 0) | (first_chunk_meta->mask_); const TSDataType data_type = first_chunk_meta->data_type_; @@ -160,13 +139,9 @@ int TSMIterator::get_next(std::shared_ptr& ret_device_id, ret_ts_index.set_data_type(data_type); ret_ts_index.init_statistic(data_type); - SimpleList::Iterator ts_chunk_meta_iter = - chunk_meta_list_of_this_ts.begin(); - for (; - IS_SUCC(ret) && ts_chunk_meta_iter != chunk_meta_list_of_this_ts.end(); - ts_chunk_meta_iter++) { - ChunkMeta* chunk_meta = ts_chunk_meta_iter.get(); + for (ChunkMeta* chunk_meta : chunk_metas) { if (RET_FAIL(ret_ts_index.add_chunk_meta(chunk_meta, multi_chunks))) { + break; } } if (IS_SUCC(ret)) { diff --git a/cpp/src/common/tsfile_common.h b/cpp/src/common/tsfile_common.h index ee3b389d5..f604cb9a4 100644 --- a/cpp/src/common/tsfile_common.h +++ b/cpp/src/common/tsfile_common.h @@ -695,9 +695,7 @@ class TSMIterator { public: explicit TSMIterator( common::SimpleList& chunk_group_meta_list) - : chunk_group_meta_list_(chunk_group_meta_list), - chunk_group_meta_iter_(), - chunk_meta_iter_() {} + : chunk_group_meta_list_(chunk_group_meta_list) {} // sort => iterate int init(); @@ -707,24 +705,22 @@ class TSMIterator { TimeseriesIndex& ret_ts_index); private: + using MeasurementChunkMetaMap = + std::map>; + using DeviceChunkMetaMap = + std::map, MeasurementChunkMetaMap, + IDeviceIDComparator>; + common::SimpleList& chunk_group_meta_list_; - common::SimpleList::Iterator chunk_group_meta_iter_; - common::SimpleList::Iterator chunk_meta_iter_; // timeseries measurenemnt chunk meta info - std::map, - std::map>, - IDeviceIDComparator> - tsm_chunk_meta_info_; + DeviceChunkMetaMap tsm_chunk_meta_info_; // device iterator - std::map, - std::map>, - IDeviceIDComparator>::iterator tsm_device_iter_; + DeviceChunkMetaMap::iterator tsm_device_iter_; // measurement iterator - std::map>::iterator - tsm_measurement_iter_; + MeasurementChunkMetaMap::iterator tsm_measurement_iter_; }; /* =============== TsFile Index ================ */ diff --git a/cpp/src/file/tsfile_io_writer.cc b/cpp/src/file/tsfile_io_writer.cc index 29ddf0d90..5b34d8e10 100644 --- a/cpp/src/file/tsfile_io_writer.cc +++ b/cpp/src/file/tsfile_io_writer.cc @@ -35,6 +35,8 @@ using namespace common; namespace storage { +const uint32_t TsFileIOWriter::WRITE_STREAM_PAGE_SIZE; + #if 0 #define OFFSET_DEBUG(msg) \ std::cout << "OFFSET_DEBUG: " << msg \ @@ -379,7 +381,7 @@ int TsFileIOWriter::write_file_index() { std::shared_ptr meta_index_entry = nullptr; std::shared_ptr cur_index_node = nullptr; SimpleList>* cur_index_node_queue = nullptr; - DeviceNodeMap device_map; + std::map table_device_nodes_map; TSMIterator tsm_iter(chunk_group_meta_list_); @@ -410,9 +412,11 @@ int TsFileIOWriter::write_file_index() { if (prev_device_id != nullptr) { if (RET_FAIL(add_cur_index_node_to_queue( cur_index_node, cur_index_node_queue))) { - } else if (RET_FAIL(add_device_node(device_map, prev_device_id, - cur_index_node_queue, - writing_mm))) { + } else if (RET_FAIL(add_device_node( + table_device_nodes_map[prev_device_id + ->get_table_name()], + prev_device_id, cur_index_node_queue, + writing_mm))) { } } if (IS_SUCC(ret)) { @@ -477,9 +481,9 @@ int TsFileIOWriter::write_file_index() { ASSERT(cur_index_node_queue != nullptr); if (RET_FAIL(add_cur_index_node_to_queue(cur_index_node, cur_index_node_queue))) { - } else if (RET_FAIL(add_device_node(device_map, prev_device_id, - cur_index_node_queue, - writing_mm))) { + } else if (RET_FAIL(add_device_node( + table_device_nodes_map[prev_device_id->get_table_name()], + prev_device_id, cur_index_node_queue, writing_mm))) { } } @@ -487,16 +491,6 @@ int TsFileIOWriter::write_file_index() { TsFileMeta tsfile_meta; tsfile_meta.meta_offset_ = meta_offset; tsfile_meta.bloom_filter_ = &filter; - // split device by table - std::map table_device_nodes_map; - for (const auto& entry : device_map) { - std::string table_name = entry.first->get_table_name(); - auto& table_map = table_device_nodes_map[table_name]; - if (table_map.empty() || - table_map.find(entry.first) == table_map.end()) { - table_map[entry.first] = entry.second; - } - } std::map> table_nodes_map; for (auto& entry : table_device_nodes_map) { auto meta_index_node = @@ -730,8 +724,10 @@ int TsFileIOWriter::add_device_node( FileIndexWritingMemManager& wmm) { ASSERT(measurement_index_node_queue->size() > 0); int ret = E_OK; - auto find_iter = device_map.find(device_id); - if (find_iter != device_map.end()) { + auto insert_pos = device_map.lower_bound(device_id); + IDeviceIDComparator comparator; + if (insert_pos != device_map.end() && + !comparator(device_id, insert_pos->first)) { return E_ALREADY_EXIST; } @@ -739,11 +735,7 @@ int TsFileIOWriter::add_device_node( if (RET_FAIL(generate_root(measurement_index_node_queue, root, INTERNAL_MEASUREMENT, wmm))) { } else { - std::pair ins_res = - device_map.insert(std::make_pair(device_id, root)); - if (!ins_res.second) { - ASSERT(false); - } + device_map.emplace_hint(insert_pos, device_id, root); } return ret; } diff --git a/cpp/src/file/tsfile_io_writer.h b/cpp/src/file/tsfile_io_writer.h index bbb1e4988..d7c656015 100644 --- a/cpp/src/file/tsfile_io_writer.h +++ b/cpp/src/file/tsfile_io_writer.h @@ -56,7 +56,8 @@ class TsFileIOWriter { typedef DeviceNodeMap::iterator DeviceNodeMapIterator; public: - static const uint32_t WRITE_STREAM_PAGE_SIZE = 512; // FIXME + static const uint32_t WRITE_STREAM_PAGE_SIZE = 16 * 1024; + public: TsFileIOWriter() : meta_allocator_(), diff --git a/cpp/test/common/tsfile_common_test.cc b/cpp/test/common/tsfile_common_test.cc index 309fdbc39..67bfa9231 100644 --- a/cpp/test/common/tsfile_common_test.cc +++ b/cpp/test/common/tsfile_common_test.cc @@ -21,8 +21,12 @@ #include #include +#include +#include + #include "common/global.h" #include "compress/compressor_factory.h" +#include "file/tsfile_io_writer.h" namespace storage { TEST(PageHeaderTest, DefaultConstructor) { @@ -177,11 +181,14 @@ class TSMIteratorTest : public ::testing::Test { char measure_name[] = "measurement_1"; common::String measurement_name(measure_name, sizeof(measure_name)); stat_ = StatisticFactory::alloc_statistic(common::TSDataType::INT32); + statistics_.push_back(stat_); + stat_->update(100, static_cast(100)); chunk_meta->init(measurement_name, common::TSDataType::INT32, 100, stat_, 1, common::PLAIN, common::UNCOMPRESSED, arena); chunk_group_meta->chunk_meta_list_.push_back(chunk_meta); chunk_group_meta_list_->push_back(chunk_group_meta); + chunk_group_meta_ = chunk_group_meta; } void TearDown() override { @@ -190,14 +197,42 @@ class TSMIteratorTest : public ::testing::Test { iter.get()->device_id_.reset(); } delete chunk_group_meta_list_; - StatisticFactory::free(stat_); + for (Statistic* statistic : statistics_) { + StatisticFactory::free(statistic); + } + } + + ChunkMeta* append_chunk(ChunkGroupMeta* chunk_group_meta, + const char* measurement, int64_t offset) { + void* buf = arena.alloc(sizeof(ChunkMeta)); + auto chunk_meta = new (buf) ChunkMeta(); + common::String measurement_name( + const_cast(measurement), + static_cast(std::strlen(measurement) + 1)); + Statistic* statistic = + StatisticFactory::alloc_statistic(common::TSDataType::INT32); + statistics_.push_back(statistic); + statistic->update(offset, static_cast(offset)); + EXPECT_EQ(chunk_meta->init(measurement_name, common::TSDataType::INT32, + offset, statistic, 1, common::PLAIN, + common::UNCOMPRESSED, arena), + common::E_OK); + EXPECT_EQ(chunk_group_meta->chunk_meta_list_.push_back(chunk_meta), + common::E_OK); + return chunk_meta; } common::PageArena arena; Statistic* stat_; + std::vector statistics_; + ChunkGroupMeta* chunk_group_meta_; common::SimpleList* chunk_group_meta_list_; }; +TEST(TsFileIOWriterTest, WriteStreamUses16KiBPages) { + EXPECT_EQ(TsFileIOWriter::WRITE_STREAM_PAGE_SIZE, 16U * 1024U); +} + TEST_F(TSMIteratorTest, InitSuccess) { TSMIterator iter(*chunk_group_meta_list_); ASSERT_EQ(iter.init(), common::E_OK); @@ -239,6 +274,54 @@ TEST_F(TSMIteratorTest, GetNext) { common::E_NO_MORE_DATA); } +TEST_F(TSMIteratorTest, InitDoesNotReorderSourceChunkMetadata) { + ChunkMeta* first = chunk_group_meta_->chunk_meta_list_.front(); + ChunkMeta* second = append_chunk(chunk_group_meta_, "measurement_2", 75); + ChunkMeta* third = append_chunk(chunk_group_meta_, "measurement_1", 50); + + TSMIterator iter(*chunk_group_meta_list_); + ASSERT_EQ(iter.init(), common::E_OK); + + auto source_iter = chunk_group_meta_->chunk_meta_list_.begin(); + ASSERT_NE(source_iter, chunk_group_meta_->chunk_meta_list_.end()); + EXPECT_EQ(source_iter.get(), first); + source_iter++; + ASSERT_NE(source_iter, chunk_group_meta_->chunk_meta_list_.end()); + EXPECT_EQ(source_iter.get(), second); + source_iter++; + ASSERT_NE(source_iter, chunk_group_meta_->chunk_meta_list_.end()); + EXPECT_EQ(source_iter.get(), third); +} + +TEST_F(TSMIteratorTest, SortsChunkOffsetsWithinChunkGroup) { + append_chunk(chunk_group_meta_, "measurement_1", 50); + + TSMIterator iter(*chunk_group_meta_list_); + ASSERT_EQ(iter.init(), common::E_OK); + + std::shared_ptr device_id; + common::String measurement_name; + TimeseriesIndex timeseries_index; + ASSERT_EQ(iter.get_next(device_id, measurement_name, timeseries_index), + common::E_OK); + + common::ByteStream serialized(1024, common::MOD_DEFAULT); + ASSERT_EQ(timeseries_index.serialize_to(serialized), common::E_OK); + common::PageArena deserialize_arena; + deserialize_arena.init(1024, common::MOD_DEFAULT); + TimeseriesIndex deserialized; + ASSERT_EQ(deserialized.deserialize_from(serialized, &deserialize_arena), + common::E_OK); + + auto* chunk_meta_list = deserialized.get_chunk_meta_list(); + ASSERT_NE(chunk_meta_list, nullptr); + ASSERT_EQ(chunk_meta_list->size(), 2U); + auto chunk_iter = chunk_meta_list->begin(); + EXPECT_EQ(chunk_iter.get()->offset_of_chunk_header_, 50); + chunk_iter++; + EXPECT_EQ(chunk_iter.get()->offset_of_chunk_header_, 100); +} + class MetaIndexEntryTest : public ::testing::Test { protected: common::PageArena pa_;