From b8578fdf059462f6ff5507f98758dce2556e6838 Mon Sep 17 00:00:00 2001 From: yy <736262857@qq.com> Date: Thu, 9 Jul 2026 10:09:49 +0000 Subject: [PATCH 1/7] storage/v3: merge continuous field reads into single system call (#5078) Signed-off-by: yy <736262857@qq.com> --- dbms/src/Storages/Page/V3/BlobStore.cpp | 73 ++++++++++++++++--------- 1 file changed, 48 insertions(+), 25 deletions(-) diff --git a/dbms/src/Storages/Page/V3/BlobStore.cpp b/dbms/src/Storages/Page/V3/BlobStore.cpp index 3ab513d21f9..348366a1e39 100644 --- a/dbms/src/Storages/Page/V3/BlobStore.cpp +++ b/dbms/src/Storages/Page/V3/BlobStore.cpp @@ -910,7 +910,6 @@ typename BlobStore::PageMap BlobStore::read(FieldReadInfos & to_re return page_map; } - // Allocate one for holding all pages data char * shared_data_buf = static_cast(alloc(buf_size)); MemHolder shared_mem_holder = createMemHolder(shared_data_buf, [&, buf_size](char * p) { free(p, buf_size); }); @@ -921,40 +920,63 @@ typename BlobStore::PageMap BlobStore::read(FieldReadInfos & to_re { size_t read_size_this_entry = 0; char * write_offset = pos; - for (const auto field_index : fields) + size_t field_idx = 0; + while (field_idx < fields.size()) { - // TODO: Continuously fields can read by one system call. - const auto [beg_offset, end_offset] = entry.getFieldOffsets(field_index); + size_t start_field_index = fields[field_idx]; + const auto [beg_offset, _1] = entry.getFieldOffsets(start_field_index); + + size_t end_field_index = start_field_index; + while (field_idx + 1 < fields.size() && fields[field_idx + 1] == fields[field_idx] + 1) + { + ++field_idx; + end_field_index = fields[field_idx]; + } + + const auto [_2, end_offset] = entry.getFieldOffsets(end_field_index); const auto size_to_read = end_offset - beg_offset; + read(page_id_v3, entry.file_id, entry.offset + beg_offset, write_offset, size_to_read, read_limiter); - fields_offset_in_page.emplace(field_index, read_size_this_entry); - if constexpr (BLOBSTORE_CHECKSUM_ON_READ) + size_t field_offset_in_buffer = 0; + for (size_t i = start_field_index; i <= end_field_index; ++i) { - const auto expect_checksum = entry.field_offsets[field_index].second; - ChecksumClass digest; - digest.update(write_offset, size_to_read); - auto field_checksum = digest.checksum(); - if (unlikely(entry.size != 0 && field_checksum != expect_checksum)) + const auto [field_beg, field_end] = entry.getFieldOffsets(i); + const auto field_size = field_end - field_beg; + const auto field_offset_in_entry = field_beg - beg_offset; + + fields_offset_in_page.emplace(i, read_size_this_entry); + + if constexpr (BLOBSTORE_CHECKSUM_ON_READ) { - throw Exception( - ErrorCodes::CHECKSUM_DOESNT_MATCH, - "Reading with fields meet checksum not match " - "page_id={} expected=0x{:X} actual=0x{:X} " - "field_index={} field_offset={} field_size={} " - "entry={}", - page_id_v3, - expect_checksum, - field_checksum, - field_index, - beg_offset, - size_to_read, - entry); + const auto expect_checksum = entry.field_offsets[i].second; + ChecksumClass digest; + digest.update(write_offset + field_offset_in_entry, field_size); + auto field_checksum = digest.checksum(); + if (unlikely(entry.size != 0 && field_checksum != expect_checksum)) + { + throw Exception( + ErrorCodes::CHECKSUM_DOESNT_MATCH, + "Reading with fields meet checksum not match " + "page_id={} expected=0x{:X} actual=0x{:X} " + "field_index={} field_offset={} field_size={} " + "entry={}", + page_id_v3, + expect_checksum, + field_checksum, + i, + field_beg, + field_size, + entry); + } } + + read_size_this_entry += field_size; + field_offset_in_buffer += field_size; } - read_size_this_entry += size_to_read; write_offset += size_to_read; + ++field_idx; } Page page(Trait::PageIdTrait::getU64ID(page_id_v3)); @@ -985,6 +1007,7 @@ typename BlobStore::PageMap BlobStore::read(FieldReadInfos & to_re pos, buf.toString()); } + return page_map; } From da8fc45b51cf1ef1027648eec0dd7113c574fb01 Mon Sep 17 00:00:00 2001 From: yy <736262857@qq.com> Date: Thu, 9 Jul 2026 12:15:48 +0000 Subject: [PATCH 2/7] storage/v3: add test ReadByFieldReadInfosContinuousFields to verify field merge optimization (#5078) Signed-off-by: yy <736262857@qq.com> --- dbms/src/Storages/Page/V3/BlobStore.cpp | 7 +++ .../Page/V3/tests/gtest_blob_store.cpp | 58 +++++++++++++++++++ 2 files changed, 65 insertions(+) diff --git a/dbms/src/Storages/Page/V3/BlobStore.cpp b/dbms/src/Storages/Page/V3/BlobStore.cpp index 348366a1e39..d25ac7c9cd9 100644 --- a/dbms/src/Storages/Page/V3/BlobStore.cpp +++ b/dbms/src/Storages/Page/V3/BlobStore.cpp @@ -847,6 +847,7 @@ void BlobStore::removePosFromStats(BlobFileId blob_id, BlobFileOffset off template typename BlobStore::PageMap BlobStore::read(FieldReadInfos & to_read, const ReadLimiterPtr & read_limiter) { + if (to_read.empty()) { return {}; @@ -861,6 +862,7 @@ typename BlobStore::PageMap BlobStore::read(FieldReadInfos & to_re // allocate data_buf that can hold all pages with specify fields + size_t buf_size = 0; for (auto & [page_id, entry, fields] : to_read) { @@ -938,6 +940,9 @@ typename BlobStore::PageMap BlobStore::read(FieldReadInfos & to_re read(page_id_v3, entry.file_id, entry.offset + beg_offset, write_offset, size_to_read, read_limiter); + + + size_t field_offset_in_buffer = 0; for (size_t i = start_field_index; i <= end_field_index; ++i) { @@ -1008,6 +1013,8 @@ typename BlobStore::PageMap BlobStore::read(FieldReadInfos & to_re buf.toString()); } + + return page_map; } diff --git a/dbms/src/Storages/Page/V3/tests/gtest_blob_store.cpp b/dbms/src/Storages/Page/V3/tests/gtest_blob_store.cpp index ceb3449f775..a5eecaf9e83 100644 --- a/dbms/src/Storages/Page/V3/tests/gtest_blob_store.cpp +++ b/dbms/src/Storages/Page/V3/tests/gtest_blob_store.cpp @@ -1923,4 +1923,62 @@ try } } CATCH + +TEST_F(BlobStoreTest, ReadByFieldReadInfosContinuousFields) +try +{ + const auto file_provider = DB::tests::TiFlashTestEnv::getDefaultFileProvider(); + PageIdU64 page_id = 50; + size_t buff_size = 20; + + auto blob_store = BlobStore( + getCurrentTestName(), + file_provider, + delegator, + BlobConfig{}, + page_type_and_config); + + std::vector c_buff(buff_size); + for (size_t j = 0; j < buff_size; ++j) + { + c_buff[j] = static_cast(j); + } + + ReadBufferPtr buff = std::make_shared(c_buff.data(), buff_size); + PageFieldSizes field_sizes{1, 2, 4, 8, (buff_size - 1 - 2 - 4 - 8)}; + WriteBatch wb; + wb.putPage(page_id, /* tag */ 0, buff, buff_size, field_sizes); + PageEntriesEdit edit = blob_store.write(std::move(wb)); + const auto & records = edit.getRecords(); + ASSERT_EQ(records.size(), 1); + auto entry = records[0].entry; + + std::cout << "=== Test 1: Continuous fields {0,1,2,3,4} ===" << std::endl; + { + BlobStore::FieldReadInfos read_infos; + read_infos.emplace_back(BlobStore::FieldReadInfo(buildV3Id(TEST_NAMESPACE_ID, page_id), entry, {0, 1, 2, 3, 4})); + auto page_map = blob_store.read(read_infos); + Page page = page_map.at(page_id); + ASSERT_EQ(page.fieldSize(), 5); + } + + std::cout << "=== Test 2: Fields with hole {0,1,3,4} ===" << std::endl; + { + BlobStore::FieldReadInfos read_infos; + read_infos.emplace_back(BlobStore::FieldReadInfo(buildV3Id(TEST_NAMESPACE_ID, page_id), entry, {0, 1, 3, 4})); + auto page_map = blob_store.read(read_infos); + Page page = page_map.at(page_id); + ASSERT_EQ(page.fieldSize(), 4); + } + + std::cout << "=== Test 3: Scattered fields {0,2,4} ===" << std::endl; + { + BlobStore::FieldReadInfos read_infos; + read_infos.emplace_back(BlobStore::FieldReadInfo(buildV3Id(TEST_NAMESPACE_ID, page_id), entry, {0, 2, 4})); + auto page_map = blob_store.read(read_infos); + Page page = page_map.at(page_id); + ASSERT_EQ(page.fieldSize(), 3); + } +} +CATCH } // namespace DB::PS::V3::tests From f1d7d137a9530db40c497739e8d1d21b9721a5b1 Mon Sep 17 00:00:00 2001 From: yy <736262857@qq.com> Date: Thu, 9 Jul 2026 13:42:14 +0000 Subject: [PATCH 3/7] storage/v3: improve code style for field offset extraction (#5078) Signed-off-by: yy <736262857@qq.com> --- dbms/src/Storages/Page/V3/BlobStore.cpp | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/dbms/src/Storages/Page/V3/BlobStore.cpp b/dbms/src/Storages/Page/V3/BlobStore.cpp index d25ac7c9cd9..fe640942c9a 100644 --- a/dbms/src/Storages/Page/V3/BlobStore.cpp +++ b/dbms/src/Storages/Page/V3/BlobStore.cpp @@ -926,7 +926,7 @@ typename BlobStore::PageMap BlobStore::read(FieldReadInfos & to_re while (field_idx < fields.size()) { size_t start_field_index = fields[field_idx]; - const auto [beg_offset, _1] = entry.getFieldOffsets(start_field_index); + const auto beg_offset = entry.getFieldOffsets(start_field_index).first; size_t end_field_index = start_field_index; while (field_idx + 1 < fields.size() && fields[field_idx + 1] == fields[field_idx] + 1) @@ -935,7 +935,7 @@ typename BlobStore::PageMap BlobStore::read(FieldReadInfos & to_re end_field_index = fields[field_idx]; } - const auto [_2, end_offset] = entry.getFieldOffsets(end_field_index); + const auto end_offset = entry.getFieldOffsets(end_field_index).second; const auto size_to_read = end_offset - beg_offset; read(page_id_v3, entry.file_id, entry.offset + beg_offset, write_offset, size_to_read, read_limiter); From 4d57c44f0d115624ae407512a12b50e7bbb34f11 Mon Sep 17 00:00:00 2001 From: yy <736262857@qq.com> Date: Fri, 10 Jul 2026 03:30:10 +0000 Subject: [PATCH 4/7] storage/v3: add field_read_call_count to verify read merge optimization (#5078) Signed-off-by: yy <736262857@qq.com> --- dbms/src/Storages/Page/V3/BlobStore.cpp | 6 +++--- dbms/src/Storages/Page/V3/BlobStore.h | 11 +++++++++++ dbms/src/Storages/Page/V3/tests/gtest_blob_store.cpp | 12 +++++++++--- 3 files changed, 23 insertions(+), 6 deletions(-) diff --git a/dbms/src/Storages/Page/V3/BlobStore.cpp b/dbms/src/Storages/Page/V3/BlobStore.cpp index fe640942c9a..7a69629fada 100644 --- a/dbms/src/Storages/Page/V3/BlobStore.cpp +++ b/dbms/src/Storages/Page/V3/BlobStore.cpp @@ -939,9 +939,9 @@ typename BlobStore::PageMap BlobStore::read(FieldReadInfos & to_re const auto size_to_read = end_offset - beg_offset; read(page_id_v3, entry.file_id, entry.offset + beg_offset, write_offset, size_to_read, read_limiter); - - - +#ifdef DBMS_PUBLIC_GTEST + ++field_read_call_count; +#endif size_t field_offset_in_buffer = 0; for (size_t i = start_field_index; i <= end_field_index; ++i) diff --git a/dbms/src/Storages/Page/V3/BlobStore.h b/dbms/src/Storages/Page/V3/BlobStore.h index 05c7f5fa299..cf9d07d7651 100644 --- a/dbms/src/Storages/Page/V3/BlobStore.h +++ b/dbms/src/Storages/Page/V3/BlobStore.h @@ -117,6 +117,12 @@ class BlobStore : private Allocator using FieldReadInfos = std::vector; PageMap read(FieldReadInfos & to_read, const ReadLimiterPtr & read_limiter = nullptr); +#ifdef DBMS_PUBLIC_GTEST + void resetFieldReadCallCount() { field_read_call_count = 0; } + size_t getFieldReadCallCount() { return field_read_call_count; } + +#endif + #ifndef DBMS_PUBLIC_GTEST private: #endif @@ -173,6 +179,11 @@ class BlobStore : private Allocator std::mutex mtx_blob_files; std::unordered_map blob_files; + +#ifdef DBMS_PUBLIC_GTEST + size_t field_read_call_count; +#endif + }; namespace u128 { diff --git a/dbms/src/Storages/Page/V3/tests/gtest_blob_store.cpp b/dbms/src/Storages/Page/V3/tests/gtest_blob_store.cpp index a5eecaf9e83..3b42cdb7128 100644 --- a/dbms/src/Storages/Page/V3/tests/gtest_blob_store.cpp +++ b/dbms/src/Storages/Page/V3/tests/gtest_blob_store.cpp @@ -1953,31 +1953,37 @@ try ASSERT_EQ(records.size(), 1); auto entry = records[0].entry; - std::cout << "=== Test 1: Continuous fields {0,1,2,3,4} ===" << std::endl; + // Test 1: Continuous fields {0,1,2,3,4} { + blob_store.resetFieldReadCallCount(); BlobStore::FieldReadInfos read_infos; read_infos.emplace_back(BlobStore::FieldReadInfo(buildV3Id(TEST_NAMESPACE_ID, page_id), entry, {0, 1, 2, 3, 4})); auto page_map = blob_store.read(read_infos); Page page = page_map.at(page_id); ASSERT_EQ(page.fieldSize(), 5); + ASSERT_EQ(blob_store.getFieldReadCallCount(), 1); } - std::cout << "=== Test 2: Fields with hole {0,1,3,4} ===" << std::endl; + // Test 2: Fields with hole {0,1,3,4} { + blob_store.resetFieldReadCallCount(); BlobStore::FieldReadInfos read_infos; read_infos.emplace_back(BlobStore::FieldReadInfo(buildV3Id(TEST_NAMESPACE_ID, page_id), entry, {0, 1, 3, 4})); auto page_map = blob_store.read(read_infos); Page page = page_map.at(page_id); ASSERT_EQ(page.fieldSize(), 4); + ASSERT_EQ(blob_store.getFieldReadCallCount(), 2); } - std::cout << "=== Test 3: Scattered fields {0,2,4} ===" << std::endl; + // Test 3: Scattered fields {0,2,4} { + blob_store.resetFieldReadCallCount(); BlobStore::FieldReadInfos read_infos; read_infos.emplace_back(BlobStore::FieldReadInfo(buildV3Id(TEST_NAMESPACE_ID, page_id), entry, {0, 2, 4})); auto page_map = blob_store.read(read_infos); Page page = page_map.at(page_id); ASSERT_EQ(page.fieldSize(), 3); + ASSERT_EQ(blob_store.getFieldReadCallCount(), 3); } } CATCH From 143a94facdbce0fb283be64e5db47d77330953ac Mon Sep 17 00:00:00 2001 From: yy <736262857@qq.com> Date: Fri, 10 Jul 2026 04:02:40 +0000 Subject: [PATCH 5/7] Add default initialization to field_read_call_count and remove unused field_offset_in_buffer Signed-off-by: yy <736262857@qq.com> --- dbms/src/Storages/Page/V3/BlobStore.cpp | 3 --- dbms/src/Storages/Page/V3/BlobStore.h | 2 +- 2 files changed, 1 insertion(+), 4 deletions(-) diff --git a/dbms/src/Storages/Page/V3/BlobStore.cpp b/dbms/src/Storages/Page/V3/BlobStore.cpp index 7a69629fada..f7f79c716db 100644 --- a/dbms/src/Storages/Page/V3/BlobStore.cpp +++ b/dbms/src/Storages/Page/V3/BlobStore.cpp @@ -942,8 +942,6 @@ typename BlobStore::PageMap BlobStore::read(FieldReadInfos & to_re #ifdef DBMS_PUBLIC_GTEST ++field_read_call_count; #endif - - size_t field_offset_in_buffer = 0; for (size_t i = start_field_index; i <= end_field_index; ++i) { const auto [field_beg, field_end] = entry.getFieldOffsets(i); @@ -977,7 +975,6 @@ typename BlobStore::PageMap BlobStore::read(FieldReadInfos & to_re } read_size_this_entry += field_size; - field_offset_in_buffer += field_size; } write_offset += size_to_read; diff --git a/dbms/src/Storages/Page/V3/BlobStore.h b/dbms/src/Storages/Page/V3/BlobStore.h index cf9d07d7651..cbaa4464ede 100644 --- a/dbms/src/Storages/Page/V3/BlobStore.h +++ b/dbms/src/Storages/Page/V3/BlobStore.h @@ -181,7 +181,7 @@ class BlobStore : private Allocator std::unordered_map blob_files; #ifdef DBMS_PUBLIC_GTEST - size_t field_read_call_count; + size_t field_read_call_count = 0; #endif }; From b5b9472a74298a06a495022aa30b6b5d971d5291 Mon Sep 17 00:00:00 2001 From: yy <736262857@qq.com> Date: Fri, 10 Jul 2026 15:02:12 +0000 Subject: [PATCH 6/7] fix: format code with clang-format Signed-off-by: yy <736262857@qq.com> --- dbms/src/Storages/Page/V3/BlobStore.cpp | 2 -- dbms/src/Storages/Page/V3/BlobStore.h | 3 +-- dbms/src/Storages/Page/V3/tests/gtest_blob_store.cpp | 10 +++------- 3 files changed, 4 insertions(+), 11 deletions(-) diff --git a/dbms/src/Storages/Page/V3/BlobStore.cpp b/dbms/src/Storages/Page/V3/BlobStore.cpp index f7f79c716db..f2429ba0871 100644 --- a/dbms/src/Storages/Page/V3/BlobStore.cpp +++ b/dbms/src/Storages/Page/V3/BlobStore.cpp @@ -847,7 +847,6 @@ void BlobStore::removePosFromStats(BlobFileId blob_id, BlobFileOffset off template typename BlobStore::PageMap BlobStore::read(FieldReadInfos & to_read, const ReadLimiterPtr & read_limiter) { - if (to_read.empty()) { return {}; @@ -1010,7 +1009,6 @@ typename BlobStore::PageMap BlobStore::read(FieldReadInfos & to_re buf.toString()); } - return page_map; } diff --git a/dbms/src/Storages/Page/V3/BlobStore.h b/dbms/src/Storages/Page/V3/BlobStore.h index cbaa4464ede..7d4bdca0aaa 100644 --- a/dbms/src/Storages/Page/V3/BlobStore.h +++ b/dbms/src/Storages/Page/V3/BlobStore.h @@ -181,9 +181,8 @@ class BlobStore : private Allocator std::unordered_map blob_files; #ifdef DBMS_PUBLIC_GTEST - size_t field_read_call_count = 0; + size_t field_read_call_count = 0; #endif - }; namespace u128 { diff --git a/dbms/src/Storages/Page/V3/tests/gtest_blob_store.cpp b/dbms/src/Storages/Page/V3/tests/gtest_blob_store.cpp index 3b42cdb7128..c989310a65f 100644 --- a/dbms/src/Storages/Page/V3/tests/gtest_blob_store.cpp +++ b/dbms/src/Storages/Page/V3/tests/gtest_blob_store.cpp @@ -1931,12 +1931,7 @@ try PageIdU64 page_id = 50; size_t buff_size = 20; - auto blob_store = BlobStore( - getCurrentTestName(), - file_provider, - delegator, - BlobConfig{}, - page_type_and_config); + auto blob_store = BlobStore(getCurrentTestName(), file_provider, delegator, BlobConfig{}, page_type_and_config); std::vector c_buff(buff_size); for (size_t j = 0; j < buff_size; ++j) @@ -1957,7 +1952,8 @@ try { blob_store.resetFieldReadCallCount(); BlobStore::FieldReadInfos read_infos; - read_infos.emplace_back(BlobStore::FieldReadInfo(buildV3Id(TEST_NAMESPACE_ID, page_id), entry, {0, 1, 2, 3, 4})); + read_infos.emplace_back( + BlobStore::FieldReadInfo(buildV3Id(TEST_NAMESPACE_ID, page_id), entry, {0, 1, 2, 3, 4})); auto page_map = blob_store.read(read_infos); Page page = page_map.at(page_id); ASSERT_EQ(page.fieldSize(), 5); From 6ee7a68105b560f74361ae49a9ea28b6711d19a1 Mon Sep 17 00:00:00 2001 From: yy <736262857@qq.com> Date: Sat, 11 Jul 2026 08:10:08 +0000 Subject: [PATCH 7/7] Modify field_read_call_count to be an atomic variable. Signed-off-by: yy <736262857@qq.com> --- dbms/src/Storages/Page/V3/BlobStore.cpp | 2 +- dbms/src/Storages/Page/V3/BlobStore.h | 6 +++--- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/dbms/src/Storages/Page/V3/BlobStore.cpp b/dbms/src/Storages/Page/V3/BlobStore.cpp index f2429ba0871..d68bb4dfbec 100644 --- a/dbms/src/Storages/Page/V3/BlobStore.cpp +++ b/dbms/src/Storages/Page/V3/BlobStore.cpp @@ -939,7 +939,7 @@ typename BlobStore::PageMap BlobStore::read(FieldReadInfos & to_re read(page_id_v3, entry.file_id, entry.offset + beg_offset, write_offset, size_to_read, read_limiter); #ifdef DBMS_PUBLIC_GTEST - ++field_read_call_count; + field_read_call_count.fetch_add(1, std::memory_order_relaxed); #endif for (size_t i = start_field_index; i <= end_field_index; ++i) { diff --git a/dbms/src/Storages/Page/V3/BlobStore.h b/dbms/src/Storages/Page/V3/BlobStore.h index 7d4bdca0aaa..71e0de79ddf 100644 --- a/dbms/src/Storages/Page/V3/BlobStore.h +++ b/dbms/src/Storages/Page/V3/BlobStore.h @@ -118,8 +118,8 @@ class BlobStore : private Allocator PageMap read(FieldReadInfos & to_read, const ReadLimiterPtr & read_limiter = nullptr); #ifdef DBMS_PUBLIC_GTEST - void resetFieldReadCallCount() { field_read_call_count = 0; } - size_t getFieldReadCallCount() { return field_read_call_count; } + void resetFieldReadCallCount() { field_read_call_count.store(0, std::memory_order_relaxed); } + size_t getFieldReadCallCount() { return field_read_call_count.load(std::memory_order_relaxed); } #endif @@ -181,7 +181,7 @@ class BlobStore : private Allocator std::unordered_map blob_files; #ifdef DBMS_PUBLIC_GTEST - size_t field_read_call_count = 0; + std::atomic field_read_call_count{0}; #endif }; namespace u128