From be6644b11dc7008bca8c99e77ede73e909ad3d05 Mon Sep 17 00:00:00 2001 From: xufei Date: Mon, 21 Apr 2025 16:22:42 +0800 Subject: [PATCH 1/2] save work Signed-off-by: xufei --- dbms/src/Flash/EstablishCall.cpp | 5 +++++ dbms/src/Flash/EstablishCall.h | 4 ++++ dbms/src/Flash/Mpp/MPPTaskManager.cpp | 10 +++++----- dbms/src/Flash/Mpp/MPPTaskManager.h | 2 +- 4 files changed, 15 insertions(+), 6 deletions(-) diff --git a/dbms/src/Flash/EstablishCall.cpp b/dbms/src/Flash/EstablishCall.cpp index 682b2eaa98d..711292ce6b4 100644 --- a/dbms/src/Flash/EstablishCall.cpp +++ b/dbms/src/Flash/EstablishCall.cpp @@ -216,6 +216,11 @@ void EstablishCallData::initRpc() } } +grpc::Alarm * EstablishCallData::getAlarm() +{ + return &alarm; +} + void EstablishCallData::tryConnectTunnel() { auto * task_manager = service->getContext()->getTMTContext().getMPPTaskManager().get(); diff --git a/dbms/src/Flash/EstablishCall.h b/dbms/src/Flash/EstablishCall.h index b015b96b632..41283968862 100644 --- a/dbms/src/Flash/EstablishCall.h +++ b/dbms/src/Flash/EstablishCall.h @@ -19,6 +19,7 @@ #include #include #include +#include #include #include @@ -95,6 +96,8 @@ class EstablishCallData final grpc::ServerContext * getGrpcContext() { return &ctx; } String getResourceGroupName() const { return resource_group_name; } + grpc::Alarm * getAlarm(); + private: /// WARNING: Since a event from one grpc completion queue may be handled by different @@ -142,5 +145,6 @@ class EstablishCallData final String resource_group_name; String connection_id; double waiting_task_time_ms = 0; + grpc::Alarm alarm{}; }; } // namespace DB diff --git a/dbms/src/Flash/Mpp/MPPTaskManager.cpp b/dbms/src/Flash/Mpp/MPPTaskManager.cpp index 9a0727b51f2..10d4a986278 100644 --- a/dbms/src/Flash/Mpp/MPPTaskManager.cpp +++ b/dbms/src/Flash/Mpp/MPPTaskManager.cpp @@ -54,7 +54,7 @@ void MPPGatherTaskSet::cancelAlarmsBySenderTaskId(const MPPTaskId & task_id) if (alarm_it != alarms.end()) { for (auto & alarm : alarm_it->second) - alarm.second.Cancel(); + alarm.second->Cancel(); alarms.erase(alarm_it); } } @@ -185,13 +185,14 @@ std::pair MPPTaskManager::findAsyncTunnel( context.getSettingsRef().auto_spill_check_min_interval_ms.get()); if (gather_task_set == nullptr) gather_task_set = query->addMPPGatherTaskSet(id.gather_id); - auto & alarm = gather_task_set->alarms[sender_task_id][receiver_task_id]; + auto * alarm = call_data->getAlarm(); call_data->setCallStateAndUpdateMetrics( EstablishCallData::WAIT_TUNNEL, GET_METRIC(tiflash_establish_calldata_count, type_wait_tunnel_calldata)); + gather_task_set->alarms[sender_task_id][receiver_task_id] = alarm; if likely (cq != nullptr) { - alarm.Set(cq, Clock::now() + std::chrono::seconds(10), call_data->asGRPCKickTag()); + alarm->Set(cq, Clock::now() + std::chrono::seconds(10), call_data->asGRPCKickTag()); } return {nullptr, ""}; } @@ -220,7 +221,6 @@ std::pair MPPTaskManager::findAsyncTunnel( } } /// don't need to delete the alarm here because registerMPPTask will delete all the related alarm - return task->getTunnel(request); } @@ -315,7 +315,7 @@ void MPPTaskManager::abortMPPGather(const MPPGatherId & gather_id, const String for (auto & alarms_per_task : gather_task_set->alarms) { for (auto & alarm : alarms_per_task.second) - alarm.second.Cancel(); + alarm.second->Cancel(); } gather_task_set->alarms.clear(); if (!gather_task_set->hasMPPTask()) diff --git a/dbms/src/Flash/Mpp/MPPTaskManager.h b/dbms/src/Flash/Mpp/MPPTaskManager.h index a46d33b2df2..b726d62d046 100644 --- a/dbms/src/Flash/Mpp/MPPTaskManager.h +++ b/dbms/src/Flash/Mpp/MPPTaskManager.h @@ -42,7 +42,7 @@ struct MPPGatherTaskSet State state = Normal; String error_message; /// > - std::unordered_map> alarms; + std::unordered_map> alarms; /// only used in scheduler std::queue waiting_tasks; bool isInNormalState() const { return state == Normal; } From 202c1eca6c2db342c11ddb47ae10bd75d8f2725f Mon Sep 17 00:00:00 2001 From: xufei Date: Tue, 22 Apr 2025 15:59:42 +0800 Subject: [PATCH 2/2] save work Signed-off-by: xufei --- dbms/src/Flash/EstablishCall.cpp | 4 ++-- dbms/src/Flash/EstablishCall.h | 2 +- dbms/src/Flash/Mpp/MPPTaskManager.cpp | 10 +++++----- dbms/src/Flash/Mpp/MPPTaskManager.h | 3 ++- 4 files changed, 10 insertions(+), 9 deletions(-) diff --git a/dbms/src/Flash/EstablishCall.cpp b/dbms/src/Flash/EstablishCall.cpp index 711292ce6b4..9e63134febe 100644 --- a/dbms/src/Flash/EstablishCall.cpp +++ b/dbms/src/Flash/EstablishCall.cpp @@ -216,9 +216,9 @@ void EstablishCallData::initRpc() } } -grpc::Alarm * EstablishCallData::getAlarm() +grpc::Alarm & EstablishCallData::getAlarm() { - return &alarm; + return alarm; } void EstablishCallData::tryConnectTunnel() diff --git a/dbms/src/Flash/EstablishCall.h b/dbms/src/Flash/EstablishCall.h index 41283968862..4067dfe6dbc 100644 --- a/dbms/src/Flash/EstablishCall.h +++ b/dbms/src/Flash/EstablishCall.h @@ -96,7 +96,7 @@ class EstablishCallData final grpc::ServerContext * getGrpcContext() { return &ctx; } String getResourceGroupName() const { return resource_group_name; } - grpc::Alarm * getAlarm(); + grpc::Alarm & getAlarm(); private: diff --git a/dbms/src/Flash/Mpp/MPPTaskManager.cpp b/dbms/src/Flash/Mpp/MPPTaskManager.cpp index 10d4a986278..3296da53c37 100644 --- a/dbms/src/Flash/Mpp/MPPTaskManager.cpp +++ b/dbms/src/Flash/Mpp/MPPTaskManager.cpp @@ -54,7 +54,7 @@ void MPPGatherTaskSet::cancelAlarmsBySenderTaskId(const MPPTaskId & task_id) if (alarm_it != alarms.end()) { for (auto & alarm : alarm_it->second) - alarm.second->Cancel(); + alarm.second.get().Cancel(); alarms.erase(alarm_it); } } @@ -185,14 +185,14 @@ std::pair MPPTaskManager::findAsyncTunnel( context.getSettingsRef().auto_spill_check_min_interval_ms.get()); if (gather_task_set == nullptr) gather_task_set = query->addMPPGatherTaskSet(id.gather_id); - auto * alarm = call_data->getAlarm(); + auto & alarm = call_data->getAlarm(); call_data->setCallStateAndUpdateMetrics( EstablishCallData::WAIT_TUNNEL, GET_METRIC(tiflash_establish_calldata_count, type_wait_tunnel_calldata)); - gather_task_set->alarms[sender_task_id][receiver_task_id] = alarm; + gather_task_set->alarms[sender_task_id].emplace(receiver_task_id, std::ref(alarm)); if likely (cq != nullptr) { - alarm->Set(cq, Clock::now() + std::chrono::seconds(10), call_data->asGRPCKickTag()); + alarm.Set(cq, Clock::now() + std::chrono::seconds(10), call_data->asGRPCKickTag()); } return {nullptr, ""}; } @@ -315,7 +315,7 @@ void MPPTaskManager::abortMPPGather(const MPPGatherId & gather_id, const String for (auto & alarms_per_task : gather_task_set->alarms) { for (auto & alarm : alarms_per_task.second) - alarm.second->Cancel(); + alarm.second.get().Cancel(); } gather_task_set->alarms.clear(); if (!gather_task_set->hasMPPTask()) diff --git a/dbms/src/Flash/Mpp/MPPTaskManager.h b/dbms/src/Flash/Mpp/MPPTaskManager.h index b726d62d046..f49d14600ab 100644 --- a/dbms/src/Flash/Mpp/MPPTaskManager.h +++ b/dbms/src/Flash/Mpp/MPPTaskManager.h @@ -24,6 +24,7 @@ #include #include +#include #include #include #include @@ -42,7 +43,7 @@ struct MPPGatherTaskSet State state = Normal; String error_message; /// > - std::unordered_map> alarms; + std::unordered_map>> alarms; /// only used in scheduler std::queue waiting_tasks; bool isInNormalState() const { return state == Normal; }