MCPcopy Create free account
hub / github.com/4paradigm/OpenMLDB / ProcessTask

Method ProcessTask

src/nameserver/name_server_impl.cc:2070–2140  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

2068}
2069
2070void NameServerImpl::ProcessTask() {
2071 while (running_.load(std::memory_order_acquire)) {
2072 {
2073 bool has_task = false;
2074 std::unique_lock<std::mutex> lock(mu_);
2075 for (const auto& op_list : task_vec_) {
2076 if (!op_list.empty()) {
2077 has_task = true;
2078 break;
2079 }
2080 }
2081 if (!has_task) {
2082 cv_.wait_for(lock, std::chrono::milliseconds(FLAGS_name_server_task_wait_time));
2083 if (!running_.load(std::memory_order_acquire)) {
2084 PDLOG(WARNING, "cur nameserver is not leader");
2085 return;
2086 }
2087 }
2088
2089 for (const auto& op_list : task_vec_) {
2090 if (op_list.empty()) {
2091 continue;
2092 }
2093 std::shared_ptr<OPData> op_data = op_list.front();
2094 if (op_data->task_list_.empty() || op_data->op_info_.task_status() == ::fedb::api::kFailed ||
2095 op_data->op_info_.task_status() == ::fedb::api::kCanceled) {
2096 continue;
2097 }
2098 if (op_data->op_info_.task_status() == ::fedb::api::kInited) {
2099 op_data->op_info_.set_start_time(::baidu::common::timer::now_time());
2100 op_data->op_info_.set_task_status(::fedb::api::kDoing);
2101 std::string value;
2102 op_data->op_info_.SerializeToString(&value);
2103 std::string node = zk_op_data_path_ + "/" + std::to_string(op_data->op_info_.op_id());
2104 if (!zk_client_->SetNodeValue(node, value)) {
2105 PDLOG(WARNING, "set zk op status value failed. node[%s] value[%s]", node.c_str(),
2106 value.c_str());
2107 op_data->op_info_.set_task_status(::fedb::api::kInited);
2108 continue;
2109 }
2110 }
2111 std::shared_ptr<Task> task = op_data->task_list_.front();
2112 if (task->task_info_->status() == ::fedb::api::kFailed) {
2113 PDLOG(WARNING, "task[%s] run failed, terminate op[%s]. op_id[%lu]",
2114 ::fedb::api::TaskType_Name(task->task_info_->task_type()).c_str(),
2115 ::fedb::api::OPType_Name(task->task_info_->op_type()).c_str(), task->task_info_->op_id());
2116 } else if (task->task_info_->status() == ::fedb::api::kInited) {
2117 DEBUGLOG("run task. opid[%lu] op_type[%s] task_type[%s]", task->task_info_->op_id(),
2118 ::fedb::api::OPType_Name(task->task_info_->op_type()).c_str(),
2119 ::fedb::api::TaskType_Name(task->task_info_->task_type()).c_str());
2120 task_thread_pool_.AddTask(task->fun_);
2121 task->task_info_->set_status(::fedb::api::kDoing);
2122 } else if (task->task_info_->status() == ::fedb::api::kDoing) {
2123 if (::baidu::common::timer::now_time() - op_data->op_info_.start_time() >
2124 FLAGS_name_server_op_execute_timeout / 1000) {
2125 PDLOG(INFO,
2126 "The execution time of op is too long. "
2127 "opid[%lu] op_type[%s] cur task_type[%s] "

Callers

nothing calls this directly

Calls 4

to_stringFunction · 0.85
SetNodeValueMethod · 0.80
AddTaskMethod · 0.80
emptyMethod · 0.45

Tested by

no test coverage detected