| 282 | } |
| 283 | |
| 284 | void tryQueueQuery(ContextMutablePtr context, ASTPtr & query_ast) |
| 285 | { |
| 286 | ASTType ast_type; |
| 287 | try { |
| 288 | ast_type = query_ast->getType(); |
| 289 | } catch (...) { |
| 290 | LOG_DEBUG(&Poco::Logger::get("executeQuery"), "only queue dml query"); |
| 291 | return; |
| 292 | } |
| 293 | auto worker_group_handler = context->tryGetCurrentWorkerGroup(); |
| 294 | if (ast_type != ASTType::ASTSelectQuery && ast_type != ASTType::ASTSelectWithUnionQuery && ast_type != ASTType::ASTInsertQuery |
| 295 | && ast_type != ASTType::ASTDeleteQuery && ast_type != ASTType::ASTUpdateQuery) |
| 296 | { |
| 297 | LOG_DEBUG(&Poco::Logger::get("executeQuery"), "only queue dml query"); |
| 298 | return; |
| 299 | } |
| 300 | if (worker_group_handler) |
| 301 | { |
| 302 | Stopwatch queue_watch; |
| 303 | queue_watch.start(); |
| 304 | auto query_queue = context->getQueueManager(); |
| 305 | auto query_id = context->getCurrentQueryId(); |
| 306 | const auto & vw_name = worker_group_handler->getVWName(); |
| 307 | const auto & wg_name = worker_group_handler->getID(); |
| 308 | context->getWorkerStatusManager()->updateVWWorkerList(worker_group_handler->getHostWithPortsVec(), vw_name, wg_name); |
| 309 | auto queue_info = std::make_shared<QueueInfo>(query_id, vw_name, wg_name, context); |
| 310 | auto queue_result = query_queue->enqueue(queue_info, context->getSettingsRef().query_queue_timeout_ms); |
| 311 | if (queue_result == QueueResultStatus::QueueSuccess) |
| 312 | { |
| 313 | auto current_vw = context->tryGetCurrentVW(); |
| 314 | if (current_vw) |
| 315 | { |
| 316 | context->setCurrentWorkerGroup(current_vw->getWorkerGroup(wg_name)); |
| 317 | } |
| 318 | LOG_DEBUG(&Poco::Logger::get("executeQuery"), "query queue run time : {} ms", queue_watch.elapsedMilliseconds()); |
| 319 | } |
| 320 | else |
| 321 | { |
| 322 | LOG_ERROR(&Poco::Logger::get("executeQuery"), "query queue result : {}", queueResultStatusToString(queue_result)); |
| 323 | throw Exception( |
| 324 | ErrorCodes::CNCH_QUEUE_QUERY_FAILURE, |
| 325 | "query queue failed for query_id {}: {}", |
| 326 | query_id, |
| 327 | queueResultStatusToString(queue_result)); |
| 328 | } |
| 329 | } |
| 330 | } |
| 331 | |
| 332 | bool needThrowRootCauseError(const Context * context, int & error_code, String & error_messge) |
| 333 | { |
no test coverage detected