diff --git a/dbms/src/Storages/Page/V3/BlobStore.cpp b/dbms/src/Storages/Page/V3/BlobStore.cpp index 3ab513d21f9..d68bb4dfbec 100644 --- a/dbms/src/Storages/Page/V3/BlobStore.cpp +++ b/dbms/src/Storages/Page/V3/BlobStore.cpp @@ -861,6 +861,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) { @@ -910,7 +911,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 +921,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 = 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) + { + ++field_idx; + end_field_index = fields[field_idx]; + } + + 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); - fields_offset_in_page.emplace(field_index, read_size_this_entry); - if constexpr (BLOBSTORE_CHECKSUM_ON_READ) + 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.fetch_add(1, std::memory_order_relaxed); +#endif + 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; } - read_size_this_entry += size_to_read; write_offset += size_to_read; + ++field_idx; } Page page(Trait::PageIdTrait::getU64ID(page_id_v3)); @@ -985,6 +1008,8 @@ typename BlobStore::PageMap BlobStore::read(FieldReadInfos & to_re pos, buf.toString()); } + + return page_map; } diff --git a/dbms/src/Storages/Page/V3/BlobStore.h b/dbms/src/Storages/Page/V3/BlobStore.h index 05c7f5fa299..71e0de79ddf 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.store(0, std::memory_order_relaxed); } + size_t getFieldReadCallCount() { return field_read_call_count.load(std::memory_order_relaxed); } + +#endif + #ifndef DBMS_PUBLIC_GTEST private: #endif @@ -173,6 +179,10 @@ class BlobStore : private Allocator std::mutex mtx_blob_files; std::unordered_map blob_files; + +#ifdef DBMS_PUBLIC_GTEST + std::atomic 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 ceb3449f775..c989310a65f 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,64 @@ 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; + + // 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); + } + + // 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); + } + + // 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 } // namespace DB::PS::V3::tests