diff --git a/store_handler/eloq_data_store_service/data_store_service.cpp b/store_handler/eloq_data_store_service/data_store_service.cpp index acec7f8f..58ee2582 100644 --- a/store_handler/eloq_data_store_service/data_store_service.cpp +++ b/store_handler/eloq_data_store_service/data_store_service.cpp @@ -85,6 +85,11 @@ thread_local ObjectPool thread_local ObjectPool local_sync_file_cache_req_pool_; +thread_local const DataStoreService *tls_read_submit_slot_owner = nullptr; +thread_local uint32_t tls_read_submit_slot_idx = UINT32_MAX; +thread_local std::array + tls_read_submit_depths{}; + TTLWrapperCache::TTLWrapperCache() { ttl_check_running_ = true; @@ -273,6 +278,91 @@ DataStoreService::~DataStoreService() } } +DataStoreService::ReadSubmitGuard::ReadSubmitGuard(DataStoreService *service, + uint32_t shard_id, + uint32_t slot_idx) + : service_(service), shard_id_(shard_id), slot_idx_(slot_idx) +{ +} + +DataStoreService::ReadSubmitGuard::ReadSubmitGuard( + ReadSubmitGuard &&other) noexcept + : service_(other.service_), + shard_id_(other.shard_id_), + slot_idx_(other.slot_idx_) +{ + other.service_ = nullptr; + other.shard_id_ = UINT32_MAX; + other.slot_idx_ = UINT32_MAX; +} + +DataStoreService::ReadSubmitGuard::~ReadSubmitGuard() +{ + if (service_ != nullptr) + { + service_->LeaveReadSubmitWindow(shard_id_, slot_idx_); + } +} + +uint32_t DataStoreService::GetReadSubmitSlotIndex() +{ + if (__builtin_expect(tls_read_submit_slot_owner == this && + tls_read_submit_slot_idx != UINT32_MAX, + 1)) + { + return tls_read_submit_slot_idx; + } + + uint32_t slot_idx = + next_read_submit_slot_idx_.fetch_add(1, std::memory_order_seq_cst); + CHECK_LT(slot_idx, kReadSubmitSlotCount) + << "Too many read submit threads for DataStoreService"; + + tls_read_submit_slot_owner = this; + tls_read_submit_slot_idx = slot_idx; + tls_read_submit_depths.fill(0); + return slot_idx; +} + +DataStoreService::ReadSubmitGuard DataStoreService::EnterReadSubmitWindow( + uint32_t shard_id) +{ + uint32_t slot_idx = GetReadSubmitSlotIndex(); + uint32_t depth = ++tls_read_submit_depths.at(shard_id); + data_shards_.at(shard_id).read_submit_slots_.at(slot_idx).depth.store( + depth, std::memory_order_seq_cst); + return ReadSubmitGuard(this, shard_id, slot_idx); +} + +void DataStoreService::LeaveReadSubmitWindow(uint32_t shard_id, + uint32_t slot_idx) +{ + uint32_t &local_depth = tls_read_submit_depths.at(shard_id); + CHECK_GT(local_depth, 0); + uint32_t depth = --local_depth; + data_shards_.at(shard_id).read_submit_slots_.at(slot_idx).depth.store( + depth, std::memory_order_seq_cst); +} + +void DataStoreService::WaitReadSubmitWindowsDrained( + const DataShard &ds_ref) const +{ + bool drained = false; + while (!drained) + { + drained = true; + for (const auto &slot : ds_ref.read_submit_slots_) + { + if (slot.depth.load(std::memory_order_seq_cst) != 0) + { + drained = false; + bthread_usleep(1000); + break; + } + } + } +} + bool DataStoreService::StartService(bool create_db_if_missing) { if (server_ != nullptr) @@ -477,7 +567,8 @@ void DataStoreService::Read(::google::protobuf::RpcController *controller, } DataShard &ds_ref = data_shards_.at(shard_id); - auto shard_status = ds_ref.shard_status_.load(std::memory_order_acquire); + auto read_submit_guard = EnterReadSubmitWindow(shard_id); + auto shard_status = ds_ref.shard_status_.load(std::memory_order_seq_cst); if (shard_status != DSShardStatus::ReadOnly && shard_status != DSShardStatus::ReadWrite) { @@ -514,7 +605,8 @@ void DataStoreService::Read(const std::string_view table_name, } DataShard &ds_ref = data_shards_.at(shard_id); - auto shard_status = ds_ref.shard_status_.load(std::memory_order_acquire); + auto read_submit_guard = EnterReadSubmitWindow(shard_id); + auto shard_status = ds_ref.shard_status_.load(std::memory_order_seq_cst); if (shard_status != DSShardStatus::ReadOnly && shard_status != DSShardStatus::ReadWrite) { @@ -559,7 +651,7 @@ void DataStoreService::FlushData( IncreaseWriteReqCount(shard_id); DataShard &ds_ref = data_shards_.at(shard_id); - auto shard_status = ds_ref.shard_status_.load(std::memory_order_acquire); + auto shard_status = ds_ref.shard_status_.load(std::memory_order_seq_cst); if (shard_status != DSShardStatus::ReadWrite) { DecreaseWriteReqCount(shard_id); @@ -713,7 +805,7 @@ void DataStoreService::FlushData(const std::vector &kv_table_names, IncreaseWriteReqCount(shard_id); DataShard &ds_ref = data_shards_.at(shard_id); - auto shard_status = ds_ref.shard_status_.load(std::memory_order_acquire); + auto shard_status = ds_ref.shard_status_.load(std::memory_order_seq_cst); if (shard_status != DSShardStatus::ReadWrite) { DecreaseWriteReqCount(shard_id); @@ -759,7 +851,7 @@ void DataStoreService::DeleteRange( IncreaseWriteReqCount(shard_id); DataShard &ds_ref = data_shards_.at(shard_id); - auto shard_status = ds_ref.shard_status_.load(std::memory_order_acquire); + auto shard_status = ds_ref.shard_status_.load(std::memory_order_seq_cst); if (shard_status != DSShardStatus::ReadWrite) { DecreaseWriteReqCount(shard_id); @@ -809,7 +901,7 @@ void DataStoreService::DeleteRange(const std::string_view table_name, IncreaseWriteReqCount(shard_id); DataShard &ds_ref = data_shards_.at(shard_id); - auto shard_status = ds_ref.shard_status_.load(std::memory_order_acquire); + auto shard_status = ds_ref.shard_status_.load(std::memory_order_seq_cst); if (shard_status != DSShardStatus::ReadWrite) { DecreaseWriteReqCount(shard_id); @@ -863,7 +955,7 @@ void DataStoreService::CreateTable( IncreaseWriteReqCount(shard_id); DataShard &ds_ref = data_shards_.at(shard_id); - auto shard_status = ds_ref.shard_status_.load(std::memory_order_acquire); + auto shard_status = ds_ref.shard_status_.load(std::memory_order_seq_cst); if (shard_status != DSShardStatus::ReadWrite) { DecreaseWriteReqCount(shard_id); @@ -907,7 +999,7 @@ void DataStoreService::CreateTable(const std::string_view table_name, IncreaseWriteReqCount(shard_id); DataShard &ds_ref = data_shards_.at(shard_id); - auto shard_status = ds_ref.shard_status_.load(std::memory_order_acquire); + auto shard_status = ds_ref.shard_status_.load(std::memory_order_seq_cst); if (shard_status != DSShardStatus::ReadWrite) { DecreaseWriteReqCount(shard_id); @@ -953,7 +1045,7 @@ void DataStoreService::DropTable( IncreaseWriteReqCount(shard_id); DataShard &ds_ref = data_shards_.at(shard_id); - auto shard_status = ds_ref.shard_status_.load(std::memory_order_acquire); + auto shard_status = ds_ref.shard_status_.load(std::memory_order_seq_cst); if (shard_status != DSShardStatus::ReadWrite) { DecreaseWriteReqCount(shard_id); @@ -997,7 +1089,7 @@ void DataStoreService::DropTable(const std::string_view table_name, IncreaseWriteReqCount(shard_id); DataShard &ds_ref = data_shards_.at(shard_id); - auto shard_status = ds_ref.shard_status_.load(std::memory_order_acquire); + auto shard_status = ds_ref.shard_status_.load(std::memory_order_seq_cst); if (shard_status != DSShardStatus::ReadWrite) { DecreaseWriteReqCount(shard_id); @@ -1044,7 +1136,7 @@ void DataStoreService::BatchWriteRecords( IncreaseWriteReqCount(shard_id); DataShard &ds_ref = data_shards_.at(shard_id); - auto shard_status = ds_ref.shard_status_.load(std::memory_order_acquire); + auto shard_status = ds_ref.shard_status_.load(std::memory_order_seq_cst); if (shard_status != DSShardStatus::ReadWrite) { DecreaseWriteReqCount(shard_id); @@ -1100,7 +1192,8 @@ void DataStoreService::ScanNext( } DataShard &ds_ref = data_shards_.at(shard_id); - auto shard_status = ds_ref.shard_status_.load(std::memory_order_acquire); + auto scan_submit_guard = EnterReadSubmitWindow(shard_id); + auto shard_status = ds_ref.shard_status_.load(std::memory_order_seq_cst); if (shard_status != DSShardStatus::ReadWrite && shard_status != DSShardStatus::ReadOnly) { @@ -1148,7 +1241,8 @@ void DataStoreService::ScanNext(::google::protobuf::RpcController *controller, } DataShard &ds_ref = data_shards_.at(shard_id); - auto shard_status = ds_ref.shard_status_.load(std::memory_order_acquire); + auto scan_submit_guard = EnterReadSubmitWindow(shard_id); + auto shard_status = ds_ref.shard_status_.load(std::memory_order_seq_cst); if (shard_status != DSShardStatus::ReadWrite && shard_status != DSShardStatus::ReadOnly) { @@ -1182,7 +1276,8 @@ void DataStoreService::ScanClose(::google::protobuf::RpcController *controller, } DataShard &ds_ref = data_shards_.at(shard_id); - auto shard_status = ds_ref.shard_status_.load(std::memory_order_acquire); + auto scan_submit_guard = EnterReadSubmitWindow(shard_id); + auto shard_status = ds_ref.shard_status_.load(std::memory_order_seq_cst); if (shard_status != DSShardStatus::ReadWrite && shard_status != DSShardStatus::ReadOnly) { @@ -1217,7 +1312,8 @@ void DataStoreService::ScanClose(const std::string_view table_name, } DataShard &ds_ref = data_shards_.at(shard_id); - auto shard_status = ds_ref.shard_status_.load(std::memory_order_acquire); + auto scan_submit_guard = EnterReadSubmitWindow(shard_id); + auto shard_status = ds_ref.shard_status_.load(std::memory_order_seq_cst); if (shard_status != DSShardStatus::ReadWrite && shard_status != DSShardStatus::ReadOnly) { @@ -1338,7 +1434,7 @@ void DataStoreService::BatchWriteRecords( IncreaseWriteReqCount(shard_id); DataShard &ds_ref = data_shards_.at(shard_id); - auto shard_status = ds_ref.shard_status_.load(std::memory_order_acquire); + auto shard_status = ds_ref.shard_status_.load(std::memory_order_seq_cst); if (shard_status != DSShardStatus::ReadWrite) { DecreaseWriteReqCount(shard_id); @@ -1399,7 +1495,7 @@ void DataStoreService::CreateSnapshotForBackup( IncreaseWriteReqCount(shard_id); DataShard &ds_ref = data_shards_.at(shard_id); - auto shard_status = ds_ref.shard_status_.load(std::memory_order_acquire); + auto shard_status = ds_ref.shard_status_.load(std::memory_order_seq_cst); if (shard_status != DSShardStatus::ReadWrite) { DecreaseWriteReqCount(shard_id); @@ -1445,7 +1541,7 @@ void DataStoreService::CreateSnapshotForBackup( IncreaseWriteReqCount(shard_id); DataShard &ds_ref = data_shards_.at(shard_id); - auto shard_status = ds_ref.shard_status_.load(std::memory_order_acquire); + auto shard_status = ds_ref.shard_status_.load(std::memory_order_seq_cst); if (shard_status != DSShardStatus::ReadWrite) { DecreaseWriteReqCount(shard_id); @@ -3089,7 +3185,7 @@ bool DataStoreService::SwitchReadWriteToReadOnly(uint32_t shard_id) auto &ds_ref = data_shards_.at(shard_id); DSShardStatus expected = DSShardStatus::ReadWrite; if (!ds_ref.shard_status_.compare_exchange_strong( - expected, DSShardStatus::ReadOnly) && + expected, DSShardStatus::ReadOnly, std::memory_order_seq_cst) && expected != DSShardStatus::ReadOnly) { DLOG(ERROR) << "SwitchReadWriteToReadOnly failed, shard status is not " @@ -3098,11 +3194,11 @@ bool DataStoreService::SwitchReadWriteToReadOnly(uint32_t shard_id) } // wait for all write requests to finish - while (ds_ref.ongoing_write_requests_.load(std::memory_order_acquire) > 0) + while (ds_ref.ongoing_write_requests_.load(std::memory_order_seq_cst) > 0) { bthread_usleep(1000); } - if (ds_ref.shard_status_.load(std::memory_order_acquire) == + if (ds_ref.shard_status_.load(std::memory_order_seq_cst) == DSShardStatus::ReadOnly) { cluster_manager_.SwitchShardToReadOnly(shard_id, expected); @@ -3126,7 +3222,7 @@ bool DataStoreService::SwitchReadOnlyToClosed(uint32_t shard_id) auto &ds_ref = data_shards_.at(shard_id); DSShardStatus expected = DSShardStatus::ReadOnly; if (!ds_ref.shard_status_.compare_exchange_strong( - expected, DSShardStatus::Starting) && + expected, DSShardStatus::Starting, std::memory_order_seq_cst) && expected != DSShardStatus::Closed) { DLOG(ERROR) << "SwitchReadOnlyToClosed failed, shard status is not " @@ -3138,10 +3234,13 @@ bool DataStoreService::SwitchReadOnlyToClosed(uint32_t shard_id) { DLOG(INFO) << "SwitchReadOnlyToClosed enter shutdown, shard " << shard_id; + WaitReadSubmitWindowsDrained(ds_ref); ds_ref.data_store_->Shutdown(); DSShardStatus expected_after_shutdown = DSShardStatus::Starting; const bool switched = ds_ref.shard_status_.compare_exchange_strong( - expected_after_shutdown, DSShardStatus::Closed); + expected_after_shutdown, + DSShardStatus::Closed, + std::memory_order_seq_cst); CHECK(switched); cluster_manager_.SwitchShardToClosed(shard_id, DSShardStatus::ReadOnly); } diff --git a/store_handler/eloq_data_store_service/data_store_service.h b/store_handler/eloq_data_store_service/data_store_service.h index fdaf2c18..4cc62f98 100644 --- a/store_handler/eloq_data_store_service/data_store_service.h +++ b/store_handler/eloq_data_store_service/data_store_service.h @@ -25,6 +25,7 @@ #include #include +#include #include #include #include @@ -194,6 +195,9 @@ class TTLWrapperCache class DataStoreService : EloqDS::remote::DataStoreRpcService { public: + static constexpr uint32_t kMaxShardCount = 1000; + static constexpr uint32_t kReadSubmitSlotCount = 1000; + DataStoreService(const DataStoreServiceClusterManager &config, const std::string &config_file_path, const std::string &migration_log_path, @@ -648,13 +652,13 @@ class DataStoreService : EloqDS::remote::DataStoreRpcService void IncreaseWriteReqCount(uint32_t shard_id) { data_shards_.at(shard_id).ongoing_write_requests_.fetch_add( - 1, std::memory_order_release); + 1, std::memory_order_seq_cst); } void DecreaseWriteReqCount(uint32_t shard_id) { data_shards_.at(shard_id).ongoing_write_requests_.fetch_sub( - 1, std::memory_order_release); + 1, std::memory_order_seq_cst); } bool IsOwnerOfShard(uint32_t shard_id) const @@ -709,6 +713,32 @@ class DataStoreService : EloqDS::remote::DataStoreRpcService } private: + struct DataShard; + + struct alignas(64) ReadSubmitSlot + { + std::atomic depth{0}; + }; + + class ReadSubmitGuard + { + public: + ReadSubmitGuard() = default; + ReadSubmitGuard(DataStoreService *service, + uint32_t shard_id, + uint32_t slot_idx); + ReadSubmitGuard(const ReadSubmitGuard &) = delete; + ReadSubmitGuard &operator=(const ReadSubmitGuard &) = delete; + ReadSubmitGuard(ReadSubmitGuard &&other) noexcept; + ReadSubmitGuard &operator=(ReadSubmitGuard &&other) = delete; + ~ReadSubmitGuard(); + + private: + DataStoreService *service_{nullptr}; + uint32_t shard_id_{UINT32_MAX}; + uint32_t slot_idx_{UINT32_MAX}; + }; + DataStore *GetDataStore(uint32_t shard_id) { if (data_shards_.at(shard_id).shard_id_ == shard_id) @@ -729,6 +759,10 @@ class DataStoreService : EloqDS::remote::DataStoreRpcService bool SwitchReadWriteToReadOnly(uint32_t shard_id); bool SwitchReadOnlyToClosed(uint32_t shard_id); bool SwitchReadOnlyToReadWrite(uint32_t shard_id); + uint32_t GetReadSubmitSlotIndex(); + ReadSubmitGuard EnterReadSubmitWindow(uint32_t shard_id); + void LeaveReadSubmitWindow(uint32_t shard_id, uint32_t slot_idx); + void WaitReadSubmitWindowsDrained(const DataShard &ds_ref) const; bool WriteMigrationLog(uint32_t shard_id, const std::string &event_id, const std::string &target_node_ip, @@ -785,13 +819,15 @@ class DataStoreService : EloqDS::remote::DataStoreRpcService std::atomic latest_snapshot_ts_{0}; std::atomic latest_delete_archive_ts_{0}; std::unique_ptr scan_iter_cache_{nullptr}; + std::array read_submit_slots_; // Whether the file cache sync is running. Used to avoid concurrent // local ssd file operations between db and file sync worker. std::atomic is_file_sync_running_{false}; }; - std::array data_shards_; + std::array data_shards_; + std::atomic next_read_submit_slot_idx_{0}; std::unique_ptr data_store_factory_; diff --git a/store_handler/eloq_data_store_service/eloqstore b/store_handler/eloq_data_store_service/eloqstore index 123e2629..acd3cb37 160000 --- a/store_handler/eloq_data_store_service/eloqstore +++ b/store_handler/eloq_data_store_service/eloqstore @@ -1 +1 @@ -Subproject commit 123e2629d870ca8d047690c9daba8bf56414f812 +Subproject commit acd3cb376c71e1567a558ee6cd7eb5189f0fe69a