diff --git a/build_tx_service.cmake b/build_tx_service.cmake index 1857afa4..f0ae597f 100644 --- a/build_tx_service.cmake +++ b/build_tx_service.cmake @@ -30,6 +30,17 @@ find_package(Protobuf REQUIRED) find_package(GFLAGS REQUIRED) find_package(MIMALLOC REQUIRED) +# boost_context for FlushDataWorker coroutine refactor (Phase 1) +if(CMAKE_BUILD_TYPE STREQUAL "Debug" AND CMAKE_CXX_FLAGS MATCHES "fsanitize=address") + find_library(Boost_CONTEXT_LIBRARY NAMES boost_context-asan) +else() + find_library(Boost_CONTEXT_LIBRARY NAMES boost_context) +endif() +if(NOT Boost_CONTEXT_LIBRARY) + message(FATAL_ERROR "libboost_context not found") +endif() +find_package(Boost 1.70 REQUIRED) + find_path(GFLAGS_INCLUDE_PATH gflags/gflags.h) find_library(GFLAGS_LIBRARY NAMES gflags libgflags) @@ -211,9 +222,9 @@ set(ELOQ_SOURCES ${ELOQ_SOURCES} ${PROTO_CC_FILES}) ADD_LIBRARY(txservice ${ELOQ_SOURCES}) -target_include_directories(txservice PUBLIC ${INCLUDE_DIR}) +target_include_directories(txservice PUBLIC ${INCLUDE_DIR} ${Boost_INCLUDE_DIRS}) -target_link_libraries(txservice PUBLIC ${LINK_LIB} ${PROTOBUF_LIBRARIES}) +target_link_libraries(txservice PUBLIC ${LINK_LIB} ${PROTOBUF_LIBRARIES} ${Boost_CONTEXT_LIBRARY}) if(WITH_JEMALLOC) target_link_libraries(txservice PUBLIC jemalloc_cfg) diff --git a/store_handler/bigtable_handler.cpp b/store_handler/bigtable_handler.cpp index 172c321f..5b51b829 100644 --- a/store_handler/bigtable_handler.cpp +++ b/store_handler/bigtable_handler.cpp @@ -613,7 +613,9 @@ void EloqDS::BigTableHandler::FetchTableRanges( LocalCcShards *shards = Sharder::Instance().GetLocalCcShards(); std::unique_lock heap_lk( shards->table_ranges_heap_mux_); - mi_override_thread(shards->GetTableRangesHeapThreadId()); + bool is_override_thd = mi_is_override_thread(); + mi_threadid_t prev_thd = + mi_override_thread(shards->GetTableRangesHeapThreadId()); mi_heap_t *prev_heap = mi_heap_set_default(shards->GetTableRangesHeap()); @@ -622,7 +624,14 @@ void EloqDS::BigTableHandler::FetchTableRanges( mono_key.size()); mi_heap_set_default(prev_heap); - mi_restore_default_thread_id(); + if (is_override_thd) + { + mi_override_thread(prev_thd); + } + else + { + mi_restore_default_thread_id(); + } } int32_t partition_id = @@ -758,7 +767,9 @@ void EloqDS::BigTableHandler::OnFetchRangeSlices( LocalCcShards *shards = Sharder::Instance().GetLocalCcShards(); std::unique_lock heap_lk( shards->table_ranges_heap_mux_); - mi_override_thread(shards->GetTableRangesHeapThreadId()); + bool is_override_thd = mi_is_override_thread(); + mi_threadid_t prev_thd = + mi_override_thread(shards->GetTableRangesHeapThreadId()); mi_heap_t *prev_heap = mi_heap_set_default(shards->GetTableRangesHeap()); @@ -773,7 +784,14 @@ void EloqDS::BigTableHandler::OnFetchRangeSlices( } mi_heap_set_default(prev_heap); - mi_restore_default_thread_id(); + if (is_override_thd) + { + mi_override_thread(prev_thd); + } + else + { + mi_restore_default_thread_id(); + } heap_lk.unlock(); } diff --git a/store_handler/data_store_service_client.cpp b/store_handler/data_store_service_client.cpp index 233864e0..1882b79d 100644 --- a/store_handler/data_store_service_client.cpp +++ b/store_handler/data_store_service_client.cpp @@ -33,6 +33,7 @@ #include #include #include +#include #include #include #include @@ -313,12 +314,25 @@ void DataStoreServiceClient::ScheduleTimerTasks() bool DataStoreServiceClient::PutAll( std::unordered_map>> - &flush_task) + &flush_task, + const std::function *yield_fptr, + const std::function *resume_fptr, + const std::function *sync_yield_fptr) +{ + return PutAllImpl(flush_task, yield_fptr, resume_fptr, sync_yield_fptr); +} + +bool DataStoreServiceClient::PutAllImpl( + std::unordered_map>> + &flush_task, + const std::function *yield_fptr, + const std::function *resume_fptr, + const std::function *sync_yield_fptr) { DLOG(INFO) << "DataStoreServiceClient::PutAll called with " << flush_task.size() << " tables to flush."; uint64_t now = txservice::LocalCcShards::ClockTsInMillseconds(); - // Create global coordinator SyncPutAllData *sync_putall = sync_putall_data_pool_.NextObject(); PoolableGuard sync_putall_guard(sync_putall); @@ -434,7 +448,17 @@ bool DataStoreServiceClient::PutAll( // Set up global coordinator sync_putall->total_partitions_ = sync_putall->partition_states_.size(); + // Set coroutine callbacks BEFORE starting async work (see plan risk + // analysis) + if (yield_fptr != nullptr && resume_fptr != nullptr) + { + sync_putall->SetCoroCallbacks(yield_fptr, resume_fptr); + } + // Start concurrent processing for each partition + constexpr size_t MAX_BATCH_WRITES_WITHOUT_YIELD = 10; + size_t batch_writes_since_yield = 0; + for (size_t i = 0; i < callback_data_list.size(); ++i) { auto *partition_state = sync_putall->partition_states_[i]; @@ -459,6 +483,13 @@ bool DataStoreServiceClient::PutAll( PartitionBatchCallback, first_batch.parts_cnt_per_key, first_batch.parts_cnt_per_record); + batch_writes_since_yield++; + if (sync_yield_fptr != nullptr && + batch_writes_since_yield >= MAX_BATCH_WRITES_WITHOUT_YIELD) + { + (*sync_yield_fptr)(); + batch_writes_since_yield = 0; + } } else { @@ -466,15 +497,14 @@ bool DataStoreServiceClient::PutAll( sync_putall->OnPartitionCompleted(); } } - // Wait for all partitions to complete + if (yield_fptr != nullptr && resume_fptr != nullptr) { - std::unique_lock lk(sync_putall->mux_); - while (sync_putall->completed_partitions_ < - sync_putall->total_partitions_) - { - sync_putall->cv_.wait(lk); - } + sync_putall->Wait(yield_fptr, resume_fptr); + } + else + { + sync_putall->Wait(); } // Check for errors @@ -490,6 +520,7 @@ bool DataStoreServiceClient::PutAll( callback_data->Clear(); callback_data->Free(); } + return false; } } @@ -505,6 +536,7 @@ bool DataStoreServiceClient::PutAll( metrics::kv_meter->Collect( metrics::NAME_KV_FLUSH_ROWS_TOTAL, records_count, "base"); } + return true; } @@ -520,14 +552,27 @@ bool DataStoreServiceClient::PutAll( * fails. */ bool DataStoreServiceClient::PersistKV( - const std::vector &kv_table_names) + const std::vector &kv_table_names, + const std::function *yield_fptr, + const std::function *resume_fptr) { SyncCallbackData *callback_data = sync_callback_data_pool_.NextObject(); PoolableGuard guard(callback_data); callback_data->Reset(); + if (yield_fptr != nullptr && resume_fptr != nullptr) + { + callback_data->SetCoroCallbacks(yield_fptr, resume_fptr); + } FlushData(kv_table_names, callback_data, &SyncCallback); - callback_data->Wait(); + if (yield_fptr != nullptr && resume_fptr != nullptr) + { + callback_data->Wait(yield_fptr, resume_fptr); + } + else + { + callback_data->Wait(); + } if (callback_data->Result().error_code() != EloqDS::remote::DataStoreError::NO_ERROR) { @@ -1529,16 +1574,7 @@ void DataStoreServiceClient::DispatchRangeSliceBatches( keys.size() > 0) { // Concurrency control: wait if limit reached, then increment - // counter - { - std::unique_lock lk(sync_concurrent->mux_); - while (sync_concurrent->unfinished_request_cnt_ >= - SyncConcurrentRequest::max_flying_write_count) - { - sync_concurrent->cv_.wait(lk); - } - sync_concurrent->unfinished_request_cnt_++; - } + sync_concurrent->WaitForCapacityAndIncrement(); // Dispatch current batch BatchWriteRecords(kv_table_name, @@ -1583,15 +1619,7 @@ void DataStoreServiceClient::DispatchRangeSliceBatches( if (keys.size() > 0) { // Concurrency control: wait if limit reached, then increment counter - { - std::unique_lock lk(sync_concurrent->mux_); - while (sync_concurrent->unfinished_request_cnt_ >= - SyncConcurrentRequest::max_flying_write_count) - { - sync_concurrent->cv_.wait(lk); - } - sync_concurrent->unfinished_request_cnt_++; - } + sync_concurrent->WaitForCapacityAndIncrement(); BatchWriteRecords(kv_table_name, kv_partition_id, @@ -1693,16 +1721,7 @@ void DataStoreServiceClient::DispatchRangeMetadataBatches( keys.size() > 0) { // Concurrency control: wait if limit reached, then increment - // counter - { - std::unique_lock lk(sync_concurrent->mux_); - while (sync_concurrent->unfinished_request_cnt_ >= - SyncConcurrentRequest::max_flying_write_count) - { - sync_concurrent->cv_.wait(lk); - } - sync_concurrent->unfinished_request_cnt_++; - } + sync_concurrent->WaitForCapacityAndIncrement(); // Dispatch current batch BatchWriteRecords(target_table_name, @@ -1747,16 +1766,7 @@ void DataStoreServiceClient::DispatchRangeMetadataBatches( if (keys.size() > 0) { // Concurrency control: wait if limit reached, then increment - // counter - { - std::unique_lock lk(sync_concurrent->mux_); - while (sync_concurrent->unfinished_request_cnt_ >= - SyncConcurrentRequest::max_flying_write_count) - { - sync_concurrent->cv_.wait(lk); - } - sync_concurrent->unfinished_request_cnt_++; - } + sync_concurrent->WaitForCapacityAndIncrement(); BatchWriteRecords(target_table_name, kv_partition_id, @@ -1776,7 +1786,10 @@ void DataStoreServiceClient::DispatchRangeMetadataBatches( } bool DataStoreServiceClient::UpdateRangeSlices( - const std::vector &update_range_slice_reqs) + const std::vector &update_range_slice_reqs, + const std::function *yield_fptr, + const std::function *resume_fptr, + const std::function *sync_yield_fptr) { if (update_range_slice_reqs.empty()) { @@ -1787,6 +1800,9 @@ bool DataStoreServiceClient::UpdateRangeSlices( slice_plans.reserve(update_range_slice_reqs.size()); RangeMetadataAccumulator meta_acc; + constexpr size_t MAX_ITERATIONS_WITHOUT_YIELD = 10; + size_t iterations_since_yield = 0; + // 1- First pass: Prepare slice batches and accumulate metadata for all // ranges for (auto &req : update_range_slice_reqs) @@ -1824,6 +1840,16 @@ bool DataStoreServiceClient::UpdateRangeSlices( segment_cnt, range_size, meta_acc); + + if (sync_yield_fptr != nullptr) + { + ++iterations_since_yield; + if (iterations_since_yield >= MAX_ITERATIONS_WITHOUT_YIELD) + { + (*sync_yield_fptr)(); + iterations_since_yield = 0; + } + } } // 2- Dispatch slice batches for all ranges concurrently (shared @@ -1832,25 +1858,33 @@ bool DataStoreServiceClient::UpdateRangeSlices( sync_concurrent_request_pool_.NextObject(); PoolableGuard slice_guard(slice_sync_concurrent); slice_sync_concurrent->Reset(); - for (const auto &[kv_partition_id, slice_plans] : slice_plans) + if (yield_fptr != nullptr && resume_fptr != nullptr) + { + slice_sync_concurrent->SetCoroCallbacks(yield_fptr, resume_fptr); + } + iterations_since_yield = 0; + for (const auto &[kv_partition_id, plans] : slice_plans) { // Call DispatchRangeSliceBatches once with all plans DispatchRangeSliceBatches(kv_range_slices_table_name, kv_partition_id, - slice_plans, + plans, slice_sync_concurrent); - } - // 3- Wait for slice requests to complete - { - std::unique_lock lk(slice_sync_concurrent->mux_); - slice_sync_concurrent->all_request_started_ = true; - while (slice_sync_concurrent->unfinished_request_cnt_ != 0) + if (sync_yield_fptr != nullptr) { - slice_sync_concurrent->cv_.wait(lk); + ++iterations_since_yield; + if (iterations_since_yield >= MAX_ITERATIONS_WITHOUT_YIELD) + { + (*sync_yield_fptr)(); + iterations_since_yield = 0; + } } } + // 3- Wait for slice requests to complete + slice_sync_concurrent->WaitForAll(); + if (slice_sync_concurrent->result_.error_code() != remote::DataStoreError::NO_ERROR) { @@ -1866,11 +1900,23 @@ bool DataStoreServiceClient::UpdateRangeSlices( sync_callback_data_pool_.NextObject(); PoolableGuard guard(flush_slices_callback_data); flush_slices_callback_data->Reset(); + if (yield_fptr != nullptr && resume_fptr != nullptr) + { + flush_slices_callback_data->SetCoroCallbacks(yield_fptr, + resume_fptr); + } std::vector kv_slices_table_names; kv_slices_table_names.emplace_back(kv_range_slices_table_name); FlushData( kv_slices_table_names, flush_slices_callback_data, &SyncCallback); - flush_slices_callback_data->Wait(); + if (yield_fptr != nullptr && resume_fptr != nullptr) + { + flush_slices_callback_data->Wait(yield_fptr, resume_fptr); + } + else + { + flush_slices_callback_data->Wait(); + } if (flush_slices_callback_data->Result().error_code() != EloqDS::remote::DataStoreError::NO_ERROR) { @@ -1885,18 +1931,15 @@ bool DataStoreServiceClient::UpdateRangeSlices( sync_concurrent_request_pool_.NextObject(); PoolableGuard meta_guard(meta_sync_concurrent); meta_sync_concurrent->Reset(); + if (yield_fptr != nullptr && resume_fptr != nullptr) + { + meta_sync_concurrent->SetCoroCallbacks(yield_fptr, resume_fptr); + } DispatchRangeMetadataBatches( kv_range_table_name, meta_acc, meta_sync_concurrent); // 5- Wait for metadata requests to complete - { - std::unique_lock lk(meta_sync_concurrent->mux_); - meta_sync_concurrent->all_request_started_ = true; - while (meta_sync_concurrent->unfinished_request_cnt_ != 0) - { - meta_sync_concurrent->cv_.wait(lk); - } - } + meta_sync_concurrent->WaitForAll(); // 6- Check for errors if (meta_sync_concurrent->result_.error_code() != @@ -1913,10 +1956,21 @@ bool DataStoreServiceClient::UpdateRangeSlices( SyncCallbackData *callback_data = sync_callback_data_pool_.NextObject(); PoolableGuard guard(callback_data); callback_data->Reset(); + if (yield_fptr != nullptr && resume_fptr != nullptr) + { + callback_data->SetCoroCallbacks(yield_fptr, resume_fptr); + } std::vector kv_range_table_names; kv_range_table_names.emplace_back(kv_range_table_name); FlushData(kv_range_table_names, callback_data, &SyncCallback); - callback_data->Wait(); + if (yield_fptr != nullptr && resume_fptr != nullptr) + { + callback_data->Wait(yield_fptr, resume_fptr); + } + else + { + callback_data->Wait(); + } if (callback_data->Result().error_code() != EloqDS::remote::DataStoreError::NO_ERROR) { @@ -1977,14 +2031,7 @@ bool DataStoreServiceClient::UpdateRangeSlices( slice_sync_concurrent); // 3- Wait for slice requests to complete. Make sure meta data is updated // after all slice info is written. - { - std::unique_lock lk(slice_sync_concurrent->mux_); - slice_sync_concurrent->all_request_started_ = true; - while (slice_sync_concurrent->unfinished_request_cnt_ != 0) - { - slice_sync_concurrent->cv_.wait(lk); - } - } + slice_sync_concurrent->WaitForAll(); if (slice_sync_concurrent->result_.error_code() != remote::DataStoreError::NO_ERROR) @@ -2035,14 +2082,7 @@ bool DataStoreServiceClient::UpdateRangeSlices( kv_range_table_name, meta_acc, meta_sync_concurrent); // 5- Wait for metadata requests to complete - { - std::unique_lock lk(meta_sync_concurrent->mux_); - meta_sync_concurrent->all_request_started_ = true; - while (meta_sync_concurrent->unfinished_request_cnt_ != 0) - { - meta_sync_concurrent->cv_.wait(lk); - } - } + meta_sync_concurrent->WaitForAll(); // 6- Check for errors if (meta_sync_concurrent->result_.error_code() != @@ -2161,14 +2201,7 @@ bool DataStoreServiceClient::UpsertRanges( } // 3- Wait for slice requests to complete - { - std::unique_lock lk(slice_sync_concurrent->mux_); - slice_sync_concurrent->all_request_started_ = true; - while (slice_sync_concurrent->unfinished_request_cnt_ != 0) - { - slice_sync_concurrent->cv_.wait(lk); - } - } + slice_sync_concurrent->WaitForAll(); if (slice_sync_concurrent->result_.error_code() != remote::DataStoreError::NO_ERROR) { @@ -2207,14 +2240,7 @@ bool DataStoreServiceClient::UpsertRanges( kv_range_table_name, meta_acc, meta_sync_concurrent); // 5- Wait for metadata requests to complete - { - std::unique_lock lk(meta_sync_concurrent->mux_); - meta_sync_concurrent->all_request_started_ = true; - while (meta_sync_concurrent->unfinished_request_cnt_ != 0) - { - meta_sync_concurrent->cv_.wait(lk); - } - } + meta_sync_concurrent->WaitForAll(); // 6- Check for errors if (meta_sync_concurrent->result_.error_code() != @@ -2908,7 +2934,19 @@ void DataStoreServiceClient::DecodeArchiveValue( bool DataStoreServiceClient::PutArchivesAll( std::unordered_map>> - &flush_task) + &flush_task, + const std::function *yield_fptr, + const std::function *resume_fptr) +{ + return PutArchivesAllImpl(flush_task, yield_fptr, resume_fptr); +} + +bool DataStoreServiceClient::PutArchivesAllImpl( + std::unordered_map>> + &flush_task, + const std::function *yield_fptr, + const std::function *resume_fptr) { std::unordered_map< uint32_t, @@ -2971,6 +3009,11 @@ bool DataStoreServiceClient::PutArchivesAll( PoolableGuard guard(sync_concurrent); sync_concurrent->Reset(); + if (yield_fptr != nullptr && resume_fptr != nullptr) + { + sync_concurrent->SetCoroCallbacks(yield_fptr, resume_fptr); + } + size_t recs_cnt = archive_ptrs.size(); keys.reserve(recs_cnt * parts_cnt_per_key); records.reserve(recs_cnt * parts_cnt_per_record); @@ -2986,15 +3029,7 @@ bool DataStoreServiceClient::PutArchivesAll( if (write_batch_size >= MAX_WRITE_BATCH_SIZE) { // Wait for in-flight requests to decrease if limit reached - { - std::unique_lock lk(sync_concurrent->mux_); - while (sync_concurrent->unfinished_request_cnt_ >= - SyncConcurrentRequest::max_flying_write_count) - { - sync_concurrent->cv_.wait(lk); - } - sync_concurrent->unfinished_request_cnt_++; - } + sync_concurrent->WaitForCapacityAndIncrement(); BatchWriteRecords(kv_mvcc_archive_name, partition_id, data_shard_id, @@ -3104,14 +3139,7 @@ bool DataStoreServiceClient::PutArchivesAll( } // Wait the result. - { - std::unique_lock lk(sync_concurrent->mux_); - sync_concurrent->all_request_started_ = true; - while (sync_concurrent->unfinished_request_cnt_ != 0) - { - sync_concurrent->cv_.wait(lk); - } - } + sync_concurrent->WaitForAll(); if (sync_concurrent->result_.error_code() != remote::DataStoreError::NO_ERROR) @@ -3147,7 +3175,19 @@ bool DataStoreServiceClient::PutArchivesAll( bool DataStoreServiceClient::CopyBaseToArchive( std::unordered_map>> - &flush_task) + &flush_task, + const std::function *yield_fptr, + const std::function *resume_fptr) +{ + return CopyBaseToArchiveImpl(flush_task, yield_fptr, resume_fptr); +} + +bool DataStoreServiceClient::CopyBaseToArchiveImpl( + std::unordered_map>> + &flush_task, + const std::function *yield_fptr, + const std::function *resume_fptr) { // Prepare for the copied base table data to be flushed to the archive table std::unordered_map callback_datas; - callback_datas.reserve(base_vec.size()); + std::atomic waiting{ + false}; // shared: last callback must notify Wait + // Use deque: ReadBaseForArchiveCallbackData contains std::atomic + // which is non-copyable/non-movable; vector requires element + // move/copy on realloc, deque does not. + std::deque callback_datas; for (size_t i = 0; i < base_vec.size(); ++i) { - callback_datas.emplace_back(mtx, cv, flying_cnt, error_code); + callback_datas.emplace_back( + mtx, cv, flying_cnt, error_code, waiting, resume_fptr); } for (size_t base_idx = 0; base_idx < base_vec.size(); ++base_idx) @@ -3212,7 +3257,14 @@ bool DataStoreServiceClient::CopyBaseToArchive( if (flying_cnt >= MAX_FLYING_READ_COUNT) { - callback_data->Wait(); + if (yield_fptr && resume_fptr) + { + callback_data->Wait(yield_fptr, resume_fptr); + } + else + { + callback_data->Wait(); + } } if (callback_data->GetErrorCode() != 0) { @@ -3223,6 +3275,11 @@ bool DataStoreServiceClient::CopyBaseToArchive( } // Wait the result all return before returning to avoid referencing // invalid memory in callback. + if (yield_fptr && resume_fptr) + { + callback_datas[0].Wait(yield_fptr, resume_fptr); + } + else { std::unique_lock lk(mtx); while (flying_cnt > 0) @@ -3321,7 +3378,7 @@ bool DataStoreServiceClient::CopyBaseToArchive( { // Put the archive records to the archive table. // This is a sync call - bool ret = PutArchivesAll(archive_flush_task); + bool ret = PutArchivesAll(archive_flush_task, yield_fptr, resume_fptr); if (!ret) { return false; @@ -4936,14 +4993,7 @@ bool DataStoreServiceClient::DeleteTableRanges( for (uint32_t kv_partition_id = 0; kv_partition_id < kv_partition_cnt; ++kv_partition_id) { - std::unique_lock lk( - delete_slices_sync_concurrent->mux_); - while (delete_slices_sync_concurrent->unfinished_request_cnt_ >= - SyncConcurrentRequest::max_flying_write_count) - { - delete_slices_sync_concurrent->cv_.wait(lk); - } - delete_slices_sync_concurrent->unfinished_request_cnt_++; + delete_slices_sync_concurrent->WaitForCapacityAndIncrement(); // get shard id uint32_t data_shard_id = @@ -4959,16 +5009,7 @@ bool DataStoreServiceClient::DeleteTableRanges( SyncConcurrentRequestCallback); } - // callback_data->Wait(); - { - std::unique_lock lk( - delete_slices_sync_concurrent->mux_); - delete_slices_sync_concurrent->all_request_started_ = true; - while (delete_slices_sync_concurrent->unfinished_request_cnt_ != 0) - { - delete_slices_sync_concurrent->cv_.wait(lk); - } - } + delete_slices_sync_concurrent->WaitForAll(); if (delete_slices_sync_concurrent->result_.error_code() != EloqDS::remote::DataStoreError::NO_ERROR) @@ -4984,14 +5025,7 @@ bool DataStoreServiceClient::DeleteTableRanges( for (int32_t kv_partition_id = 0; kv_partition_id < kv_partition_cnt; ++kv_partition_id) { - std::unique_lock lk( - delete_slices_sync_concurrent->mux_); - while (delete_slices_sync_concurrent->unfinished_request_cnt_ >= - SyncConcurrentRequest::max_flying_write_count) - { - delete_slices_sync_concurrent->cv_.wait(lk); - } - delete_slices_sync_concurrent->unfinished_request_cnt_++; + delete_slices_sync_concurrent->WaitForCapacityAndIncrement(); // get shard id uint32_t data_shard_id = @@ -5006,15 +5040,7 @@ bool DataStoreServiceClient::DeleteTableRanges( SyncConcurrentRequestCallback); } - { - std::unique_lock lk( - delete_slices_sync_concurrent->mux_); - delete_slices_sync_concurrent->all_request_started_ = true; - while (delete_slices_sync_concurrent->unfinished_request_cnt_ != 0) - { - delete_slices_sync_concurrent->cv_.wait(lk); - } - } + delete_slices_sync_concurrent->WaitForAll(); if (delete_slices_sync_concurrent->result_.error_code() != EloqDS::remote::DataStoreError::NO_ERROR) diff --git a/store_handler/data_store_service_client.h b/store_handler/data_store_service_client.h index e086a17f..b3b02792 100644 --- a/store_handler/data_store_service_client.h +++ b/store_handler/data_store_service_client.h @@ -229,10 +229,14 @@ class DataStoreServiceClient : public txservice::store::DataStoreHandler * @param node_group * @return whether all entries are written to data store successfully */ - bool PutAll(std::unordered_map< - std::string_view, - std::vector>> - &flush_task) override; + bool PutAll( + std::unordered_map< + std::string_view, + std::vector>> + &flush_task, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr, + const std::function *sync_yield_fptr = nullptr) override; bool NeedPersistKV() override { @@ -247,7 +251,9 @@ class DataStoreServiceClient : public txservice::store::DataStoreHandler * @param node_group * @return whether all entries are written to data store successfully */ - bool PersistKV(const std::vector &kv_table_names) override; + bool PersistKV(const std::vector &kv_table_names, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr) override; void UpsertTable( const txservice::TableSchema *old_table_schema, @@ -336,8 +342,12 @@ class DataStoreServiceClient : public txservice::store::DataStoreHandler uint32_t range_partition_id, txservice::FillStoreSliceCc *load_slice_req) override; - bool UpdateRangeSlices(const std::vector - &update_range_slice_reqs) override; + bool UpdateRangeSlices( + const std::vector + &update_range_slice_reqs, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr, + const std::function *sync_yield_fptr = nullptr) override; bool UpdateRangeSlices(const txservice::TableName &table_name, uint64_t version, @@ -423,10 +433,14 @@ class DataStoreServiceClient : public txservice::store::DataStoreHandler * @brief Write batch historical versions into DataStore. * */ - bool PutArchivesAll(std::unordered_map< - std::string_view, - std::vector>> - &flush_task) override; + bool PutArchivesAll( + std::unordered_map< + std::string_view, + std::vector>> + &flush_task, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr) override; + /** * @brief Copy record from base/sk table to mvcc_archives. */ @@ -434,7 +448,9 @@ class DataStoreServiceClient : public txservice::store::DataStoreHandler std::unordered_map< std::string_view, std::vector>> - &flush_task) override; + &flush_task, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr) override; /** * @brief Get the latest visible(commit_ts <= upper_bound_ts) @@ -598,6 +614,30 @@ class DataStoreServiceClient : public txservice::store::DataStoreHandler bool is_range_partition) const; private: + bool PutArchivesAllImpl( + std::unordered_map< + std::string_view, + std::vector>> + &flush_task, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr); + + bool PutAllImpl(std::unordered_map< + std::string_view, + std::vector>> + &flush_task, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr, + const std::function *sync_yield_fptr = nullptr); + + bool CopyBaseToArchiveImpl( + std::unordered_map< + std::string_view, + std::vector>> + &flush_task, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr); + int32_t MapKeyHashToPartitionId(const txservice::TxKey &key) const { return txservice::Sharder::MapKeyHashToHashPartitionId(key.Hash()); diff --git a/store_handler/data_store_service_client_closure.cpp b/store_handler/data_store_service_client_closure.cpp index bdddbec3..950cabc6 100644 --- a/store_handler/data_store_service_client_closure.cpp +++ b/store_handler/data_store_service_client_closure.cpp @@ -513,6 +513,76 @@ void SyncConcurrentRequestCallback(void *data, callback_data->Finish(result); } +void SyncPutAllData::Wait() +{ + std::unique_lock lk(mux_); + while (completed_partitions_ < total_partitions_) + { + cv_.wait(lk); + } +} + +void SyncPutAllData::Wait(const std::function *yield_fn, + const std::function *resume_fn) +{ + if (yield_fn == nullptr || resume_fn == nullptr) + { + Wait(); + return; + } + std::unique_lock lk(mux_); + while (completed_partitions_ < total_partitions_) + { + waiting_.store(true, std::memory_order_release); + lk.unlock(); + (*yield_fn)(); + lk.lock(); + waiting_.store(false, std::memory_order_release); + } +} + +void SyncConcurrentRequest::WaitForCapacityAndIncrement() +{ + std::unique_lock lk(mux_); + while (unfinished_request_cnt_ >= max_flying_write_count) + { + if (yield_fn_ != nullptr && resume_fn_ != nullptr) + { + waiting_.store(true, std::memory_order_release); + lk.unlock(); + (*yield_fn_)(); + lk.lock(); + waiting_.store(false, std::memory_order_release); + } + else + { + cv_.wait(lk); + } + } + unfinished_request_cnt_++; +} + +void SyncConcurrentRequest::WaitForAll() +{ + std::unique_lock lk(mux_); + all_request_started_ = true; + while (unfinished_request_cnt_ != 0) + { + if (yield_fn_ != nullptr && resume_fn_ != nullptr) + { + waiting_.store(true, std::memory_order_release); + lk.unlock(); + (*yield_fn_)(); + lk.lock(); + waiting_.store(false, std::memory_order_release); + } + else + { + cv_.wait(lk); + } + } +} + void PartitionBatchCallback(void *data, ::google::protobuf::Closure *closure, DataStoreServiceClient &client, diff --git a/store_handler/data_store_service_client_closure.h b/store_handler/data_store_service_client_closure.h index b8c3813c..00de7d0a 100644 --- a/store_handler/data_store_service_client_closure.h +++ b/store_handler/data_store_service_client_closure.h @@ -25,6 +25,8 @@ #include #include +#include +#include #include #include #include @@ -82,19 +84,42 @@ struct SyncCallbackData : public Poolable { finished_ = false; result_.Clear(); + waiting_.store(false); + yield_fn_ = nullptr; + resume_fn_ = nullptr; } virtual void Clear() override { finished_ = false; result_.Clear(); + waiting_.store(false); + yield_fn_ = nullptr; + resume_fn_ = nullptr; + } + + void SetCoroCallbacks(const std::function *yield_fn, + const std::function *resume_fn) + { + yield_fn_ = yield_fn; + resume_fn_ = resume_fn; } virtual void Notify() { std::unique_lock lk(mtx_); finished_ = true; - cv_.notify_one(); + if (resume_fn_ != nullptr && waiting_.load(std::memory_order_acquire)) + { + waiting_.store(false, std::memory_order_release); + auto *fn = resume_fn_; + lk.unlock(); + (*fn)(); + } + else if (resume_fn_ == nullptr) + { + cv_.notify_one(); + } } virtual void Wait() @@ -106,6 +131,25 @@ struct SyncCallbackData : public Poolable } } + void Wait(const std::function *yield_fn, + const std::function *resume_fn) + { + if (yield_fn == nullptr || resume_fn == nullptr) + { + Wait(); + return; + } + std::unique_lock lk(mtx_); + while (!finished_) + { + waiting_.store(true, std::memory_order_release); + lk.unlock(); + (*yield_fn)(); + lk.lock(); + waiting_.store(false, std::memory_order_release); + } + } + remote::CommonResult &Result() { return result_; @@ -123,6 +167,11 @@ struct SyncCallbackData : public Poolable bool finished_; remote::CommonResult result_; + + // Coroutine yield/resume support + const std::function *yield_fn_{nullptr}; + const std::function *resume_fn_{nullptr}; + std::atomic waiting_{false}; }; /** @@ -215,12 +264,18 @@ struct SyncPutAllData : public Poolable partition_states_.clear(); completed_partitions_ = 0; total_partitions_ = 0; + waiting_.store(false); + yield_fn_ = nullptr; + resume_fn_ = nullptr; } virtual void Clear() override { completed_partitions_ = 0; total_partitions_ = 0; + waiting_.store(false); + yield_fn_ = nullptr; + resume_fn_ = nullptr; for (auto *partition_state : partition_states_) { partition_state->Clear(); @@ -228,15 +283,42 @@ struct SyncPutAllData : public Poolable } partition_states_.clear(); } + + void SetCoroCallbacks(const std::function *yield_fn, + const std::function *resume_fn) + { + yield_fn_ = yield_fn; + resume_fn_ = resume_fn; + } + + void Wait(); + + void Wait(const std::function *yield_fn, + const std::function *resume_fn); + void OnPartitionCompleted() { std::unique_lock lk(mux_); completed_partitions_++; if (completed_partitions_ >= total_partitions_) { - cv_.notify_one(); + if (resume_fn_ && waiting_.load(std::memory_order_acquire)) + { + waiting_.store(false, std::memory_order_release); + auto *fn = resume_fn_; + lk.unlock(); + (*fn)(); + // Do not re-lock: resume_fn may acquire flush_worker_mux; + // coroutine will need mux_ when it resumes. Unlock before + // resume_fn avoids deadlock. + } + else if (!resume_fn_) + { + cv_.notify_one(); + } } } + mutable bthread::Mutex mux_; bthread::ConditionVariable cv_; @@ -244,6 +326,11 @@ struct SyncPutAllData : public Poolable std::vector partition_states_; int32_t completed_partitions_{0}; int32_t total_partitions_{0}; + + // Coroutine yield/resume support + const std::function *yield_fn_{nullptr}; + const std::function *resume_fn_{nullptr}; + std::atomic waiting_{false}; }; /** @@ -273,6 +360,9 @@ struct SyncConcurrentRequest : public Poolable unfinished_request_cnt_ = 0; all_request_started_ = false; result_.Clear(); + waiting_.store(false); + yield_fn_ = nullptr; + resume_fn_ = nullptr; } virtual void Clear() override @@ -280,8 +370,22 @@ struct SyncConcurrentRequest : public Poolable unfinished_request_cnt_ = 0; all_request_started_ = false; result_.Clear(); + waiting_.store(false); + yield_fn_ = nullptr; + resume_fn_ = nullptr; + } + + void SetCoroCallbacks(const std::function *yield_fn, + const std::function *resume_fn) + { + yield_fn_ = yield_fn; + resume_fn_ = resume_fn; } + void WaitForCapacityAndIncrement(); + + void WaitForAll(); + void Finish(const remote::CommonResult &res) { std::unique_lock lk(mux_); @@ -295,7 +399,18 @@ struct SyncConcurrentRequest : public Poolable if ((all_request_started_ && unfinished_request_cnt_ == 0) || unfinished_request_cnt_ == max_flying_write_count - 1) { - cv_.notify_one(); + if (resume_fn_ && waiting_.load(std::memory_order_acquire)) + { + waiting_.store(false, std::memory_order_release); + auto *fn = resume_fn_; + lk.unlock(); + (*fn)(); + // Do not re-lock: resume_fn may acquire flush_worker_mux + } + else if (!resume_fn_) + { + cv_.notify_one(); + } } } @@ -303,6 +418,11 @@ struct SyncConcurrentRequest : public Poolable int32_t unfinished_request_cnt_{0}; bool all_request_started_{false}; remote::CommonResult result_; + + // Coroutine yield/resume support + const std::function *yield_fn_{nullptr}; + const std::function *resume_fn_{nullptr}; + std::atomic waiting_{false}; mutable bthread::Mutex mux_; bthread::ConditionVariable cv_; }; @@ -446,14 +566,19 @@ void SyncCallback(void *data, struct ReadBaseForArchiveCallbackData { - ReadBaseForArchiveCallbackData(bthread::Mutex &mtx, - bthread::ConditionVariable &cv, - size_t &flying_read_cnt, - int &error_code) + ReadBaseForArchiveCallbackData( + bthread::Mutex &mtx, + bthread::ConditionVariable &cv, + size_t &flying_read_cnt, + int &error_code, + std::atomic &waiting, + const std::function *resume_fn = nullptr) : mtx_(mtx), cv_(cv), flying_read_cnt_(flying_read_cnt), error_code_(error_code), + waiting_(waiting), + resume_fn_(resume_fn), partition_id_(0), key_str_(), value_str_(), @@ -462,6 +587,11 @@ struct ReadBaseForArchiveCallbackData { } + void SetCoroCallbacks(const std::function *resume_fn) + { + resume_fn_ = resume_fn; + } + void ResetResult() { partition_id_ = 0; @@ -480,6 +610,28 @@ struct ReadBaseForArchiveCallbackData } } + void Wait(const std::function *yield_fn, + const std::function *resume_fn) + { + if (!yield_fn || !resume_fn) + { + Wait(); + return; + } + std::unique_lock lk(mtx_); + resume_fn_ = resume_fn; + while (flying_read_cnt_ > 0) + { + waiting_.store(true, std::memory_order_release); + lk.unlock(); + (*yield_fn)(); + lk.lock(); + waiting_.store(false, std::memory_order_release); + // If flying_read_cnt_ still > 0 after resume, loop and yield again. + // Do not use cv_.wait(): resume_fn path never notifies cv_. + } + } + size_t AddFlyingReadCount() { std::unique_lock lk(mtx_); @@ -493,7 +645,20 @@ struct ReadBaseForArchiveCallbackData flying_read_cnt_--; if (flying_read_cnt_ == 0) { - cv_.notify_one(); + if (resume_fn_ && waiting_.load(std::memory_order_acquire)) + { + waiting_.store(false, std::memory_order_release); + auto *fn = resume_fn_; + lk.unlock(); + (*fn)(); + // Do not re-lock: resume_fn may acquire flush_worker_mux. + // Do not access members after fn(): object may be destroyed. + return 0; + } + else if (!resume_fn_) + { + cv_.notify_one(); + } } return flying_read_cnt_; } @@ -537,6 +702,9 @@ struct ReadBaseForArchiveCallbackData bthread::ConditionVariable &cv_; size_t &flying_read_cnt_; int &error_code_; + std::atomic + &waiting_; // shared: any callback's DecreaseFlyingReadCount can notify + const std::function *resume_fn_{nullptr}; int32_t partition_id_; std::string_view key_str_; std::string value_str_; diff --git a/store_handler/eloq_data_store_service/eloq_store_data_store.cpp b/store_handler/eloq_data_store_service/eloq_store_data_store.cpp index fe796c7a..e6b49ca5 100644 --- a/store_handler/eloq_data_store_service/eloq_store_data_store.cpp +++ b/store_handler/eloq_data_store_service/eloq_store_data_store.cpp @@ -59,6 +59,12 @@ inline void BuildKey(const WriteRecordsRequest &write_req, uint16_t key_parts, std::string &key_out) { + size_t total_size = 0; + for (uint16_t i = 0; i < key_parts; ++i) + { + total_size += write_req.GetKeyPart(key_first_idx + i).size(); + } + key_out.reserve(key_out.size() + total_size); size_t part_idx = key_first_idx; for (uint16_t i = 0; i < key_parts; ++i, ++part_idx) { @@ -72,6 +78,12 @@ inline void BuildValue(const WriteRecordsRequest &write_req, uint16_t rec_parts, std::string &rec_out) { + size_t total_size = 0; + for (uint16_t i = 0; i < rec_parts; ++i) + { + total_size += write_req.GetRecordPart(rec_first_idx + i).size(); + } + rec_out.reserve(rec_out.size() + total_size); size_t part_idx = rec_first_idx; for (uint16_t i = 0; i < rec_parts; ++i, ++part_idx) { @@ -214,26 +226,28 @@ void EloqStoreDataStore::BatchWriteRecords(WriteRecordsRequest *write_req) ::eloqstore::BatchWriteRequest &kv_write_req = write_op->EloqStoreRequest(); - std::vector<::eloqstore::WriteDataEntry> entries; - size_t rec_cnt = write_req->RecordsCount(); - entries.reserve(rec_cnt); + const size_t rec_cnt = write_req->RecordsCount(); const uint16_t parts_per_key = write_req->PartsCountPerKey(); const uint16_t parts_per_record = write_req->PartsCountPerRecord(); - size_t first_idx = 0; + + std::vector<::eloqstore::WriteDataEntry> entries(rec_cnt); + size_t key_offset = 0; + size_t val_offset = 0; + for (size_t i = 0; i < rec_cnt; ++i) { - ::eloqstore::WriteDataEntry entry; - first_idx = i * parts_per_key; - BuildKey(*write_req, first_idx, parts_per_key, entry.key_); - first_idx = i * parts_per_record; - BuildValue(*write_req, first_idx, parts_per_record, entry.val_); + ::eloqstore::WriteDataEntry &entry = entries[i]; + BuildKey(*write_req, key_offset, parts_per_key, entry.key_); + BuildValue(*write_req, val_offset, parts_per_record, entry.val_); + key_offset += parts_per_key; + val_offset += parts_per_record; + entry.timestamp_ = write_req->GetRecordTs(i); entry.op_ = (write_req->KeyOpType(i) == WriteOpType::PUT ? ::eloqstore::WriteOp::Upsert : ::eloqstore::WriteOp::Delete); - uint64_t ttl = write_req->GetRecordTtl(i); - entry.expire_ts_ = ttl == UINT64_MAX ? 0 : ttl; - entries.emplace_back(std::move(entry)); + const uint64_t ttl = write_req->GetRecordTtl(i); + entry.expire_ts_ = (ttl == UINT64_MAX) ? 0u : ttl; } if (!std::ranges::is_sorted( @@ -250,6 +264,7 @@ void EloqStoreDataStore::BatchWriteRecords(WriteRecordsRequest *write_req) kv_write_req.SetArgs(eloq_store_table_id, std::move(entries)); uint64_t user_data = reinterpret_cast(write_op); + if (!eloq_store_service_->ExecAsyn(&kv_write_req, user_data, OnBatchWrite)) { LOG(ERROR) << "Send write request to EloqStore failed for table: " diff --git a/store_handler/rocksdb_handler.cpp b/store_handler/rocksdb_handler.cpp index fb81f021..19d0f9b2 100644 --- a/store_handler/rocksdb_handler.cpp +++ b/store_handler/rocksdb_handler.cpp @@ -26,6 +26,7 @@ #include #include #include +#include #include #include #include @@ -87,6 +88,83 @@ namespace EloqKV namespace { +/** + * Minimal sync callback for RocksDB PutAll/PersistKV yield/resume offload. + * Mirrors SyncCallbackData pattern from data_store_service_client_closure.h. + */ +struct RocksDBWriteSyncCallback +{ + bthread::Mutex mtx_; + bthread::ConditionVariable cv_; + bool finished_{false}; + bool success_{true}; + const std::function *yield_fn_{nullptr}; + const std::function *resume_fn_{nullptr}; + std::atomic waiting_{false}; + + void Reset() + { + finished_ = false; + success_ = true; + waiting_.store(false); + yield_fn_ = nullptr; + resume_fn_ = nullptr; + } + + void SetCoroCallbacks(const std::function *yield_fn, + const std::function *resume_fn) + { + yield_fn_ = yield_fn; + resume_fn_ = resume_fn; + } + + void Notify(bool ok) + { + std::unique_lock lk(mtx_); + success_ = ok; + finished_ = true; + if (resume_fn_ != nullptr && waiting_.load(std::memory_order_acquire)) + { + waiting_.store(false, std::memory_order_release); + auto *fn = resume_fn_; + lk.unlock(); + (*fn)(); + } + else if (resume_fn_ == nullptr) + { + cv_.notify_one(); + } + } + + void Wait() + { + std::unique_lock lk(mtx_); + while (!finished_) + { + cv_.wait(lk); + } + } + + void Wait(const std::function *yield_fn, + const std::function *resume_fn) + { + if (yield_fn == nullptr || resume_fn == nullptr) + { + Wait(); + return; + } + std::unique_lock lk(mtx_); + while (!finished_) + { + waiting_.store(true, std::memory_order_release); + lk.unlock(); + (*yield_fn)(); + lk.lock(); + waiting_.store(false, std::memory_order_release); + } + } +}; + std::string NormalizeDbPath(const std::string &path) { std::string normalized = path; @@ -433,13 +511,147 @@ void RocksDBHandler::DeserializeRecord(const char *payload, bool RocksDBHandler::PutAll( std::unordered_map>> - &batch) + &batch, + const std::function *yield_fptr, + const std::function *resume_fptr, + const std::function *sync_yield_fptr) { + (void) sync_yield_fptr; std::thread::id this_id = std::this_thread::get_id(); if (batch.empty()) { return true; } + + if (yield_fptr != nullptr && resume_fptr != nullptr) + { + RocksDBWriteSyncCallback callback; + callback.Reset(); + callback.SetCoroCallbacks(yield_fptr, resume_fptr); + + query_worker_pool_->SubmitWork( + [this, &batch, &callback, this_id](size_t) + { + std::shared_lock db_lk(db_mux_); + auto db = GetDBPtr(); + if (!db) + { + callback.Notify(false); + return; + } + rocksdb::WriteOptions write_options; + write_options.disableWAL = true; + write_options.no_slowdown = false; + rocksdb::WriteBatch write_batch; + uint64_t write_batch_size = 0; + + for (auto &[kv_cf_name, flush_task_entries] : batch) + { + rocksdb::ColumnFamilyHandle *cfh = + GetColumnFamilyHandler(kv_cf_name.data()); + if (cfh == nullptr) + { + LOG(ERROR) << "Failed to get column family, cf name: " + << kv_cf_name; + callback.Notify(false); + return; + } + uint64_t now = + txservice::LocalCcShards::ClockTsInMillseconds(); + for (auto &flush_task_entry : flush_task_entries) + { + for (auto &flush_rec : + *flush_task_entry->data_sync_vec_) + { + txservice::TxKey key = flush_rec.Key(); + std::string rocksdb_key; + EncodeToKvKey(key, rocksdb_key); + + if (flush_rec.payload_status_ == + txservice::RecordStatus::Normal && + flush_rec.Payload()->GetTTL() > now) + { + std::vector rec_buf; + SerializeFlushRecord(flush_rec, rec_buf); + write_batch_size += rocksdb_key.size(); + write_batch_size += rec_buf.size(); + write_batch.Put( + cfh, + rocksdb::Slice(rocksdb_key.data(), + rocksdb_key.size()), + rocksdb::Slice(rec_buf.data(), + rec_buf.size())); + } + else + { + write_batch_size += rocksdb_key.size(); + write_batch.Delete( + cfh, + rocksdb::Slice(rocksdb_key.data(), + rocksdb_key.size())); + } + + if (write_batch_size >= batch_write_size_) + { + auto status = + db->Write(write_options, &write_batch); + if (!status.ok()) + { + LOG(ERROR) + << "PutAll end failed " + << ", thread id: " << this_id + << ", result:" + << static_cast(status.ok()) + << ", batch size:" << batch.size() + << ", error: " << status.ToString() + << ", error code: " << status.code(); + callback.Notify(false); + return; + } + if (metrics::enable_kv_metrics) + { + metrics::kv_meter->Collect( + metrics::NAME_KV_FLUSH_ROWS_TOTAL, + write_batch.Count(), + "base"); + } + write_batch.Clear(); + write_batch_size = 0; + } + } + } + } + + if (write_batch_size > 0) + { + auto status = db->Write(write_options, &write_batch); + if (!status.ok()) + { + LOG(ERROR) + << "PutAll end failed " + << ", thread id: " << this_id + << ", result:" << static_cast(status.ok()) + << ", batch size:" << batch.size() + << ", error: " << status.ToString() + << ", error code: " << status.code(); + callback.Notify(false); + return; + } + if (metrics::enable_kv_metrics) + { + metrics::kv_meter->Collect( + metrics::NAME_KV_FLUSH_ROWS_TOTAL, + write_batch.Count(), + "base"); + } + } + callback.Notify(true); + }); + + callback.Wait(yield_fptr, resume_fptr); + return callback.success_; + } + std::shared_lock db_lk(db_mux_); auto db = GetDBPtr(); if (!db) @@ -545,8 +757,56 @@ bool RocksDBHandler::PutAll( return true; } -bool RocksDBHandler::PersistKV(const std::vector &kv_table_names) +bool RocksDBHandler::PersistKV(const std::vector &kv_table_names, + const std::function *yield_fptr, + const std::function *resume_fptr) { + if (yield_fptr != nullptr && resume_fptr != nullptr) + { + RocksDBWriteSyncCallback callback; + callback.Reset(); + callback.SetCoroCallbacks(yield_fptr, resume_fptr); + + query_worker_pool_->SubmitWork( + [this, &kv_table_names, &callback](size_t) + { + std::shared_lock db_lk(db_mux_); + auto db = GetDBPtr(); + if (!db) + { + callback.Notify(false); + return; + } + for (const std::string &kv_cf_name : kv_table_names) + { + rocksdb::ColumnFamilyHandle *cfh = + GetColumnFamilyHandler(kv_cf_name); + if (cfh == nullptr) + { + LOG(ERROR) << "Failed to get column family, cf name: " + << kv_cf_name; + callback.Notify(false); + return; + } + rocksdb::FlushOptions flush_options; + flush_options.allow_write_stall = true; + flush_options.wait = true; + auto status = db->Flush(flush_options, cfh); + if (!status.ok()) + { + LOG(ERROR) << "Unable to flush db with error: " + << status.ToString(); + callback.Notify(false); + return; + } + } + callback.Notify(true); + }); + + callback.Wait(yield_fptr, resume_fptr); + return callback.success_; + } + std::shared_lock db_lk(db_mux_); auto db = GetDBPtr(); if (!db) @@ -578,7 +838,7 @@ bool RocksDBHandler::PersistKV(const std::vector &kv_table_names) } return true; -}; +} void RocksDBHandler::UpsertTableInternal( txservice::TableEngine table_engine, @@ -1512,8 +1772,14 @@ bool RocksDBHandler::UpdateRangeSlices( } bool RocksDBHandler::UpdateRangeSlices( - const std::vector &update_range_slice_reqs) + const std::vector &update_range_slice_reqs, + const std::function *yield_fptr, + const std::function *resume_fptr, + const std::function *sync_yield_fptr) { + (void) yield_fptr; + (void) resume_fptr; + (void) sync_yield_fptr; LOG(ERROR) << "RocksDBHandler::UpdateRangeSlices not implemented"; // Not implemented assert(false); @@ -1663,8 +1929,12 @@ std::string RocksDBHandler::CreateNewKVCatalogInfo( bool RocksDBHandler::PutArchivesAll( std::unordered_map>> - &batch) + &batch, + const std::function *yield_fptr, + const std::function *resume_fptr) { + (void) yield_fptr; + (void) resume_fptr; LOG(ERROR) << "RocksDBHandler::PutArchivesAll not implemented"; // Not implemented assert(false); @@ -1676,8 +1946,12 @@ bool RocksDBHandler::PutArchivesAll( bool RocksDBHandler::CopyBaseToArchive( std::unordered_map>> - &batch) + &batch, + const std::function *yield_fptr, + const std::function *resume_fptr) { + (void) yield_fptr; + (void) resume_fptr; LOG(ERROR) << "RocksDBHandler::CopyBaseToArchive not implemented"; // Not implemented assert(false); diff --git a/store_handler/rocksdb_handler.h b/store_handler/rocksdb_handler.h index ed441153..af59eaac 100644 --- a/store_handler/rocksdb_handler.h +++ b/store_handler/rocksdb_handler.h @@ -272,10 +272,13 @@ class RocksDBHandler : public txservice::store::DataStoreHandler * @param node_group * @return whether all entries are written to data store successfully */ - bool PutAll(std::unordered_map< - std::string_view, - std::vector>> &batch) - override; + bool PutAll( + std::unordered_map< + std::string_view, + std::vector>> &batch, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr, + const std::function *sync_yield_fptr = nullptr) override; /** * @brief indicate end of flush entries in a single ckpt for \@param @@ -285,7 +288,9 @@ class RocksDBHandler : public txservice::store::DataStoreHandler * @param node_group * @return whether all entries are written to data store successfully */ - bool PersistKV(const std::vector &kv_table_names) override; + bool PersistKV(const std::vector &kv_table_names, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr) override; bool NeedPersistKV() override { @@ -411,8 +416,12 @@ class RocksDBHandler : public txservice::store::DataStoreHandler int32_t partition_id, uint64_t range_version) override; - bool UpdateRangeSlices(const std::vector - &update_range_slice_reqs) override; + bool UpdateRangeSlices( + const std::vector + &update_range_slice_reqs, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr, + const std::function *sync_yield_fptr = nullptr) override; bool UpsertRanges(const txservice::TableName &table_name, std::vector range_info, @@ -477,18 +486,21 @@ class RocksDBHandler : public txservice::store::DataStoreHandler * @brief Write batch historical versions into DataStore. * */ - bool PutArchivesAll(std::unordered_map< - std::string_view, - std::vector>> - &batch) override; + bool PutArchivesAll( + std::unordered_map< + std::string_view, + std::vector>> &batch, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr) override; /** * @brief Copy record from base/sk table to mvcc_archives. */ bool CopyBaseToArchive( std::unordered_map< std::string_view, - std::vector>> &batch) - override; + std::vector>> &batch, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr) override; /** * @brief Get the latest visible(commit_ts <= upper_bound_ts) diff --git a/tx_service/include/cc/catalog_cc_map.h b/tx_service/include/cc/catalog_cc_map.h index a66f5699..69613ec1 100644 --- a/tx_service/include/cc/catalog_cc_map.h +++ b/tx_service/include/cc/catalog_cc_map.h @@ -1289,8 +1289,8 @@ class CatalogCcMap { LOG(INFO) << "ReplayLogCc, node_group(#" << req.NodeGroupId() << ") term < 0, tx:" << req.Txn(); - req.Result()->SetError(CcErrorCode::REQUESTED_NODE_NOT_LEADER); - return false; + req.AbortCcRequest(CcErrorCode::REQUESTED_NODE_NOT_LEADER); + return true; } ::txlog::SchemaOpMessage schema_op_msg; @@ -1351,7 +1351,7 @@ class CatalogCcMap << ", catalog entry of the same or higher " "version exists, stop replaying this schema op"; req.SetFinish(); - return false; + return true; } catalog_entry = new_catalog_entry; } @@ -1424,8 +1424,8 @@ class CatalogCcMap else if (!need_to_upsert_kv_table && upsert_kv_err_code.second != CcErrorCode::NO_ERROR) { - req.Result()->SetError(upsert_kv_err_code.second); - return false; + req.AbortCcRequest(upsert_kv_err_code.second); + return true; } auto [success, new_catalog_entry] = shard_->CreateCatalog( @@ -1447,7 +1447,7 @@ class CatalogCcMap << ", catalog entry of the same or higher version " "exists, stop replaying this schema op"; req.SetFinish(); - return false; + return true; } catalog_entry = new_catalog_entry; } @@ -1735,7 +1735,7 @@ class CatalogCcMap } } - return false; + return true; } bool Execute(BroadcastStatisticsCc &req) override diff --git a/tx_service/include/cc/cc_req_misc.h b/tx_service/include/cc/cc_req_misc.h index eedae7e7..51f4d050 100644 --- a/tx_service/include/cc/cc_req_misc.h +++ b/tx_service/include/cc/cc_req_misc.h @@ -29,11 +29,13 @@ #include #include #include +#include #include #include #include #include #include +#include #include #include #include @@ -876,6 +878,13 @@ struct WaitableCc : public RunOnTxProcessorCc error_code_ = CcErrorCode::NO_ERROR; } + void SetCoroCallbacks(const std::function *yield_fn, + const std::function *resume_fn) + { + yield_fn_ = yield_fn; + resume_fn_ = resume_fn; + } + void Wait() { std::unique_lock lk(mux_); @@ -885,6 +894,25 @@ struct WaitableCc : public RunOnTxProcessorCc } } + void Wait(const std::function *yield_fn, + const std::function *resume_fn) + { + if (yield_fn == nullptr || resume_fn == nullptr) + { + Wait(); + return; + } + std::unique_lock lk(mux_); + while (unfinished_cnt_) + { + waiting_.store(true, std::memory_order_release); + lk.unlock(); + (*yield_fn)(); + lk.lock(); + waiting_.store(false, std::memory_order_release); + } + } + bool IsFinished() const { std::lock_guard lk(mux_); @@ -910,7 +938,18 @@ struct WaitableCc : public RunOnTxProcessorCc error_code_ = error_code; if (unfinished_cnt_ == 0) { - cv_.notify_one(); + if (resume_fn_ != nullptr && + waiting_.load(std::memory_order_acquire)) + { + waiting_.store(false, std::memory_order_release); + auto *fn = resume_fn_; + lk.unlock(); + (*fn)(); + } + else if (resume_fn_ == nullptr) + { + cv_.notify_one(); + } } } @@ -922,7 +961,18 @@ struct WaitableCc : public RunOnTxProcessorCc error_code_ = CcErrorCode::NO_ERROR; if (--unfinished_cnt_ == 0) { - cv_.notify_one(); + if (resume_fn_ != nullptr && + waiting_.load(std::memory_order_acquire)) + { + waiting_.store(false, std::memory_order_release); + auto *fn = resume_fn_; + lk.unlock(); + (*fn)(); + } + else if (resume_fn_ == nullptr) + { + cv_.notify_one(); + } } } return false; @@ -944,6 +994,11 @@ struct WaitableCc : public RunOnTxProcessorCc uint32_t unfinished_cnt_{0}; CcErrorCode error_code_; + + // Coroutine yield/resume support + const std::function *yield_fn_{nullptr}; + const std::function *resume_fn_{nullptr}; + std::atomic waiting_{false}; }; struct UpdateCceCkptTsCc : public CcRequestBase { @@ -990,13 +1045,31 @@ struct UpdateCceCkptTsCc : public CcRequestBase bool Execute(CcShard &ccs) override; + void SetCoroCallbacks(const std::function *yield_fn, + const std::function *resume_fn) + { + yield_fn_ = yield_fn; + resume_fn_ = resume_fn; + } + void SetFinished() { - std::lock_guard lk(mux_); + std::unique_lock lk(mux_); unfinished_core_cnt_--; if (unfinished_core_cnt_ == 0) { - cv_.notify_one(); + if (resume_fn_ != nullptr && + waiting_.load(std::memory_order_acquire)) + { + waiting_.store(false, std::memory_order_release); + auto *fn = resume_fn_; + lk.unlock(); + (*fn)(); + } + else if (resume_fn_ == nullptr) + { + cv_.notify_one(); + } } } @@ -1005,7 +1078,26 @@ struct UpdateCceCkptTsCc : public CcRequestBase std::unique_lock lk(mux_); while (unfinished_core_cnt_ > 0) { - cv_.wait_for(lk, 10000); + cv_.wait_for(lk, 10000L); // timeout_us, preserve original value + } + } + + void Wait(const std::function *yield_fn, + const std::function *resume_fn) + { + if (yield_fn == nullptr || resume_fn == nullptr) + { + Wait(); + return; + } + std::unique_lock lk(mux_); + while (unfinished_core_cnt_ > 0) + { + waiting_.store(true, std::memory_order_release); + lk.unlock(); + (*yield_fn)(); + lk.lock(); + waiting_.store(false, std::memory_order_release); } } @@ -1015,6 +1107,12 @@ struct UpdateCceCkptTsCc : public CcRequestBase return cce_entries_; } + bool IsFinished() const + { + std::lock_guard lk(mux_); + return unfinished_core_cnt_ == 0; + } + private: absl::flat_hash_map> &cce_entries_; // key: core_idx, value: entry_index @@ -1024,8 +1122,13 @@ struct UpdateCceCkptTsCc : public CcRequestBase NodeGroupId node_group_id_; int64_t term_; TableName table_name_; - bthread::Mutex mux_; + mutable bthread::Mutex mux_; bthread::ConditionVariable cv_; + + // Coroutine yield/resume support + const std::function *yield_fn_{nullptr}; + const std::function *resume_fn_{nullptr}; + std::atomic waiting_{false}; }; struct WaitNoNakedBucketRefCc : public CcRequestBase diff --git a/tx_service/include/cc/cc_request.h b/tx_service/include/cc/cc_request.h index 47ebfce8..388663c9 100644 --- a/tx_service/include/cc/cc_request.h +++ b/tx_service/include/cc/cc_request.h @@ -8865,11 +8865,7 @@ struct ScanSliceDeltaSizeCcForRangePartition : public CcRequestBase struct ScanDeltaSizeCcForHashPartition : public CcRequestBase { - static constexpr size_t ScanBatchSize = 128; - ScanDeltaSizeCcForHashPartition(const TableName &table_name, - uint64_t last_datasync_ts, - uint64_t scan_ts, uint64_t ng_id, int64_t ng_term, uint64_t txn, @@ -8877,8 +8873,6 @@ struct ScanDeltaSizeCcForHashPartition : public CcRequestBase : table_name_(table_name), node_group_id_(ng_id), node_group_term_(ng_term), - last_datasync_ts_(last_datasync_ts), - scan_ts_(scan_ts), schema_version_(schema_version) { tx_number_ = txn; @@ -8970,44 +8964,27 @@ struct ScanDeltaSizeCcForHashPartition : public CcRequestBase return node_group_term_; } - uint64_t LastDataSyncTs() const - { - return last_datasync_ts_; - } - - uint64_t ScanTs() const - { - return scan_ts_; - } - void AbortCcRequest(CcErrorCode err_code) override { assert(err_code != CcErrorCode::NO_ERROR); SetError(err_code); } - TxKey &PausedKey() + void SetKeyCounts(const size_t data_key_count, const size_t dirty_key_count) { - return pause_key_; - } - - void UpdateKeyCount(const bool need_ckpt) - { - ++scanned_key_count_; - if (need_ckpt) - { - ++updated_key_count_; - } + data_key_count_ = data_key_count; + dirty_key_count_ = dirty_key_count; } size_t UpdatedMemory() const { - const uint64_t scanned = scanned_key_count_; - if (scanned == 0) + if (data_key_count_ == 0) return 0; + // Use cc map's data_key_count_ and dirty_key_count_ for estimation. // integer math with rounding up to avoid systematic underestimation return static_cast( - (updated_key_count_ * memory_usage_ + scanned - 1) / scanned); + (dirty_key_count_ * memory_usage_ + data_key_count_ - 1) / + data_key_count_); } void SetMemoryUsage(const uint64_t memory) @@ -9024,20 +9001,11 @@ struct ScanDeltaSizeCcForHashPartition : public CcRequestBase const TableName &table_name_; uint32_t node_group_id_; int64_t node_group_term_; - // It is used as a hint to decide if a page has dirty data since last round - // of checkpoint. It is guaranteed that all entries committed before this ts - // are synced into data store. - uint64_t last_datasync_ts_; - // Target ts. Collect all data changes committed before this ts into data - // sync vec. - uint64_t scan_ts_; - // Position that we left off during last round of scan. - TxKey pause_key_; uint64_t schema_version_; - // Number of keys scanned / updated, and per-core memory usage. - uint64_t scanned_key_count_{0}; - uint64_t updated_key_count_{0}; + // From cc map: total data keys and dirty keys needing checkpoint. + size_t data_key_count_{0}; + size_t dirty_key_count_{0}; uint64_t memory_usage_{0}; CcErrorCode err_{CcErrorCode::NO_ERROR}; diff --git a/tx_service/include/cc/local_cc_shards.h b/tx_service/include/cc/local_cc_shards.h index 316d8042..c28229a9 100644 --- a/tx_service/include/cc/local_cc_shards.h +++ b/tx_service/include/cc/local_cc_shards.h @@ -23,6 +23,9 @@ #include #include +#include +#include +#include #include #include #include @@ -85,6 +88,19 @@ class SkGenerator; class UploadBatchSlicesClosure; struct FlushDataTask; +struct CoroCtx +{ + ~CoroCtx() + { + } + + boost::context::continuation coro_; + std::unique_ptr task_; + std::function sync_yield_func; + std::function yield_fn; + std::function resume_fn; +}; + struct DataMigrationStatus { public: @@ -457,8 +473,13 @@ class LocalCcShards void InitializeHashPartitionCkptHeap() { - hash_partition_ckpt_heap_ = mi_heap_new(); - hash_partition_main_thread_id_ = mi_thread_id(); + std::unique_lock lk(hash_partition_ckpt_heap_mux_); + if (!hash_partition_ckpt_heap_) + { + hash_partition_main_thread_id_ = mi_thread_id(); + hash_partition_ckpt_heap_ = mi_heap_new(); + } + #if defined(WITH_JEMALLOC) // create hash partition ckpt arena size_t sz = sizeof(uint32_t); @@ -559,6 +580,7 @@ class LocalCcShards << (bool) (static_cast(allocated) >= range_slice_memory_limit_); #else + std::unique_lock heap_lk(table_ranges_heap_mux_); bool is_override_thd = mi_is_override_thread(); mi_threadid_t prev_thd = @@ -858,9 +880,27 @@ class LocalCcShards new_range_entries.push_back(new_range); } } - if (TableRangesMemoryFull()) { - KickoutRangeSlices(); + std::unique_lock heap_lk(table_ranges_heap_mux_); + bool is_override_thd = mi_is_override_thread(); + mi_threadid_t prev_thd = + mi_override_thread(GetTableRangesHeapThreadId()); + mi_heap_t *prev_heap = mi_heap_set_default(GetTableRangesHeap()); + bool range_slice_mem_full = TableRangesMemoryFull(); + mi_heap_set_default(prev_heap); + if (is_override_thd) + { + mi_override_thread(prev_thd); + } + else + { + mi_restore_default_thread_id(); + } + heap_lk.unlock(); + if (range_slice_mem_full) + { + KickoutRangeSlices(); + } } return new_range_entries; } @@ -970,6 +1010,7 @@ class LocalCcShards } mi_heap_set_default(prev_heap); + if (is_override_thd) { mi_override_thread(prev_thd); @@ -1017,6 +1058,8 @@ class LocalCcShards range_entry->UpdateRangeEntry(version, std::move(range_slices)); } + mi_heap_set_default(prev_heap); + if (is_override_thd) { mi_override_thread(prev_thd); @@ -1025,7 +1068,6 @@ class LocalCcShards { mi_restore_default_thread_id(); } - mi_heap_set_default(prev_heap); #if defined(WITH_JEMALLOC) JemallocArenaSwitcher::SwitchToArena(prev_arena_id); @@ -1884,6 +1926,8 @@ class LocalCcShards } void TimerRun(); + void TableRangesHeapThreadRun(); + void HashPartitionCkptHeapThreadRun(); // Internal interface that exposes non const return type and does // not acquire mutex lock. TableRangeEntry *GetTableRangeEntryInternal( @@ -1973,7 +2017,10 @@ class LocalCcShards std::string_view, std::vector>> &flush_task_entries, - DataSyncTask::CkptErrorCode ckpt_err); + DataSyncTask::CkptErrorCode ckpt_err, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr, + const std::function *sync_yield_fptr = nullptr); const uint32_t node_id_; // Native node group @@ -1988,6 +2035,22 @@ class LocalCcShards std::condition_variable timer_terminate_cv_; // std::atomic timer_terminate_; + // Background thread that initializes table_ranges_heap. After init, it + // sleeps until Shutdown. + std::thread table_ranges_heap_thd_; + bool table_ranges_heap_ready_{false}; + bool table_ranges_heap_terminate_{false}; + std::mutex table_ranges_heap_thd_mux_; + std::condition_variable table_ranges_heap_thd_cv_; + + // Background thread that initializes hash_partition_ckpt_heap. After init, + // it sleeps until Shutdown. + std::thread hash_partition_ckpt_heap_thd_; + bool hash_partition_ckpt_heap_ready_{false}; + bool hash_partition_ckpt_heap_terminate_{false}; + std::mutex hash_partition_ckpt_heap_thd_mux_; + std::condition_variable hash_partition_ckpt_heap_thd_cv_; + // When ccshard is full and no ccentry can be kicked-out, it will notify // checkpointer to do checkpoint and set flag is_wait_ckpt_ to true. // Subsequent ccrequest is able to skip checking freeable ccentry when @@ -2469,13 +2532,24 @@ class LocalCcShards const TxKey *end_key, bool flush_res); - bool UpdateStoreSlices(std::vector &task); + bool UpdateStoreSlices( + std::vector &task, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr, + const std::function *sync_yield_fptr = nullptr); bool GetNextRangePartitionId(const TableName &tablename, const TableSchema *table_schema, uint32_t range_cnt, int32_t &out_next_partition_id); + /** + * @brief Map a data_sync_worker index to the fixed flush_data_worker index. + * Used when data_sync_worker count != flush_data_worker count so that each + * data_sync_worker consistently targets one flush_data_worker. + */ + size_t DataSyncWorkerToFlushDataWorker(size_t data_sync_worker_id) const; + /** * @brief Add a flush task entry to the flush task. If the there's no * pending flush task, create a new flush task and add the entry to it. @@ -2487,16 +2561,39 @@ class LocalCcShards WorkerThreadContext flush_data_worker_ctx_; - // The flush task that has not reached the max pending flush size. - // New flush task entry will be added to this buffer. This task will - // be appended to pending_flush_work_ when it reaches the max pending flush - // size, which will then be processed by flush data worker. - FlushDataTask cur_flush_buffer_; - // Flush task queue for flush data worker to process. - std::deque> pending_flush_work_; + // Per-worker flush buffers. Each DataSyncWorker has its own buffer. + // New flush task entry will be added to the corresponding buffer. This task + // will be appended to pending_flush_work_[worker_idx] when it reaches the + // max pending flush size, which will then be processed by the corresponding + // flush data worker. Store as pointers because FlushDataTask contains a + // bthread::Mutex and is non-movable/non-copyable, which cannot be stored + // directly in a vector that may reallocate. + std::vector> cur_flush_buffers_; + // Per-worker flush task queues. Each FlushDataWorker processes its + // corresponding queue. + std::vector>> pending_flush_work_; + // Per-worker queues of coroutine contexts ready to resume (resume_fn or + // sync_yield_func). + std::vector>> resume_queue_; + +#ifndef NDEBUG + boost::context::protected_fixedsize_stack flush_coro_stack_allocator_{ + 128 * 1024}; // 128KB, guard page for stack overflow detection +#else + boost::context::pooled_fixedsize_stack flush_coro_stack_allocator_{ + 128 * 1024}; // 128KB for FlushData call chain +#endif + + void FlushDataWorker(size_t worker_idx); + void FlushData(std::unique_lock &flush_worker_lk, + size_t worker_idx); - void FlushDataWorker(); - void FlushData(std::unique_lock &flush_worker_lk); + bool ShouldYieldFlushData(size_t worker_idx); + void FlushDataImpl(FlushDataTask *cur_work, + size_t worker_idx, + const std::function &sync_yield_func, + const std::function &yield_fn, + const std::function &resume_fn); // Memory controller for data sync. DataSyncMemoryController data_sync_mem_controller_; diff --git a/tx_service/include/cc/range_cc_map.h b/tx_service/include/cc/range_cc_map.h index c6a7767e..3a518093 100644 --- a/tx_service/include/cc/range_cc_map.h +++ b/tx_service/include/cc/range_cc_map.h @@ -826,7 +826,7 @@ class RangeCcMap : public TemplateCcMap int64_t ng_term = Sharder::Instance().CandidateLeaderTerm(group_id); if (ng_term < 0) { - req.Result()->SetError(CcErrorCode::REQUESTED_NODE_NOT_LEADER); + req.AbortCcRequest(CcErrorCode::REQUESTED_NODE_NOT_LEADER); return true; } diff --git a/tx_service/include/cc/template_cc_map.h b/tx_service/include/cc/template_cc_map.h index 1870cbe5..6d69eb0f 100644 --- a/tx_service/include/cc/template_cc_map.h +++ b/tx_service/include/cc/template_cc_map.h @@ -7981,132 +7981,11 @@ class TemplateCcMap : public CcMap return false; } - auto &paused_key = req.PausedKey(); - - auto deduce_iterator = [this](const KeyT &search_key) -> Iterator - { - Iterator it; - std::pair search_pair = - ForwardScanStart(search_key, true); - it = search_pair.first; - if (search_pair.second == ScanType::ScanGap) - { - ++it; - } - return it; - }; - - auto next_page_it = [this](Iterator &end_it) -> Iterator - { - Iterator it = end_it; - if (it != End()) - { - CcPage *ccp = - it.GetPage(); - assert(ccp != nullptr); - if (ccp->next_page_ == PagePosInf()) - { - it = End(); - } - else - { - it = Iterator(ccp->next_page_, 0, &neg_inf_); - } - } - return it; - }; - - Iterator key_it; - Iterator req_end_it; - - // The key iterator. - const KeyT *const search_start_key = - paused_key.GetKey() == nullptr ? KeyT::NegativeInfinity() - : paused_key.GetKey(); - key_it = deduce_iterator(*search_start_key); - - // The request end iterator - req_end_it = deduce_iterator(*KeyT::PositiveInfinity()); - - // Since we might skip the page that end_it is on if it's not updated - // since last ckpt, it might skip end_it. If the last page is skipped it - // will be set as the first entry on the next page. Also check if - // (key_it == end_next_page_it). - Iterator req_end_next_page_it = next_page_it(req_end_it); - - // The current slice end iterator - Iterator slice_end_it = req_end_it; - Iterator slice_end_next_page_it = req_end_next_page_it; - - // ScanSliceDeltaSizeCcForRangePartition is running on TxProcessor - // thread. To avoid blocking other transaction for a long time, we only - // process ScanBatchSize number of keys in each round. - for (size_t scan_cnt = 0; - scan_cnt < ScanDeltaSizeCcForHashPartition::ScanBatchSize && - key_it != req_end_it && key_it != req_end_next_page_it; - ++scan_cnt) - { - CcEntry *cce = - key_it->second; - CcPage *ccp = - key_it.GetPage(); - assert(ccp); - - if (ccp->last_dirty_commit_ts_ <= req.LastDataSyncTs()) - { - // Skip the pages that have no updates since last data sync. - if (ccp->next_page_ == PagePosInf()) - { - key_it = End(); - } - else - { - key_it = Iterator(ccp->next_page_, 0, &neg_inf_); - } - - // Check the slice iterator. - if (key_it == slice_end_it || key_it == slice_end_next_page_it) - { - // Reset the slice end iterator - slice_end_it = req_end_it; - slice_end_next_page_it = next_page_it(slice_end_it); - } - continue; - } - - const uint64_t commit_ts = cce->CommitTs(); - - // The commit_ts <= 1 means the key is non-existed or a new inserted - // key that the tx has not finished post-processing. - req.UpdateKeyCount(cce->NeedCkpt() && commit_ts > 1 && - commit_ts <= req.ScanTs()); - - // Forward key iterator - ++key_it; - - if (key_it == slice_end_it) - { - // Update the end it. - slice_end_it = req_end_it; - slice_end_next_page_it = next_page_it(slice_end_it); - } - } - - if (key_it == req_end_it || key_it == req_end_next_page_it) - { - int64_t allocated, committed; - mi_thread_stats(&allocated, &committed); - req.SetMemoryUsage(allocated); - req.SetFinish(); - } - else - { - paused_key = key_it->first->CloneTxKey(); - shard_->EnqueueLowPriorityCcRequest(&req); - } - - // Access ScanSliceDeltaSizeCcForRangePartition member variable is - // unsafe after SetFinished(). + req.SetKeyCounts(data_key_count_, dirty_data_key_count_); + int64_t allocated, committed; + mi_thread_stats(&allocated, &committed); + req.SetMemoryUsage(allocated); + req.SetFinish(); return false; } @@ -8243,10 +8122,35 @@ class TemplateCcMap : public CcMap return true; } + void AdjustDataKeyStats(int64_t size_delta, int64_t dirty_delta) + { + if (table_name_.IsMeta()) + { + return; + } + + if (size_delta != 0) + { + assert(size_delta >= 0 || + data_key_count_ >= static_cast(-size_delta)); + data_key_count_ = static_cast( + static_cast(data_key_count_) + size_delta); + } + + if (dirty_delta != 0) + { + assert(dirty_delta >= 0 || + dirty_data_key_count_ >= static_cast(-dirty_delta)); + dirty_data_key_count_ = static_cast( + static_cast(dirty_data_key_count_) + dirty_delta); + } + } + void OnEntryFlushed(bool was_dirty, bool is_persistent) override { if (was_dirty && is_persistent) { + AdjustDataKeyStats(0, -1); shard_->AdjustDataKeyStats(table_name_, 0, -1); } } @@ -8266,6 +8170,7 @@ class TemplateCcMap : public CcMap page_it = ccmp_.find(ccpage->FirstKey()); } const bool was_dirty = cce->IsDirty(); + AdjustDataKeyStats(-1, was_dirty ? -1 : 0); shard_->AdjustDataKeyStats(table_name_, -1, was_dirty ? -1 : 0); ccpage->Remove(cce); @@ -8342,6 +8247,8 @@ class TemplateCcMap : public CcMap if (total_freed > 0 || dirty_freed > 0) { + AdjustDataKeyStats(-static_cast(total_freed), + -static_cast(dirty_freed)); shard_->AdjustDataKeyStats(table_name_, -static_cast(total_freed), -static_cast(dirty_freed)); @@ -8383,6 +8290,8 @@ class TemplateCcMap : public CcMap cce->ClearLocks(*shard_, cc_ng_id_); } + AdjustDataKeyStats(-static_cast(page->Size()), + -static_cast(dirty_freed)); shard_->AdjustDataKeyStats(table_name_, -static_cast(page->Size()), -static_cast(dirty_freed)); @@ -8800,6 +8709,7 @@ class TemplateCcMap : public CcMap { if (!was_dirty && cce->IsDirty()) { + AdjustDataKeyStats(0, +1); shard_->AdjustDataKeyStats(table_name_, 0, +1); } } @@ -9671,6 +9581,7 @@ class TemplateCcMap : public CcMap { normal_obj_sz_ += normal_rec_change; } + AdjustDataKeyStats(+(end_idx - first_index), 0); shard_->AdjustDataKeyStats( table_name_, +(end_idx - first_index), 0); @@ -9709,6 +9620,7 @@ class TemplateCcMap : public CcMap assert(new_key_cnt > 0); // Update CCMap size + AdjustDataKeyStats(+new_key_cnt, 0); shard_->AdjustDataKeyStats(table_name_, +new_key_cnt, 0); } else @@ -9802,6 +9714,7 @@ class TemplateCcMap : public CcMap assert(new_key_cnt > 0); // Update Ccmap size + AdjustDataKeyStats(+new_key_cnt, 0); shard_->AdjustDataKeyStats(table_name_, +new_key_cnt, 0); } @@ -10081,6 +9994,7 @@ class TemplateCcMap : public CcMap // update lru list shard_->UpdateLruList(target_page, true); + AdjustDataKeyStats(+1, 0); shard_->AdjustDataKeyStats(table_name_, +1, 0); return Iterator(target_page, idx_in_page, &neg_inf_); @@ -11254,6 +11168,9 @@ class TemplateCcMap : public CcMap { normal_obj_sz_ -= clean_guard->CleanObjectCount(); } + AdjustDataKeyStats( + -static_cast(clean_guard->FreedCount()), + -static_cast(clean_guard->DirtyFreedCount())); shard_->AdjustDataKeyStats( table_name_, -static_cast(clean_guard->FreedCount()), @@ -11713,6 +11630,8 @@ class TemplateCcMap : public CcMap TemplateCcMapSamplePool *sample_pool_; size_t normal_obj_sz_{ 0}; // The count of all normal status objects, only used for redis + size_t data_key_count_{0}; + size_t dirty_data_key_count_{0}; }; template diff --git a/tx_service/include/fault/log_replay_service.h b/tx_service/include/fault/log_replay_service.h index 6dcc9677..8d69f06d 100644 --- a/tx_service/include/fault/log_replay_service.h +++ b/tx_service/include/fault/log_replay_service.h @@ -29,6 +29,7 @@ #include #include +#include #include #include #include @@ -270,8 +271,9 @@ class RecoveryService : public brpc::StreamInputHandler, LocalCcShards &local_shards_; // Each ConnectionInfo is uniquely identified by - // pair. - std::unordered_map inbound_connections_; + // pair. unique_ptr ensures pointer stability across map rehashing. + std::unordered_map> + inbound_connections_; int active_stream_cnt_ = 0; bthread::Mutex inbound_mux_; bthread::ConditionVariable inbound_cv_; diff --git a/tx_service/include/range_record.h b/tx_service/include/range_record.h index 701b904d..a7af99d6 100644 --- a/tx_service/include/range_record.h +++ b/tx_service/include/range_record.h @@ -542,6 +542,8 @@ struct TemplateTableRangeEntry : public TableRangeEntry { } + ~TemplateTableRangeEntry() = default; + void UpdateRangeEntry(uint64_t version_ts, std::unique_ptr> slices) { diff --git a/tx_service/include/store/data_store_handler.h b/tx_service/include/store/data_store_handler.h index abb7b0ab..e7f8efb4 100644 --- a/tx_service/include/store/data_store_handler.h +++ b/tx_service/include/store/data_store_handler.h @@ -84,10 +84,15 @@ class DataStoreHandler * data_sync_vec_ in each flush task entry. * @return whether all entries are written to data store successfully */ - virtual bool PutAll(std::unordered_map< - std::string_view, - std::vector>> - &flush_task) = 0; + + virtual bool PutAll( + std::unordered_map< + std::string_view, + std::vector>> + &flush_task, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr, + const std::function *sync_yield_fptr = nullptr) = 0; /** * @brief indicate end of flush entries in a single ckpt for \@param batch @@ -97,8 +102,12 @@ class DataStoreHandler * @param node_group * @return whether all entries are written to data store successfully */ - virtual bool PersistKV(const std::vector &kv_table_names) + virtual bool PersistKV(const std::vector &kv_table_names, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr) { + (void) yield_fptr; + (void) resume_fptr; return true; } @@ -227,8 +236,14 @@ class DataStoreHandler } virtual bool UpdateRangeSlices( - const std::vector &update_range_slice_reqs) + const std::vector &update_range_slice_reqs, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr, + const std::function *sync_yield_fptr = nullptr) { + (void) yield_fptr; + (void) resume_fptr; + (void) sync_yield_fptr; return false; } @@ -285,7 +300,9 @@ class DataStoreHandler std::unordered_map< std::string_view, std::vector>> - &flush_task) = 0; + &flush_task, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr) = 0; /** * @brief Copy record from base/sk table to mvcc_archives. */ @@ -293,7 +310,9 @@ class DataStoreHandler std::unordered_map< std::string_view, std::vector>> - &flush_task) = 0; + &flush_task, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr) = 0; /** * @brief Get the latest visible(commit_ts <= upper_bound_ts) historical diff --git a/tx_service/include/store/int_mem_store.h b/tx_service/include/store/int_mem_store.h index eb3dffa6..39c12c8f 100644 --- a/tx_service/include/store/int_mem_store.h +++ b/tx_service/include/store/int_mem_store.h @@ -55,10 +55,16 @@ class IntMemoryStore : public DataStoreHandler * @return whether all entries are written to data store successfully */ bool PutAll(std::unordered_map< - std::string_view, - std::vector>> - &flush_task) override - { + std::string_view, + std::vector>> + &flush_task, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr, + const std::function *sync_yield_fptr = nullptr) override + { + (void) yield_fptr; + (void) resume_fptr; + (void) sync_yield_fptr; assert(false); // for (const auto &ref : batch) // { @@ -243,11 +249,16 @@ class IntMemoryStore : public DataStoreHandler * @brief Write batch historical versions into DataStore. * */ - bool PutArchivesAll(std::unordered_map< - std::string_view, - std::vector>> - &flush_task) override + bool PutArchivesAll( + std::unordered_map< + std::string_view, + std::vector>> + &flush_task, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr) override { + (void) yield_fptr; + (void) resume_fptr; assert(false); return true; } @@ -258,8 +269,12 @@ class IntMemoryStore : public DataStoreHandler std::unordered_map< std::string_view, std::vector>> - &flush_task) override + &flush_task, + const std::function *yield_fptr = nullptr, + const std::function *resume_fptr = nullptr) override { + (void) yield_fptr; + (void) resume_fptr; assert(false); return true; } diff --git a/tx_service/src/cc/local_cc_shards.cpp b/tx_service/src/cc/local_cc_shards.cpp index 280719e4..a2fdedc1 100644 --- a/tx_service/src/cc/local_cc_shards.cpp +++ b/tx_service/src/cc/local_cc_shards.cpp @@ -132,18 +132,12 @@ LocalCcShards::LocalCcShards( data_sync_worker_ctx_(conf.at("core_num")), #ifdef EXT_TX_PROC_ENABLED - flush_data_worker_ctx_( - conf.at("core_num") >= 2 - ? std::min(conf.at("core_num") / 2, (uint32_t) 10) - : 1), + flush_data_worker_ctx_(1), #else - flush_data_worker_ctx_( - std::min(static_cast(conf.at("core_num")), 10)), + flush_data_worker_ctx_(1), #endif - - cur_flush_buffer_( - static_cast(MB(conf.at("node_memory_limit_mb")) * - (FLAGS_ckpt_buffer_ratio - 0.025))), + // cur_flush_buffers_ will be resized in StartBackgroudWorkers() + // when worker_num_ is determined data_sync_mem_controller_(static_cast( MB(conf.at("node_memory_limit_mb")) * FLAGS_ckpt_buffer_ratio)), statistics_worker_ctx_(1), @@ -184,10 +178,28 @@ LocalCcShards::LocalCcShards( timer_thd_ = std::thread([this] { TimerRun(); }); pthread_setname_np(timer_thd_.native_handle(), "tx_timer"); - // For mariadb, this thread is the main thread of the mariadb process. - InitializeTableRangesHeap(); + // Start background thread to initialize table_ranges_heap. offload heap + // init to avoid concurrent access to heap with main thread. + table_ranges_heap_thd_ = + std::thread([this] { TableRangesHeapThreadRun(); }); + pthread_setname_np(table_ranges_heap_thd_.native_handle(), + "table_ranges_heap"); + { + std::unique_lock lk(table_ranges_heap_thd_mux_); + table_ranges_heap_thd_cv_.wait( + lk, [this] { return table_ranges_heap_ready_; }); + } - InitializeHashPartitionCkptHeap(); + // Start background thread to initialize hash_partition_ckpt_heap. + hash_partition_ckpt_heap_thd_ = + std::thread([this] { HashPartitionCkptHeapThreadRun(); }); + pthread_setname_np(hash_partition_ckpt_heap_thd_.native_handle(), + "hash_part_ckpt_heap"); + { + std::unique_lock lk(hash_partition_ckpt_heap_thd_mux_); + hash_partition_ckpt_heap_thd_cv_.wait( + lk, [this] { return hash_partition_ckpt_heap_ready_; }); + } std::set ng_ids; for (auto &[ng_id, _] : *ng_configs) @@ -287,15 +299,44 @@ LocalCcShards::~LocalCcShards() } timer_thd_.join(); cc_shards_.clear(); + + { + std::scoped_lock lk(hash_partition_ckpt_heap_thd_mux_); + hash_partition_ckpt_heap_terminate_ = true; + hash_partition_ckpt_heap_thd_cv_.notify_one(); + } + hash_partition_ckpt_heap_thd_.join(); + + { + std::scoped_lock lk(table_ranges_heap_thd_mux_); + table_ranges_heap_terminate_ = true; + table_ranges_heap_thd_cv_.notify_one(); + } + table_ranges_heap_thd_.join(); } void LocalCcShards::StartBackgroudWorkers() { + // cur_flush_buffers_ sized by data_sync_worker_num (one buffer per data + // sync worker). pending_flush_work_ and resume_queue_ sized by flush worker + // num. + const size_t data_sync_worker_num = data_sync_worker_ctx_.worker_num_; + const uint64_t buffer_size = + data_sync_mem_controller_.FlushMemoryQuota() / data_sync_worker_num; + cur_flush_buffers_.clear(); + for (size_t i = 0; i < data_sync_worker_num; ++i) + { + cur_flush_buffers_.emplace_back( + std::make_unique(buffer_size)); + } + pending_flush_work_.resize(flush_data_worker_ctx_.worker_num_); + resume_queue_.resize(flush_data_worker_ctx_.worker_num_); + // Starts flush worker threads firstly. for (int id = 0; id < flush_data_worker_ctx_.worker_num_; id++) { std::thread &thd = flush_data_worker_ctx_.worker_thd_.emplace_back( - [this] { FlushDataWorker(); }); + [this, id] { FlushDataWorker(id); }); std::string thread_name = "flush_data_" + std::to_string(id); pthread_setname_np(thd.native_handle(), thread_name.c_str()); } @@ -407,6 +448,28 @@ void LocalCcShards::BindThreadToFastMetaDataShard(size_t shard_idx) } } +void LocalCcShards::TableRangesHeapThreadRun() +{ + InitializeTableRangesHeap(); + + std::unique_lock lk(table_ranges_heap_thd_mux_); + table_ranges_heap_ready_ = true; + table_ranges_heap_thd_cv_.notify_one(); + table_ranges_heap_thd_cv_.wait( + lk, [this] { return table_ranges_heap_terminate_; }); +} + +void LocalCcShards::HashPartitionCkptHeapThreadRun() +{ + InitializeHashPartitionCkptHeap(); + + std::unique_lock lk(hash_partition_ckpt_heap_thd_mux_); + hash_partition_ckpt_heap_ready_ = true; + hash_partition_ckpt_heap_thd_cv_.notify_one(); + hash_partition_ckpt_heap_thd_cv_.wait( + lk, [this] { return hash_partition_ckpt_heap_terminate_; }); +} + void LocalCcShards::TimerRun() { std::unique_lock lk(timer_terminate_mux_); @@ -2324,6 +2387,7 @@ std::pair LocalCcShards::PinStoreRange( range_table_name, cc_request, ng_id, ng_term, cc_shard); return {range_entry, nullptr}; } + return {range_entry, store_range}; } @@ -3147,7 +3211,10 @@ void LocalCcShards::PostProcessFlushTaskEntries( std::unordered_map>> &flush_task_entries, - DataSyncTask::CkptErrorCode ckpt_err) + DataSyncTask::CkptErrorCode ckpt_err, + const std::function *yield_fptr, + const std::function *resume_fptr, + const std::function *sync_yield_fptr) { assert(ckpt_err != DataSyncTask::CkptErrorCode::SCAN_ERROR); @@ -3306,7 +3373,8 @@ void LocalCcShards::PostProcessFlushTaskEntries( if (!pending_update_entries.empty()) { - bool success = UpdateStoreSlices(pending_update_entries); + bool success = UpdateStoreSlices( + pending_update_entries, yield_fptr, resume_fptr, sync_yield_fptr); for (auto &entry : pending_update_entries) { @@ -4715,16 +4783,57 @@ void LocalCcShards::DataSyncForHashPartition( assert(table_name.Type() == TableType::Primary); + ScanDeltaSizeCcForHashPartition scan_delta_size_cc( + table_name, + ng_id, + ng_term, + data_sync_txm->TxNumber(), + catalog_rec.Schema()->Version()); + + EnqueueToCcShard(worker_idx, &scan_delta_size_cc); + scan_delta_size_cc.Wait(); + + if (scan_delta_size_cc.IsError()) + { + LOG(ERROR) << "DataSync scan slice delta size failed on table " + << table_name.StringView() << " with error code: " + << static_cast(scan_delta_size_cc.ErrorCode()); + + AbortTx(data_sync_txm); + std::lock_guard task_worker_lk(data_sync_worker_ctx_.mux_); + data_sync_task_queue_[worker_idx].emplace_front(data_sync_task); + data_sync_worker_ctx_.cv_.notify_one(); + return; + } + auto core_number = cc_shards_.size(); - auto approximate_partition_number_this_node = std::max( - 1U, total_hash_partitions / Sharder::Instance().NodeGroupCount()); - auto approximate_partition_number_this_core = - std::max(1UL, approximate_partition_number_this_node / core_number); - const size_t flush_buffer_size = cur_flush_buffer_.GetFlushBufferSize(); - constexpr size_t default_data_file_size = 8ULL * 1024 * 1024; - const size_t scan_concurrency = - flush_buffer_size / - (default_data_file_size * approximate_partition_number_this_core); + const size_t partition_number = total_hash_partitions; + auto partition_number_this_core = + partition_number / core_number + + (worker_idx < partition_number % core_number); + std::vector partition_ids; + partition_ids.reserve(partition_number_this_core); + for (size_t i = 0; i < partition_number_this_core; ++i) + { + partition_ids.emplace_back(worker_idx + core_number * i); + } + assert(partition_number_this_core == partition_ids.size()); + const auto updated_memory = scan_delta_size_cc.UpdatedMemory(); + auto updated_memory_per_partition = + partition_number_this_core ? updated_memory / partition_number_this_core + : 0; + const size_t data_sync_worker_id = static_cast( + data_sync_task->id_ % data_sync_worker_ctx_.worker_num_); + const size_t flush_buffer_size = + cur_flush_buffers_[data_sync_worker_id]->GetFlushBufferSize(); + + const size_t partition_number_per_scan = + std::max(1UL, + updated_memory_per_partition != 0 + ? (flush_buffer_size / updated_memory_per_partition) + : partition_number_this_core); + const size_t scan_concurrency = core_number; + if (scan_concurrency > 0) { bool need_notify = scan_concurrency > @@ -4744,432 +4853,416 @@ void LocalCcShards::DataSyncForHashPartition( data_sync_task->flight_task_cnt_ += 1; } - HashPartitionDataSyncScanCc scan_cc(table_name, - data_sync_task->data_sync_ts_, - ng_id, - ng_term, - DATA_SYNC_SCAN_BATCH_SIZE, - data_sync_txm->TxNumber(), - data_sync_task->forward_cache_, - last_sync_ts, - data_sync_task->filter_lambda_, - catalog_rec.Schema()->Version()); - bool scan_data_drained = false; - - auto data_sync_vec = std::make_unique>(); - auto archive_vec = std::make_unique>(); - auto mv_base_vec = - std::make_unique>>(); - uint64_t vec_mem_usage = 0; - - // Note: `DataSyncScanCc` needs to ensure that no two ckpt_rec with the - // same Key can be generated. Our subsequent algorithms are based on - // this assumption. - - assert(worker_idx < cc_shards_.size()); + for (size_t i = 0; i < partition_number_this_core; + i += partition_number_per_scan) + { + size_t min_partition_id_this_scan = partition_ids[i]; + size_t max_partition_id_this_scan = + partition_ids[std::min(i + partition_number_per_scan, + partition_number_this_core) - + 1]; + std::function filter_lambda = + [min_partition_id_this_scan, + max_partition_id_this_scan, + &filter_func = + data_sync_task->filter_lambda_](const size_t hash_code) + { + return (hash_code % total_hash_partitions) >= + min_partition_id_this_scan && + (hash_code % total_hash_partitions) <= + max_partition_id_this_scan && + (!filter_func || filter_func(hash_code)); + }; - while (!scan_data_drained) - { - EnqueueLowPriorityCcRequestToShard(worker_idx, &scan_cc); - scan_cc.Wait(); + HashPartitionDataSyncScanCc scan_cc(table_name, + data_sync_task->data_sync_ts_, + ng_id, + ng_term, + DATA_SYNC_SCAN_BATCH_SIZE, + data_sync_txm->TxNumber(), + data_sync_task->forward_cache_, + last_sync_ts, + filter_lambda, + catalog_rec.Schema()->Version()); + bool scan_data_drained = false; + + auto data_sync_vec = std::make_unique>(); + auto archive_vec = std::make_unique>(); + auto mv_base_vec = + std::make_unique>>(); + uint64_t vec_mem_usage = 0; - if (scan_cc.IsError() && - scan_cc.ErrorCode() != CcErrorCode::LOG_NOT_TRUNCATABLE) - { - LOG(ERROR) << "DataSync scan failed on table " - << table_name.StringView() << " with error code: " - << static_cast(scan_cc.ErrorCode()); + // Note: `DataSyncScanCc` needs to ensure that no two ckpt_rec with the + // same Key can be generated. Our subsequent algorithms are based on + // this assumption. - PostProcessHashPartitionDataSyncTask( - std::move(data_sync_task), - data_sync_txm, - DataSyncTask::CkptErrorCode::SCAN_ERROR); + assert(worker_idx < cc_shards_.size()); - return; - } - if (scan_cc.ErrorCode() == CcErrorCode::LOG_NOT_TRUNCATABLE) + while (!scan_data_drained) { - data_sync_task->status_->SetEntriesSkippedAndNoTruncateLog(); - } + EnqueueLowPriorityCcRequestToShard(worker_idx, &scan_cc); + scan_cc.Wait(); - scan_data_drained = true; - - // Send cache to target node group if needed. - if (data_sync_task->forward_cache_) - { - std::shared_lock meta_lk(meta_data_mux_); - const auto bucket_infos = GetAllBucketInfosNoLocking(ng_id); - if (bucket_infos == nullptr) + if (scan_cc.IsError() && + scan_cc.ErrorCode() != CcErrorCode::LOG_NOT_TRUNCATABLE) { - // no longer node group owner, abort the task - LOG(ERROR) << "DataSync: Failed to get bucket infos for " - "ng#" - << ng_id; + LOG(ERROR) << "DataSync scan failed on table " + << table_name.StringView() << " with error code: " + << static_cast(scan_cc.ErrorCode()); + PostProcessHashPartitionDataSyncTask( std::move(data_sync_task), data_sync_txm, DataSyncTask::CkptErrorCode::SCAN_ERROR); + return; } + if (scan_cc.ErrorCode() == CcErrorCode::LOG_NOT_TRUNCATABLE) + { + data_sync_task->status_->SetEntriesSkippedAndNoTruncateLog(); + } - std::unordered_map - send_cache_closures; - for (size_t idx = 0; idx < scan_cc.accumulated_scan_cnt_; idx++) + scan_data_drained = true; + + // Send cache to target node group if needed. + if (data_sync_task->forward_cache_) { - FlushRecord &ref = scan_cc.DataSyncVec()[idx]; - uint16_t bucket_id = - Sharder::MapKeyHashToBucketId(ref.Key().Hash()); - NodeGroupId dest_ng = - bucket_infos->at(bucket_id)->DirtyBucketOwner(); - assert(dest_ng != UINT32_MAX); + std::shared_lock meta_lk(meta_data_mux_); + const auto bucket_infos = GetAllBucketInfosNoLocking(ng_id); + if (bucket_infos == nullptr) + { + // no longer node group owner, abort the task + LOG(ERROR) << "DataSync: Failed to get bucket infos for " + "ng#" + << ng_id; + PostProcessHashPartitionDataSyncTask( + std::move(data_sync_task), + data_sync_txm, + DataSyncTask::CkptErrorCode::SCAN_ERROR); + return; + } - // Put the record into the request for this node group. - auto ins_res = send_cache_closures.try_emplace(dest_ng); - remote::UploadBatchRequest *req_ptr = nullptr; - if (ins_res.second) + std::unordered_map + send_cache_closures; + for (size_t idx = 0; idx < scan_cc.accumulated_scan_cnt_; idx++) { - uint32_t node_id = - Sharder::Instance().LeaderNodeId(dest_ng); - std::shared_ptr channel = - Sharder::Instance().GetCcNodeServiceChannel(node_id); - if (channel == nullptr) + FlushRecord &ref = scan_cc.DataSyncVec()[idx]; + uint16_t bucket_id = + Sharder::MapKeyHashToBucketId(ref.Key().Hash()); + NodeGroupId dest_ng = + bucket_infos->at(bucket_id)->DirtyBucketOwner(); + assert(dest_ng != UINT32_MAX); + + // Put the record into the request for this node group. + auto ins_res = send_cache_closures.try_emplace(dest_ng); + remote::UploadBatchRequest *req_ptr = nullptr; + if (ins_res.second) { - // Fail to establish the channel to the tx node. - // Just skip the cache sending since it is a best - // effort try to performance improvement. - LOG(ERROR) << "UploadBatch: Failed to init the " - "channel of ng#" - << dest_ng; - send_cache_closures.erase(ins_res.first); + uint32_t node_id = + Sharder::Instance().LeaderNodeId(dest_ng); + std::shared_ptr channel = + Sharder::Instance().GetCcNodeServiceChannel( + node_id); + if (channel == nullptr) + { + // Fail to establish the channel to the tx node. + // Just skip the cache sending since it is a best + // effort try to performance improvement. + LOG(ERROR) << "UploadBatch: Failed to init the " + "channel of ng#" + << dest_ng; + send_cache_closures.erase(ins_res.first); + } + else + { + // Create a closure for the first time. + UploadBatchClosure *upload_batch_closure = + new UploadBatchClosure( + [this, data_sync_task, data_sync_txm]( + CcErrorCode res_code, + int32_t dest_ng_term) + { + // We don't care if + // the cache send + // was succeed or + // not since it's a + // best effort + // move. Just pass in no error + // so that it won't cause data sync + // failure. + PostProcessHashPartitionDataSyncTask( + std::move(data_sync_task), + data_sync_txm, + DataSyncTask::CkptErrorCode:: + NO_ERROR); + }, + 10000, + false); + + upload_batch_closure->SetChannel(node_id, channel); + + ins_res.first->second = upload_batch_closure; + req_ptr = + upload_batch_closure->UploadBatchRequest(); + req_ptr->set_node_group_id(dest_ng); + req_ptr->set_node_group_term(-1); + req_ptr->set_table_name_str(table_name.String()); + req_ptr->set_table_type( + remote::ToRemoteType::ConvertTableType( + table_name.Type())); + req_ptr->set_table_engine( + remote::ToRemoteType::ConvertTableEngine( + table_name.Engine())); + req_ptr->set_kind( + remote::UploadBatchKind::DIRTY_BUCKET_DATA); + req_ptr->set_batch_size(0); + // keys + req_ptr->clear_keys(); + // records + req_ptr->clear_records(); + // commit_ts + req_ptr->clear_commit_ts(); + // rec_status + req_ptr->clear_rec_status(); + } } else { - // Create a closure for the first time. - UploadBatchClosure *upload_batch_closure = - new UploadBatchClosure( - [this, data_sync_task, data_sync_txm]( - CcErrorCode res_code, int32_t dest_ng_term) - { - // We don't care if - // the cache send - // was succeed or - // not since it's a - // best effort - // move. Just pass in no error - // so that it won't cause data sync - // failure. - PostProcessHashPartitionDataSyncTask( - std::move(data_sync_task), - data_sync_txm, - DataSyncTask::CkptErrorCode::NO_ERROR); - }, - 10000, - false); - - upload_batch_closure->SetChannel(node_id, channel); - - ins_res.first->second = upload_batch_closure; - req_ptr = upload_batch_closure->UploadBatchRequest(); - req_ptr->set_node_group_id(dest_ng); - req_ptr->set_node_group_term(-1); - req_ptr->set_partition_id(-1); - req_ptr->set_table_name_str(table_name.String()); - req_ptr->set_table_type( - remote::ToRemoteType::ConvertTableType( - table_name.Type())); - req_ptr->set_table_engine( - remote::ToRemoteType::ConvertTableEngine( - table_name.Engine())); - req_ptr->set_kind( - remote::UploadBatchKind::DIRTY_BUCKET_DATA); - req_ptr->set_batch_size(0); - // keys - req_ptr->clear_keys(); - // records - req_ptr->clear_records(); - // commit_ts - req_ptr->clear_commit_ts(); - // rec_status - req_ptr->clear_rec_status(); + req_ptr = ins_res.first->second->UploadBatchRequest(); + } + + if (req_ptr) + { + std::string *keys_str = req_ptr->mutable_keys(); + std::string *rec_status_str = + req_ptr->mutable_rec_status(); + std::string *commit_ts_str = + req_ptr->mutable_commit_ts(); + size_t len_sizeof = sizeof(uint64_t); + const char *val_ptr = nullptr; + ref.Key().Serialize(*keys_str); + if (ref.payload_status_ == RecordStatus::Normal) + { + std::string *recs_str = req_ptr->mutable_records(); + ref.Payload()->Serialize(*recs_str); + } + const char *status_ptr = reinterpret_cast( + &(ref.payload_status_)); + rec_status_str->append(status_ptr, + sizeof(RecordStatus)); + val_ptr = + reinterpret_cast(&(ref.commit_ts_)); + commit_ts_str->append(val_ptr, len_sizeof); + req_ptr->set_batch_size(req_ptr->batch_size() + 1); } } - else + meta_lk.unlock(); + { - req_ptr = ins_res.first->second->UploadBatchRequest(); + std::unique_lock flight_lk( + data_sync_task->flight_task_mux_); + data_sync_task->flight_task_cnt_ += + send_cache_closures.size(); } - - if (req_ptr) + // Send cache to target node groups. + for (auto &[ng, upload_batch_closure] : send_cache_closures) { - std::string *keys_str = req_ptr->mutable_keys(); - std::string *rec_status_str = req_ptr->mutable_rec_status(); - std::string *commit_ts_str = req_ptr->mutable_commit_ts(); - size_t len_sizeof = sizeof(uint64_t); - const char *val_ptr = nullptr; - ref.Key().Serialize(*keys_str); - if (ref.payload_status_ == RecordStatus::Normal) - { - std::string *recs_str = req_ptr->mutable_records(); - ref.Payload()->Serialize(*recs_str); - } - const char *status_ptr = - reinterpret_cast(&(ref.payload_status_)); - rec_status_str->append(status_ptr, sizeof(RecordStatus)); - val_ptr = reinterpret_cast(&(ref.commit_ts_)); - commit_ts_str->append(val_ptr, len_sizeof); - req_ptr->set_batch_size(req_ptr->batch_size() + 1); + remote::CcRpcService_Stub stub( + upload_batch_closure->Channel()); + brpc::Controller *cntl_ptr = + upload_batch_closure->Controller(); + cntl_ptr->set_timeout_ms(10000); + // Asynchronous mode + stub.UploadBatch( + upload_batch_closure->Controller(), + upload_batch_closure->UploadBatchRequest(), + upload_batch_closure->UploadBatchResponse(), + upload_batch_closure); } } - meta_lk.unlock(); + uint64_t flush_data_size = scan_cc.accumulated_flush_data_size_; + + // nothing to flush + if (scan_cc.accumulated_scan_cnt_ == 0) { - std::unique_lock flight_lk( - data_sync_task->flight_task_mux_); - data_sync_task->flight_task_cnt_ += send_cache_closures.size(); - } - // Send cache to target node groups. - for (auto &[ng, upload_batch_closure] : send_cache_closures) - { - remote::CcRpcService_Stub stub(upload_batch_closure->Channel()); - brpc::Controller *cntl_ptr = upload_batch_closure->Controller(); - cntl_ptr->set_timeout_ms(10000); - // Asynchronous mode - stub.UploadBatch(upload_batch_closure->Controller(), - upload_batch_closure->UploadBatchRequest(), - upload_batch_closure->UploadBatchResponse(), - upload_batch_closure); + scan_cc.Reset(); + continue; } - } - - uint64_t flush_data_size = scan_cc.accumulated_flush_data_size_; - - // nothing to flush - if (scan_cc.accumulated_scan_cnt_ == 0) - { - scan_cc.Reset(); - continue; - } - // The cost of FlushRecord also needs to be considered. + // The cost of FlushRecord also needs to be considered. #ifdef WITH_JEMALLOC - flush_data_size += (scan_cc.DataSyncVec().size() * sizeof(FlushRecord) + - scan_cc.ArchiveVec().size() * sizeof(FlushRecord) + - scan_cc.MoveBaseIdxVec().size() * - sizeof(std::pair)); + flush_data_size += + (scan_cc.DataSyncVec().size() * sizeof(FlushRecord) + + scan_cc.ArchiveVec().size() * sizeof(FlushRecord) + + scan_cc.MoveBaseIdxVec().size() * + sizeof(std::pair)); #else - // Check if vectors are empty before calling malloc_usable_size - // to avoid SEGV on nullptr or invalid pointers. - // Use malloc_usable_size when ASan is enabled (vectors may be - // allocated by ASan's allocator), otherwise use - // mi_malloc_usable_size for mimalloc-allocated memory. - { - auto &data_sync_vec_ref = scan_cc.DataSyncVec(); - auto &archive_vec_ref = scan_cc.ArchiveVec(); - auto &move_base_idx_vec_ref = scan_cc.MoveBaseIdxVec(); + // Check if vectors are empty before calling malloc_usable_size + // to avoid SEGV on nullptr or invalid pointers. + // Use malloc_usable_size when ASan is enabled (vectors may be + // allocated by ASan's allocator), otherwise use + // mi_malloc_usable_size for mimalloc-allocated memory. + { + auto &data_sync_vec_ref = scan_cc.DataSyncVec(); + auto &archive_vec_ref = scan_cc.ArchiveVec(); + auto &move_base_idx_vec_ref = scan_cc.MoveBaseIdxVec(); #ifdef __SANITIZE_ADDRESS__ - // When ASan is enabled, use standard malloc_usable_size - flush_data_size += - (data_sync_vec_ref.empty() - ? 0 - : malloc_usable_size(data_sync_vec_ref.data())) + - (archive_vec_ref.empty() - ? 0 - : malloc_usable_size(archive_vec_ref.data())) + - (move_base_idx_vec_ref.empty() - ? 0 - : malloc_usable_size(move_base_idx_vec_ref.data())); + // When ASan is enabled, use standard malloc_usable_size + flush_data_size += + (data_sync_vec_ref.empty() + ? 0 + : malloc_usable_size(data_sync_vec_ref.data())) + + (archive_vec_ref.empty() + ? 0 + : malloc_usable_size(archive_vec_ref.data())) + + (move_base_idx_vec_ref.empty() + ? 0 + : malloc_usable_size(move_base_idx_vec_ref.data())); #else - // When ASan is not enabled, use mimalloc's API - flush_data_size += - (data_sync_vec_ref.empty() - ? 0 - : mi_malloc_usable_size(data_sync_vec_ref.data())) + - (archive_vec_ref.empty() - ? 0 - : mi_malloc_usable_size(archive_vec_ref.data())) + - (move_base_idx_vec_ref.empty() - ? 0 - : mi_malloc_usable_size(move_base_idx_vec_ref.data())); + // When ASan is not enabled, use mimalloc's API + flush_data_size += + (data_sync_vec_ref.empty() + ? 0 + : mi_malloc_usable_size(data_sync_vec_ref.data())) + + (archive_vec_ref.empty() + ? 0 + : mi_malloc_usable_size(archive_vec_ref.data())) + + (move_base_idx_vec_ref.empty() + ? 0 + : mi_malloc_usable_size(move_base_idx_vec_ref.data())); #endif - } + } #endif - // this thread will wait in AllocatePendingFlushDataMemQuota if - // quota is not available - uint64_t old_usage = - data_sync_mem_controller_.AllocateFlushDataMemQuota( - flush_data_size); - DLOG(INFO) << "AllocateFlushDataMemQuota old_usage: " << old_usage - << " new_usage: " << old_usage + flush_data_size - << " quota: " << data_sync_mem_controller_.FlushMemoryQuota() - << " flight_tasks: " << data_sync_task->flight_task_cnt_ - << " record count: " << scan_cc.accumulated_scan_cnt_; - - std::unique_lock heap_lk(hash_partition_ckpt_heap_mux_); - mi_threadid_t prev_thd = - mi_override_thread(hash_partition_main_thread_id_); - mi_heap_t *prev_heap = mi_heap_set_default(hash_partition_ckpt_heap_); + uint64_t old_usage = + data_sync_mem_controller_.AllocateFlushDataMemQuota( + flush_data_size); + DLOG(INFO) << "AllocateFlushDataMemQuota old_usage: " << old_usage + << " new_usage: " << old_usage + flush_data_size + << " quota: " + << data_sync_mem_controller_.FlushMemoryQuota() + << " flight_tasks: " << data_sync_task->flight_task_cnt_ + << " record count: " << scan_cc.accumulated_scan_cnt_; + + std::unique_lock heap_lk(hash_partition_ckpt_heap_mux_); + mi_threadid_t prev_thd = + mi_override_thread(hash_partition_main_thread_id_); + mi_heap_t *prev_heap = + mi_heap_set_default(hash_partition_ckpt_heap_); #if defined(WITH_JEMALLOC) - uint32_t prev_arena; - JemallocArenaSwitcher::ReadCurrentArena(prev_arena); - JemallocArenaSwitcher::SwitchToArena(hash_partition_ckpt_arena_id_); + uint32_t prev_arena; + JemallocArenaSwitcher::ReadCurrentArena(prev_arena); + JemallocArenaSwitcher::SwitchToArena(hash_partition_ckpt_arena_id_); #endif - data_sync_vec->reserve(scan_cc.accumulated_scan_cnt_); - for (size_t j = 0; j < scan_cc.accumulated_scan_cnt_; ++j) - { - auto &rec = scan_cc.DataSyncVec()[j]; - // Note. Clone key instead of move key. The memory of - // rec.Key() will be reused to avoid memory allocation. - if (rec.cce_) + data_sync_vec->reserve(scan_cc.accumulated_scan_cnt_); + for (size_t j = 0; j < scan_cc.accumulated_scan_cnt_; ++j) { - // cce_ is null means the key is already persisted on - // kv, so we don't need to put it into the flush vec. - int32_t part_id = - Sharder::MapKeyHashToHashPartitionId(rec.Key().Hash()); - if (table_name.Engine() == TableEngine::EloqKv) - { - data_sync_vec->emplace_back(rec.Key().Clone(), - rec.GetNonVersionedPayload(), - rec.payload_status_, - rec.commit_ts_, - rec.cce_, - rec.post_flush_size_, - part_id); - } - else + auto &rec = scan_cc.DataSyncVec()[j]; + // Note. Clone key instead of move key. The memory of + // rec.Key() will be reused to avoid memory allocation. + if (rec.cce_) { - data_sync_vec->emplace_back(rec.Key().Clone(), - rec.ReleaseVersionedPayload(), - rec.payload_status_, - rec.commit_ts_, - rec.cce_, - rec.post_flush_size_, - part_id); + // cce_ is null means the key is already persisted on + // kv, so we don't need to put it into the flush vec. + int32_t part_id = + Sharder::MapKeyHashToHashPartitionId(rec.Key().Hash()); + if (table_name.Engine() == TableEngine::EloqKv) + { + data_sync_vec->emplace_back( + rec.Key().Clone(), + rec.GetNonVersionedPayload(), + rec.payload_status_, + rec.commit_ts_, + rec.cce_, + rec.post_flush_size_, + part_id); + } + else + { + data_sync_vec->emplace_back( + rec.Key().Clone(), + rec.ReleaseVersionedPayload(), + rec.payload_status_, + rec.commit_ts_, + rec.cce_, + rec.post_flush_size_, + part_id); + } } } - } - for (size_t j = 0; j < scan_cc.ArchiveVec().size(); ++j) - { - auto &rec = scan_cc.ArchiveVec()[j]; - // Note. We need to ensure the copy constructor of - // FlushRecord could not be called. - rec.SetKey((*data_sync_vec)[rec.GetKeyIndex()].Key()); - } + for (size_t j = 0; j < scan_cc.ArchiveVec().size(); ++j) + { + auto &rec = scan_cc.ArchiveVec()[j]; + // Note. We need to ensure the copy constructor of + // FlushRecord could not be called. + rec.SetKey((*data_sync_vec)[rec.GetKeyIndex()].Key()); + } - for (size_t j = 0; j < scan_cc.MoveBaseIdxVec().size(); ++j) - { - size_t key_idx = scan_cc.MoveBaseIdxVec()[j]; - TxKey key_raw = (*data_sync_vec)[key_idx].Key(); - int32_t part_id = - Sharder::MapKeyHashToHashPartitionId(key_raw.Hash()); - mv_base_vec->emplace_back(std::move(key_raw), part_id); - } - mi_override_thread(prev_thd); - mi_heap_set_default(prev_heap); + for (size_t j = 0; j < scan_cc.MoveBaseIdxVec().size(); ++j) + { + size_t key_idx = scan_cc.MoveBaseIdxVec()[j]; + TxKey key_raw = (*data_sync_vec)[key_idx].Key(); + int32_t part_id = + Sharder::MapKeyHashToHashPartitionId(key_raw.Hash()); + mv_base_vec->emplace_back(std::move(key_raw), part_id); + } + mi_override_thread(prev_thd); + mi_heap_set_default(prev_heap); #if defined(WITH_JEMALLOC) - // override arena id - JemallocArenaSwitcher::SwitchToArena(prev_arena); + // override arena id + JemallocArenaSwitcher::SwitchToArena(prev_arena); #endif - heap_lk.unlock(); + heap_lk.unlock(); - std::move(scan_cc.ArchiveVec().begin(), - scan_cc.ArchiveVec().end(), - std::back_inserter(*archive_vec)); + std::move(scan_cc.ArchiveVec().begin(), + scan_cc.ArchiveVec().end(), + std::back_inserter(*archive_vec)); - vec_mem_usage += flush_data_size; + vec_mem_usage += flush_data_size; - scan_data_drained = scan_cc.IsDrained() && scan_data_drained; + scan_data_drained = scan_cc.IsDrained() && scan_data_drained; - { - std::unique_lock flight_task_lk( - data_sync_task->flight_task_mux_); - if (data_sync_task->ckpt_err_ == - DataSyncTask::CkptErrorCode::FLUSH_ERROR) { - flight_task_lk.unlock(); - - LOG(WARNING) - << "There are error during flush for this data sync: " - << data_sync_txm->TxNumber() << " on worker#" << worker_idx - << ". Terminal this datasync task."; - // 1. Release read intent on paused key - if (!scan_cc.IsDrained()) + std::unique_lock flight_task_lk( + data_sync_task->flight_task_mux_); + if (data_sync_task->ckpt_err_ == + DataSyncTask::CkptErrorCode::FLUSH_ERROR) { - scan_cc.Reset( - HashPartitionDataSyncScanCc::OpType::Terminated); - EnqueueLowPriorityCcRequestToShard(worker_idx, &scan_cc); - scan_cc.Wait(); + flight_task_lk.unlock(); + + LOG(WARNING) + << "There are error during flush for this data sync: " + << data_sync_txm->TxNumber() << " on worker#" + << worker_idx << ". Terminal this datasync task."; + // 1. Release read intent on paused key + if (!scan_cc.IsDrained()) + { + scan_cc.Reset( + HashPartitionDataSyncScanCc::OpType::Terminated); + EnqueueLowPriorityCcRequestToShard(worker_idx, + &scan_cc); + scan_cc.Wait(); + } + // 2. Release memory usage on this datasync worker. + data_sync_vec->clear(); + archive_vec->clear(); + mv_base_vec->clear(); + data_sync_mem_controller_.DeallocateFlushMemQuota( + vec_mem_usage); + break; } - // 2. Release memory usage on this datasync worker. - data_sync_vec->clear(); - archive_vec->clear(); - mv_base_vec->clear(); - data_sync_mem_controller_.DeallocateFlushMemQuota( - vec_mem_usage); - break; - } - - // Flush worker will call PostProcessDataSyncTask() to - // decrement flight task count. - data_sync_task->flight_task_cnt_ += 1; - } - AddFlushTaskEntry( - std::make_unique(std::move(data_sync_vec), - std::move(archive_vec), - std::move(mv_base_vec), - data_sync_txm, - data_sync_task, - catalog_rec.CopySchema(), - vec_mem_usage)); - - data_sync_vec = std::make_unique>(); - - archive_vec = std::make_unique>(); - - mv_base_vec = - std::make_unique>>(); - - vec_mem_usage = 0; - - if (scan_cc.scan_heap_is_full_ == 1) - { - // Clear the FlushRecords' memory of scan cc since the - // DataSyncScan heap is full. - auto &data_sync_vec_ref = scan_cc.DataSyncVec(); - auto &archive_vec_ref = scan_cc.ArchiveVec(); - ReleaseDataSyncScanHeapCc release_scan_heap_cc(&data_sync_vec_ref, - &archive_vec_ref); - EnqueueLowPriorityCcRequestToShard(worker_idx, - &release_scan_heap_cc); - release_scan_heap_cc.Wait(); - } - - scan_cc.Reset(); - } - - // release scan heap memory after scan finish - auto &data_sync_vec_ref = scan_cc.DataSyncVec(); - auto &archive_vec_ref = scan_cc.ArchiveVec(); - ReleaseDataSyncScanHeapCc release_scan_heap_cc(&data_sync_vec_ref, - &archive_vec_ref); - EnqueueLowPriorityCcRequestToShard(worker_idx, &release_scan_heap_cc); - release_scan_heap_cc.Wait(); - - if (!data_sync_vec->empty() || !archive_vec->empty() || - !mv_base_vec->empty()) - { - std::unique_lock flight_task_lk( - data_sync_task->flight_task_mux_); - if (data_sync_task->ckpt_err_ == DataSyncTask::CkptErrorCode::NO_ERROR) - { - data_sync_task->flight_task_cnt_ += 1; - flight_task_lk.unlock(); + // Flush worker will call PostProcessDataSyncTask() to + // decrement flight task count. + data_sync_task->flight_task_cnt_ += 1; + } AddFlushTaskEntry( std::make_unique(std::move(data_sync_vec), @@ -5179,13 +5272,67 @@ void LocalCcShards::DataSyncForHashPartition( data_sync_task, catalog_rec.CopySchema(), vec_mem_usage)); + + data_sync_vec = std::make_unique>(); + + archive_vec = std::make_unique>(); + + mv_base_vec = + std::make_unique>>(); + + vec_mem_usage = 0; + + if (scan_cc.scan_heap_is_full_ == 1) + { + // Clear the FlushRecords' memory of scan cc since the + // DataSyncScan heap is full. + auto &data_sync_vec_ref = scan_cc.DataSyncVec(); + auto &archive_vec_ref = scan_cc.ArchiveVec(); + ReleaseDataSyncScanHeapCc release_scan_heap_cc( + &data_sync_vec_ref, &archive_vec_ref); + EnqueueLowPriorityCcRequestToShard(worker_idx, + &release_scan_heap_cc); + release_scan_heap_cc.Wait(); + } + + scan_cc.Reset(); } - else + + // release scan heap memory after scan finish + auto &data_sync_vec_ref = scan_cc.DataSyncVec(); + auto &archive_vec_ref = scan_cc.ArchiveVec(); + ReleaseDataSyncScanHeapCc release_scan_heap_cc(&data_sync_vec_ref, + &archive_vec_ref); + EnqueueLowPriorityCcRequestToShard(worker_idx, &release_scan_heap_cc); + release_scan_heap_cc.Wait(); + + if (!data_sync_vec->empty() || !archive_vec->empty() || + !mv_base_vec->empty()) { - // There are error during flush, and if we do not put the - // current batch data into flush worker, should release the - // memory usage. - data_sync_mem_controller_.DeallocateFlushMemQuota(vec_mem_usage); + std::unique_lock flight_task_lk( + data_sync_task->flight_task_mux_); + if (data_sync_task->ckpt_err_ == + DataSyncTask::CkptErrorCode::NO_ERROR) + { + data_sync_task->flight_task_cnt_ += 1; + flight_task_lk.unlock(); + AddFlushTaskEntry( + std::make_unique(std::move(data_sync_vec), + std::move(archive_vec), + std::move(mv_base_vec), + data_sync_txm, + data_sync_task, + catalog_rec.CopySchema(), + vec_mem_usage)); + } + else + { + // There are error during flush, and if we do not put the + // current batch data into flush worker, should release the + // memory usage. + data_sync_mem_controller_.DeallocateFlushMemQuota( + vec_mem_usage); + } } } @@ -5663,93 +5810,120 @@ void LocalCcShards::SplitFlushRange( txservice::CommitTx(split_txm); } +size_t LocalCcShards::DataSyncWorkerToFlushDataWorker( + size_t data_sync_worker_id) const +{ + return data_sync_worker_id % flush_data_worker_ctx_.worker_num_; +} + void LocalCcShards::AddFlushTaskEntry(std::unique_ptr &&entry) { - cur_flush_buffer_.AddFlushTaskEntry(std::move(entry)); + assert(cur_flush_buffers_.size() == + static_cast(data_sync_worker_ctx_.worker_num_)); + assert(pending_flush_work_.size() == + static_cast(flush_data_worker_ctx_.worker_num_)); + assert(entry->data_sync_task_ != nullptr); + + const auto &data_sync_task = entry->data_sync_task_; + const size_t data_sync_worker_id = static_cast( + data_sync_task->id_ % data_sync_worker_ctx_.worker_num_); + const size_t flush_target = + DataSyncWorkerToFlushDataWorker(data_sync_worker_id); + + std::unique_lock worker_lk(flush_data_worker_ctx_.mux_); + auto &cur_flush_buffer = *cur_flush_buffers_[data_sync_worker_id]; + + cur_flush_buffer.AddFlushTaskEntry(std::move(entry)); std::unique_ptr flush_data_task = - cur_flush_buffer_.MoveFlushData(false); + cur_flush_buffer.MoveFlushData(false); if (flush_data_task != nullptr) { - std::unique_lock worker_lk(flush_data_worker_ctx_.mux_); + auto &pending_flush_work = pending_flush_work_[flush_target]; // Try to merge with the last task if queue is not empty - if (!pending_flush_work_.empty()) + if (!pending_flush_work.empty()) { - auto &last_task = pending_flush_work_.back(); + auto &last_task = pending_flush_work.back(); if (last_task->MergeFrom(std::move(flush_data_task))) { // Merge successful, task was merged into last_task - flush_data_worker_ctx_.cv_.notify_one(); + flush_data_worker_ctx_.cv_.notify_all(); return; } } // Could not merge, wait if queue is full - while (pending_flush_work_.size() >= - static_cast(flush_data_worker_ctx_.worker_num_)) + /* + while (pending_flush_work.size() >= 2) { flush_data_worker_ctx_.cv_.wait(worker_lk); } + */ // Add as new task - pending_flush_work_.emplace_back(std::move(flush_data_task)); - flush_data_worker_ctx_.cv_.notify_one(); + pending_flush_work.emplace_back(std::move(flush_data_task)); + flush_data_worker_ctx_.cv_.notify_all(); } } void LocalCcShards::FlushCurrentFlushBuffer() { - std::unique_ptr flush_data_task = - cur_flush_buffer_.MoveFlushData(true); - if (flush_data_task != nullptr) + assert(cur_flush_buffers_.size() == + static_cast(data_sync_worker_ctx_.worker_num_)); + assert(pending_flush_work_.size() == + static_cast(flush_data_worker_ctx_.worker_num_)); + + std::unique_lock worker_lk(flush_data_worker_ctx_.mux_); + for (int i = 0; i < data_sync_worker_ctx_.worker_num_; ++i) { - std::unique_lock worker_lk(flush_data_worker_ctx_.mux_); + auto &cur_flush_buffer = *cur_flush_buffers_[i]; + size_t flush_target = + DataSyncWorkerToFlushDataWorker(static_cast(i)); + auto &pending_flush_work = pending_flush_work_[flush_target]; - // Try to merge with the last task if queue is not empty - if (!pending_flush_work_.empty()) + std::unique_ptr flush_data_task = + cur_flush_buffer.MoveFlushData(true); + if (flush_data_task != nullptr) { - auto &last_task = pending_flush_work_.back(); - if (last_task->MergeFrom(std::move(flush_data_task))) + // Try to merge with the last task if queue is not empty + if (!pending_flush_work.empty()) { - // Merge successful, task was merged into last_task - flush_data_worker_ctx_.cv_.notify_one(); - return; + auto &last_task = pending_flush_work.back(); + if (last_task->MergeFrom(std::move(flush_data_task))) + { + // Merge successful, task was merged into last_task + flush_data_worker_ctx_.cv_.notify_all(); + continue; + } } - } - // Could not merge, wait if queue is full - while (pending_flush_work_.size() >= - static_cast(flush_data_worker_ctx_.worker_num_)) - { - flush_data_worker_ctx_.cv_.wait(worker_lk); + // Add as new task + pending_flush_work.emplace_back(std::move(flush_data_task)); + flush_data_worker_ctx_.cv_.notify_all(); } - - // Add as new task - pending_flush_work_.emplace_back(std::move(flush_data_task)); - flush_data_worker_ctx_.cv_.notify_one(); } } -void LocalCcShards::FlushData(std::unique_lock &flush_worker_lk) +bool LocalCcShards::ShouldYieldFlushData(size_t worker_idx) { - // Retrieve first pending work and pop it (FIFO). - std::unique_ptr cur_work = - std::move(pending_flush_work_.front()); - pending_flush_work_.pop_front(); - - // Notify any threads waiting for queue space - flush_data_worker_ctx_.cv_.notify_all(); - - flush_worker_lk.unlock(); + std::lock_guard lk(flush_data_worker_ctx_.mux_); + return !pending_flush_work_[worker_idx].empty(); +} +void LocalCcShards::FlushDataImpl(FlushDataTask *cur_work, + size_t worker_idx, + const std::function &sync_yield_func, + const std::function &yield_fn, + const std::function &resume_fn) +{ auto &flush_task_entries = cur_work->flush_task_entries_; - bool succ = true; - // Flushes to the data store + if (EnableMvcc()) { - succ = store_hd_->CopyBaseToArchive(flush_task_entries); + succ = store_hd_->CopyBaseToArchive( + flush_task_entries, &yield_fn, &resume_fn); if (!succ) { LOG(ERROR) << "DataSync CopyBaseToArchive flush to kv " @@ -5759,7 +5933,8 @@ void LocalCcShards::FlushData(std::unique_lock &flush_worker_lk) if (succ) { - succ = store_hd_->PutAll(flush_task_entries); + succ = store_hd_->PutAll( + flush_task_entries, &yield_fn, &resume_fn, &sync_yield_func); if (!succ) { LOG(ERROR) << "DataSync PutAll flush to kv " @@ -5769,7 +5944,8 @@ void LocalCcShards::FlushData(std::unique_lock &flush_worker_lk) if (succ && EnableMvcc()) { - succ = store_hd_->PutArchivesAll(flush_task_entries); + succ = store_hd_->PutArchivesAll( + flush_task_entries, &yield_fn, &resume_fn); if (!succ) { LOG(ERROR) << "DataSync PutArchivesAll flush to " @@ -5777,7 +5953,6 @@ void LocalCcShards::FlushData(std::unique_lock &flush_worker_lk) } } - // Persist data in kv store if needed if (succ && store_hd_->NeedPersistKV()) { std::vector kv_table_names; @@ -5785,7 +5960,7 @@ void LocalCcShards::FlushData(std::unique_lock &flush_worker_lk) { kv_table_names.push_back(table_name.data()); } - succ = store_hd_->PersistKV(kv_table_names); + succ = store_hd_->PersistKV(kv_table_names, &yield_fn, &resume_fn); } // Record that data was written in DataSyncStatus if flush succeeded. @@ -5817,6 +5992,9 @@ void LocalCcShards::FlushData(std::unique_lock &flush_worker_lk) // Update cce ckpt ts in memory if (succ) { + size_t iterations_since_yield = 0; + constexpr size_t MAX_ITERATIONS_WITHOUT_YIELD = 10; + for (auto &[kv_table_name, entries] : flush_task_entries) { for (auto &entry : entries) @@ -5864,12 +6042,26 @@ void LocalCcShards::FlushData(std::unique_lock &flush_worker_lk) entry->data_sync_task_->node_group_term_, table_name, cce_entries_map); + update_cce_req.SetCoroCallbacks(&yield_fn, &resume_fn); for (auto &[core_idx, cce_entries] : cce_entries_map) { updated_ckpt_ts_core_ids.insert(core_idx); EnqueueToCcShard(core_idx, &update_cce_req); } - update_cce_req.Wait(); + + bool is_finished = update_cce_req.IsFinished(); + if (is_finished) + { + iterations_since_yield++; + if (iterations_since_yield >= + MAX_ITERATIONS_WITHOUT_YIELD) + { + sync_yield_func(); + iterations_since_yield = 0; + } + } + + update_cce_req.Wait(&yield_fn, &resume_fn); } } } @@ -5884,11 +6076,12 @@ void LocalCcShards::FlushData(std::unique_lock &flush_worker_lk) return true; }, updated_ckpt_ts_core_ids.size()); + reset_cc.SetCoroCallbacks(&yield_fn, &resume_fn); for (uint16_t core_idx : updated_ckpt_ts_core_ids) { EnqueueToCcShard(core_idx, &reset_cc); } - reset_cc.Wait(); + reset_cc.Wait(&yield_fn, &resume_fn); auto ckpt_err = succ ? DataSyncTask::CkptErrorCode::NO_ERROR : DataSyncTask::CkptErrorCode::FLUSH_ERROR; @@ -5901,23 +6094,136 @@ void LocalCcShards::FlushData(std::unique_lock &flush_worker_lk) << " new_usage: " << old_usage - cur_work->pending_flush_size_ << " quota: " << data_sync_mem_controller_.FlushMemoryQuota(); - PostProcessFlushTaskEntries(flush_task_entries, ckpt_err); + PostProcessFlushTaskEntries( + flush_task_entries, ckpt_err, &yield_fn, &resume_fn, &sync_yield_func); +} + +void LocalCcShards::FlushData(std::unique_lock &flush_worker_lk, + size_t worker_idx) +{ + assert(worker_idx < pending_flush_work_.size()); + auto &pending_flush_work = pending_flush_work_[worker_idx]; + + // Retrieve first pending work and pop it (FIFO). + std::unique_ptr cur_work = + std::move(pending_flush_work.front()); + pending_flush_work.pop_front(); + + // Notify any threads waiting for queue space + flush_data_worker_ctx_.cv_.notify_all(); + + flush_worker_lk.unlock(); + + FlushDataImpl(cur_work.get(), worker_idx, []() {}, []() {}, []() {}); + flush_worker_lk.lock(); } -void LocalCcShards::FlushDataWorker() +void LocalCcShards::FlushDataWorker(size_t worker_idx) { - size_t previous_flush_size = 0; - auto previous_size_update_time = std::chrono::steady_clock::now(); + assert(worker_idx < + static_cast(flush_data_worker_ctx_.worker_num_)); + assert(worker_idx < static_cast(pending_flush_work_.size())); + assert(worker_idx < static_cast(resume_queue_.size())); + + auto &pending_flush_work = pending_flush_work_[worker_idx]; + auto &resume_queue = resume_queue_[worker_idx]; + + using clock = std::chrono::steady_clock; + using continuation = boost::context::continuation; + std::vector previous_flush_sizes(data_sync_worker_ctx_.worker_num_, + 0); + auto previous_size_update_time = clock::now(); + std::unique_lock flush_worker_lk(flush_data_worker_ctx_.mux_); while (flush_data_worker_ctx_.status_ == WorkerStatus::Active) { + if (!pending_flush_work.empty()) + { + std::unique_ptr cur_work = + std::move(pending_flush_work.front()); + pending_flush_work.pop_front(); + flush_data_worker_ctx_.cv_.notify_all(); + flush_worker_lk.unlock(); + + auto ctx = std::make_shared(); + ctx->task_ = std::move(cur_work); + ctx->coro_ = boost::context::callcc( + std::allocator_arg, + flush_coro_stack_allocator_, + [this, ctx, worker_idx](continuation &&sink) + { + ctx->yield_fn = [&sink]() { sink = sink.resume(); }; + ctx->sync_yield_func = + [&sink, + weak_ctx = std::weak_ptr(ctx), + this, + worker_idx]() + { + if (auto c = weak_ctx.lock()) + { + { + std::lock_guard lk( + flush_data_worker_ctx_.mux_); + resume_queue_[worker_idx].push_back( + std::move(c)); + flush_data_worker_ctx_.cv_.notify_all(); + } + sink = sink.resume(); + } + }; + ctx->resume_fn = [this, + worker_idx, + weak_ctx = std::weak_ptr(ctx)]() + { + if (auto c = weak_ctx.lock()) + { + std::lock_guard lk( + flush_data_worker_ctx_.mux_); + resume_queue_[worker_idx].push_back(std::move(c)); + flush_data_worker_ctx_.cv_.notify_all(); + } + else + { + LOG(FATAL) + << "CoroCtx resume_fn: weak_ctx is expired"; + } + }; + FlushDataImpl(ctx->task_.get(), + worker_idx, + ctx->sync_yield_func, + ctx->yield_fn, + ctx->resume_fn); + return std::move(sink); + }); + + flush_worker_lk.lock(); + continue; + } + + if (!resume_queue.empty()) + { + std::shared_ptr ctx = std::move(resume_queue.front()); + resume_queue.pop_front(); + flush_worker_lk.unlock(); + + ctx->coro_ = ctx->coro_.resume(); + + flush_worker_lk.lock(); + continue; + } + flush_data_worker_ctx_.cv_.wait_for( flush_worker_lk, 10s, - [this, &previous_flush_size, &previous_size_update_time] + [this, + worker_idx, + &pending_flush_work, + &resume_queue, + &previous_flush_sizes, + &previous_size_update_time] { - if (!pending_flush_work_.empty() || + if (!pending_flush_work.empty() || !resume_queue.empty() || flush_data_worker_ctx_.status_ == WorkerStatus::Terminated) { return true; @@ -5925,42 +6231,41 @@ void LocalCcShards::FlushDataWorker() auto current_time = std::chrono::steady_clock::now(); if (current_time - previous_size_update_time > 10s) { - size_t current_flush_size = - cur_flush_buffer_.GetPendingFlushSize(); - bool flush_size_changed = - current_flush_size != previous_flush_size; - previous_flush_size = current_flush_size; previous_size_update_time = current_time; - if (!flush_size_changed && current_flush_size > 0) + for (int i = 0; i < data_sync_worker_ctx_.worker_num_; ++i) { - // data sync might be stuck due to lock conflict with - // DDL. Flush current flush buffer to release catalog - // read lock held by ongoing data sync tx, which might - // block the DDL. - std::unique_ptr flush_data_task = - cur_flush_buffer_.MoveFlushData(true); - if (flush_data_task != nullptr) + if (DataSyncWorkerToFlushDataWorker( + static_cast(i)) != worker_idx) + { + continue; + } + size_t current_flush_size = + cur_flush_buffers_[i]->GetPendingFlushSize(); + bool flush_size_changed = + current_flush_size != previous_flush_sizes[i]; + previous_flush_sizes[i] = current_flush_size; + if (!flush_size_changed && current_flush_size > 0) { - // Try to merge with the last task if queue is not - // empty Note: flush_worker_lk is already held here - // (inside condition variable predicate) - if (!pending_flush_work_.empty()) + // data sync might be stuck due to lock conflict + // with DDL. Flush buffer to release catalog read + // lock held by ongoing data sync tx. + std::unique_ptr flush_data_task = + cur_flush_buffers_[i]->MoveFlushData(true); + if (flush_data_task != nullptr) { - auto &last_task = pending_flush_work_.back(); - if (last_task->MergeFrom( - std::move(flush_data_task))) + if (!pending_flush_work.empty()) { - // Merge successful, task was merged into - // last_task - return true; + auto &last_task = pending_flush_work.back(); + if (last_task->MergeFrom( + std::move(flush_data_task))) + { + return true; + } } + pending_flush_work.emplace_back( + std::move(flush_data_task)); + return true; } - - // Add as new task. We just checked that - // pending_flush_work_ is empty, - pending_flush_work_.emplace_back( - std::move(flush_data_task)); - return true; } } } @@ -5968,17 +6273,78 @@ void LocalCcShards::FlushDataWorker() return false; }); - if (pending_flush_work_.empty()) + continue; + } + + while (!resume_queue.empty() || !pending_flush_work.empty()) + { + if (!pending_flush_work.empty()) { + std::unique_ptr cur_work = + std::move(pending_flush_work.front()); + pending_flush_work.pop_front(); + flush_data_worker_ctx_.cv_.notify_all(); + flush_worker_lk.unlock(); + + auto ctx = std::make_shared(); + ctx->task_ = std::move(cur_work); + ctx->coro_ = boost::context::callcc( + std::allocator_arg, + flush_coro_stack_allocator_, + [this, ctx, worker_idx](continuation &&sink) + { + ctx->yield_fn = [&sink]() { sink = sink.resume(); }; + ctx->sync_yield_func = + [&sink, + weak_ctx = std::weak_ptr(ctx), + this, + worker_idx]() + { + if (auto c = weak_ctx.lock()) + { + { + std::lock_guard lk( + flush_data_worker_ctx_.mux_); + resume_queue_[worker_idx].push_back( + std::move(c)); + flush_data_worker_ctx_.cv_.notify_all(); + } + sink = sink.resume(); + } + }; + ctx->resume_fn = [this, + worker_idx, + weak_ctx = std::weak_ptr(ctx)]() + { + if (auto c = weak_ctx.lock()) + { + std::lock_guard lk( + flush_data_worker_ctx_.mux_); + resume_queue_[worker_idx].push_back(std::move(c)); + flush_data_worker_ctx_.cv_.notify_all(); + } + }; + FlushDataImpl(ctx->task_.get(), + worker_idx, + ctx->sync_yield_func, + ctx->yield_fn, + ctx->resume_fn); + return std::move(sink); + }); + + flush_worker_lk.lock(); continue; } - FlushData(flush_worker_lk); - } - - while (!pending_flush_work_.empty()) - { - FlushData(flush_worker_lk); + if (!resume_queue.empty()) + { + std::shared_ptr ctx = std::move(resume_queue.front()); + resume_queue.pop_front(); + flush_worker_lk.unlock(); + ctx->coro_ = ctx->coro_.resume(); + flush_worker_lk.lock(); + continue; + } } } @@ -6012,7 +6378,10 @@ void LocalCcShards::RangeSplitWorker() } bool LocalCcShards::UpdateStoreSlices( - std::vector &flush_tasks) + std::vector &flush_tasks, + const std::function *yield_fptr, + const std::function *resume_fptr, + const std::function *sync_yield_fptr) { std::vector update_range_slice_reqs; @@ -6055,7 +6424,8 @@ bool LocalCcShards::UpdateStoreSlices( if (!update_range_slice_reqs.empty()) { - bool success = store_hd_->UpdateRangeSlices(update_range_slice_reqs); + bool success = store_hd_->UpdateRangeSlices( + update_range_slice_reqs, yield_fptr, resume_fptr, sync_yield_fptr); return success; } diff --git a/tx_service/src/checkpointer.cpp b/tx_service/src/checkpointer.cpp index 8ae1ea04..1031e6b0 100644 --- a/tx_service/src/checkpointer.cpp +++ b/tx_service/src/checkpointer.cpp @@ -568,7 +568,7 @@ bool Checkpointer::CkptEntryForTest( std::vector>> &flush_task_entries) { - return store_hd_->PutAll(flush_task_entries); + return store_hd_->PutAll(flush_task_entries, nullptr, nullptr); } bool Checkpointer::FlushArchiveForTest( @@ -576,6 +576,6 @@ bool Checkpointer::FlushArchiveForTest( std::vector>> &flush_task_entries) { - return store_hd_->PutArchivesAll(flush_task_entries); + return store_hd_->PutArchivesAll(flush_task_entries, nullptr, nullptr); } } // namespace txservice diff --git a/tx_service/src/fault/log_replay_service.cpp b/tx_service/src/fault/log_replay_service.cpp index 77c35ec5..81335512 100644 --- a/tx_service/src/fault/log_replay_service.cpp +++ b/tx_service/src/fault/log_replay_service.cpp @@ -128,10 +128,10 @@ RecoveryService::RecoveryService(LocalCcShards &local_shards, inbound_mux_); for (const auto &entry : inbound_connections_) { - if (entry.second.cc_ng_id_ == ng_id && - entry.second.cc_ng_term_ == + if (entry.second->cc_ng_id_ == ng_id && + entry.second->cc_ng_term_ == candidate_term && - entry.second.recovering_) + entry.second->recovering_) { has_replay_connection = true; break; @@ -287,6 +287,7 @@ void RecoveryService::Connect(::google::protobuf::RpcController *controller, // Invalid connection request. return; } + // node group is trying to recover cc_ng_term. Check if this log group // has finished replay. if (Sharder::Instance().CheckLogGroupReplayFinished( @@ -302,14 +303,14 @@ void RecoveryService::Connect(::google::protobuf::RpcController *controller, // Clean up old log group stream connection. for (auto &[stream_id, info] : inbound_connections_) { - if (info.cc_ng_id_ == cc_ng_id && - info.log_group_id_ == log_group_id) + if (info->cc_ng_id_ == cc_ng_id && + info->log_group_id_ == log_group_id) { - if (cc_ng_term > info.cc_ng_term_) + if (cc_ng_term > info->cc_ng_term_) { // cc_ng_term > info.cc_ng_term_, the cc_ng has failed over } - else if (cc_ng_term == info.cc_ng_term_) + else if (cc_ng_term == info->cc_ng_term_) { // the cc_ng is still recovering, and this request is a // response for RecoveryService's stream timeout and resend @@ -344,12 +345,10 @@ void RecoveryService::Connect(::google::protobuf::RpcController *controller, return; } - inbound_connections_.try_emplace(stream_socket, - log_group_id, - cc_ng_id, - cc_ng_term, - replay_start_ts, - recovering); + inbound_connections_.try_emplace( + stream_socket, + std::make_unique( + log_group_id, cc_ng_id, cc_ng_term, replay_start_ts, recovering)); // Initialize/refresh the global DDL barrier for recovering streams. if (recovering) @@ -441,7 +440,12 @@ int RecoveryService::on_received_messages(brpc::StreamId stream_id, NodeGroupReplayBarrier *barrier; { std::lock_guard lk(inbound_mux_); - info = &inbound_connections_.find(stream_id)->second; + auto it = inbound_connections_.find(stream_id); + if (it == inbound_connections_.end()) + { + return -1; + } + info = it->second.get(); barrier = GetReplayBarrier(info->cc_ng_id_); } bthread::Mutex &mux = info->mux_; @@ -836,7 +840,7 @@ void RecoveryService::on_idle_timeout(brpc::StreamId id) // this stream has been replaced by a newer one return; } - info = &it->second; + info = it->second.get(); } BAIDU_SCOPED_LOCK(info->mux_); if (info->recovering_) @@ -871,8 +875,14 @@ void RecoveryService::on_closed(brpc::StreamId id) ConnectionInfo *info; { std::lock_guard lk(inbound_mux_); - info = &inbound_connections_.find(id)->second; + auto it = inbound_connections_.find(id); + if (it == inbound_connections_.end()) + { + return; + } + info = it->second.get(); } + WaitAndClearRequests(id, info->mux_, info->on_fly_cnt_, @@ -928,18 +938,18 @@ void RecoveryService::WaitAndClearRequests(brpc::StreamId stream_id, return; } - error_node_group_id = it->second.cc_ng_id_; - error_term = it->second.cc_ng_term_; + error_node_group_id = it->second->cc_ng_id_; + error_term = it->second->cc_ng_term_; uint64_t replay_start_ts = 0; // close all the streams belonging to the current node group and term. for (const auto &[stream_id, info] : inbound_connections_) { assert(replay_start_ts == 0 || - replay_start_ts == info.replay_start_ts_); - replay_start_ts = info.replay_start_ts_; - if (info.cc_ng_id_ == error_node_group_id && - info.cc_ng_term_ == error_term) + replay_start_ts == info->replay_start_ts_); + replay_start_ts = info->replay_start_ts_; + if (info->cc_ng_id_ == error_node_group_id && + info->cc_ng_term_ == error_term) { brpc::StreamClose(stream_id); } diff --git a/tx_service/src/tx_index_operation.cpp b/tx_service/src/tx_index_operation.cpp index 7ef05670..1c852020 100644 --- a/tx_service/src/tx_index_operation.cpp +++ b/tx_service/src/tx_index_operation.cpp @@ -462,6 +462,36 @@ void UpsertTableIndexOp::Forward(TransactionExecution *txm) { if (txm->CheckLeaderTerm()) { + const CcEntryAddr *cluster_config_addr = nullptr; + const auto &meta_rset = txm->rw_set_.MetaDataReadSet(); + for (const auto &[cce_addr, read_entry_pair] : meta_rset) + { + if (read_entry_pair.second == cluster_config_ccm_name_sv) + { + cluster_config_addr = &cce_addr; + break; + } + } + + if (cluster_config_addr) + { + DLOG(INFO) + << "Alter Table Index transaction unlock cluster " + "config " + "lock failed, txn: " + << txm->TxNumber() << ", error code: " + << (int) + unlock_cluster_config_op_.hd_result_.ErrorCode() + << ", error message: " + << unlock_cluster_config_op_.hd_result_.ErrorMsg() + << ", cluster config addr term" + << cluster_config_addr->Term() + << ", cluster config addr node group id: " + << cluster_config_addr->NodeGroupId() + << ", node group term: " + << Sharder::Instance().LeaderTerm( + cluster_config_addr->NodeGroupId()); + } // Releasing a local read lock should never fail assert(false); }