From cb6c2cbcc7e40158f257ffd64187df6cf71f4528 Mon Sep 17 00:00:00 2001 From: xufei Date: Wed, 23 Apr 2025 03:23:02 +0800 Subject: [PATCH] This is an automated cherry-pick of #10124 Signed-off-by: ti-chi-bot --- dbms/src/Flash/EstablishCall.cpp | 5 +++++ dbms/src/Flash/EstablishCall.h | 4 ++++ dbms/src/Flash/Mpp/MPPTaskManager.cpp | 13 ++++++++++--- dbms/src/Flash/Mpp/MPPTaskManager.h | 3 ++- 4 files changed, 21 insertions(+), 4 deletions(-) diff --git a/dbms/src/Flash/EstablishCall.cpp b/dbms/src/Flash/EstablishCall.cpp index 99e8959a0d9..c360bc24250 100644 --- a/dbms/src/Flash/EstablishCall.cpp +++ b/dbms/src/Flash/EstablishCall.cpp @@ -180,6 +180,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 469b3241543..e24285e11b4 100644 --- a/dbms/src/Flash/EstablishCall.h +++ b/dbms/src/Flash/EstablishCall.h @@ -19,6 +19,7 @@ #include #include #include +#include #include namespace DB @@ -81,6 +82,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 @@ -138,5 +141,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 a548c52f4b0..72506c95609 100644 --- a/dbms/src/Flash/Mpp/MPPTaskManager.cpp +++ b/dbms/src/Flash/Mpp/MPPTaskManager.cpp @@ -56,7 +56,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); } } @@ -247,8 +247,16 @@ 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); +<<<<<<< HEAD auto & alarm = gather_task_set->alarms[sender_task_id][receiver_task_id]; call_data->setToWaitingTunnelState(); +======= + 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].emplace(receiver_task_id, std::ref(alarm)); +>>>>>>> 262b942077 (Avoid data race in grpc alarm (#10124)) if likely (cq != nullptr) { alarm.Set(cq, Clock::now() + std::chrono::seconds(10), call_data->asGRPCKickTag()); @@ -280,7 +288,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); } @@ -375,7 +382,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 e6e9fcb543f..118632fc0c8 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; }