mirror of https://github.com/milvus-io/milvus.git
MS-562 Add JobMgr and TaskCreator in Scheduler
Former-commit-id: bbaf2b649843e86fc7bbd28ec61e5a990fc5951fpull/191/head
parent
91192242d9
commit
5a2743946b
|
@ -12,6 +12,7 @@ Please mark all change in change log and use the ticket from JIRA.
|
|||
- MS-557 - Merge Log.h
|
||||
- MS-556 - Add Job Definition in Scheduler
|
||||
- MS-558 - Refine status code
|
||||
- MS-562 - Add JobMgr and TaskCreator in Scheduler
|
||||
|
||||
## New Feature
|
||||
|
||||
|
|
|
@ -0,0 +1,89 @@
|
|||
// Licensed to the Apache Software Foundation (ASF) under one
|
||||
// or more contributor license agreements. See the NOTICE file
|
||||
// distributed with this work for additional information
|
||||
// regarding copyright ownership. The ASF licenses this file
|
||||
// to you under the Apache License, Version 2.0 (the
|
||||
// "License"); you may not use this file except in compliance
|
||||
// with the License. You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing,
|
||||
// software distributed under the License is distributed on an
|
||||
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
|
||||
// KIND, either express or implied. See the License for the
|
||||
// specific language governing permissions and limitations
|
||||
// under the License.
|
||||
|
||||
#include "JobMgr.h"
|
||||
#include "task/Task.h"
|
||||
#include "TaskCreator.h"
|
||||
|
||||
|
||||
namespace zilliz {
|
||||
namespace milvus {
|
||||
namespace scheduler {
|
||||
|
||||
using namespace engine;
|
||||
|
||||
JobMgr::JobMgr(ResourceMgrPtr res_mgr)
|
||||
: res_mgr_(std::move(res_mgr)) {}
|
||||
|
||||
void
|
||||
JobMgr::Start() {
|
||||
if (not running_) {
|
||||
worker_thread_ = std::thread(&JobMgr::worker_function, this);
|
||||
running_ = true;
|
||||
}
|
||||
}
|
||||
|
||||
void
|
||||
JobMgr::Stop() {
|
||||
if (running_) {
|
||||
this->Put(nullptr);
|
||||
worker_thread_.join();
|
||||
running_ = false;
|
||||
}
|
||||
}
|
||||
|
||||
void
|
||||
JobMgr::Put(const JobPtr &job) {
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(mutex_);
|
||||
queue_.push(job);
|
||||
}
|
||||
cv_.notify_one();
|
||||
}
|
||||
|
||||
void
|
||||
JobMgr::worker_function() {
|
||||
while (running_) {
|
||||
std::unique_lock<std::mutex> lock(mutex_);
|
||||
cv_.wait(lock, [this] { return !queue_.empty(); });
|
||||
auto job = queue_.front();
|
||||
queue_.pop();
|
||||
lock.unlock();
|
||||
if (job == nullptr) {
|
||||
break;
|
||||
}
|
||||
|
||||
auto tasks = build_task(job);
|
||||
auto disk_list = res_mgr_->GetDiskResources();
|
||||
if (!disk_list.empty()) {
|
||||
if (auto disk = disk_list[0].lock()) {
|
||||
for (auto &task : tasks) {
|
||||
disk->task_table().Put(task);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
std::vector<TaskPtr>
|
||||
JobMgr::build_task(const JobPtr &job) {
|
||||
return TaskCreator::Create(job);
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
}
|
|
@ -0,0 +1,78 @@
|
|||
// Licensed to the Apache Software Foundation (ASF) under one
|
||||
// or more contributor license agreements. See the NOTICE file
|
||||
// distributed with this work for additional information
|
||||
// regarding copyright ownership. The ASF licenses this file
|
||||
// to you under the Apache License, Version 2.0 (the
|
||||
// "License"); you may not use this file except in compliance
|
||||
// with the License. You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing,
|
||||
// software distributed under the License is distributed on an
|
||||
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
|
||||
// KIND, either express or implied. See the License for the
|
||||
// specific language governing permissions and limitations
|
||||
// under the License.
|
||||
#pragma once
|
||||
|
||||
#include <string>
|
||||
#include <vector>
|
||||
#include <list>
|
||||
#include <queue>
|
||||
#include <deque>
|
||||
#include <unordered_map>
|
||||
#include <thread>
|
||||
#include <mutex>
|
||||
#include <condition_variable>
|
||||
#include <memory>
|
||||
|
||||
#include "job/Job.h"
|
||||
#include "task/Task.h"
|
||||
#include "ResourceMgr.h"
|
||||
|
||||
|
||||
namespace zilliz {
|
||||
namespace milvus {
|
||||
namespace scheduler {
|
||||
|
||||
using engine::TaskPtr;
|
||||
using engine::ResourceMgrPtr;
|
||||
|
||||
class JobMgr {
|
||||
public:
|
||||
explicit
|
||||
JobMgr(ResourceMgrPtr res_mgr);
|
||||
|
||||
void
|
||||
Start();
|
||||
|
||||
void
|
||||
Stop();
|
||||
|
||||
public:
|
||||
void
|
||||
Put(const JobPtr &job);
|
||||
|
||||
private:
|
||||
void
|
||||
worker_function();
|
||||
|
||||
std::vector<TaskPtr>
|
||||
build_task(const JobPtr &job);
|
||||
|
||||
private:
|
||||
bool running_ = false;
|
||||
std::queue<JobPtr> queue_;
|
||||
|
||||
std::thread worker_thread_;
|
||||
|
||||
std::mutex mutex_;
|
||||
std::condition_variable cv_;
|
||||
|
||||
ResourceMgrPtr res_mgr_ = nullptr;
|
||||
};
|
||||
|
||||
}
|
||||
}
|
||||
}
|
|
@ -0,0 +1,65 @@
|
|||
// Licensed to the Apache Software Foundation (ASF) under one
|
||||
// or more contributor license agreements. See the NOTICE file
|
||||
// distributed with this work for additional information
|
||||
// regarding copyright ownership. The ASF licenses this file
|
||||
// to you under the Apache License, Version 2.0 (the
|
||||
// "License"); you may not use this file except in compliance
|
||||
// with the License. You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing,
|
||||
// software distributed under the License is distributed on an
|
||||
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
|
||||
// KIND, either express or implied. See the License for the
|
||||
// specific language governing permissions and limitations
|
||||
// under the License.
|
||||
|
||||
#include "TaskCreator.h"
|
||||
|
||||
|
||||
namespace zilliz {
|
||||
namespace milvus {
|
||||
namespace scheduler {
|
||||
|
||||
std::vector<TaskPtr>
|
||||
TaskCreator::Create(const JobPtr &job) {
|
||||
switch (job->type()) {
|
||||
case JobType::SEARCH: {
|
||||
return Create(std::static_pointer_cast<SearchJob>(job));
|
||||
}
|
||||
case JobType::DELETE: {
|
||||
return Create(std::static_pointer_cast<DeleteJob>(job));
|
||||
}
|
||||
default: {
|
||||
// TODO: error
|
||||
return std::vector<TaskPtr>();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
std::vector<TaskPtr>
|
||||
TaskCreator::Create(const SearchJobPtr &job) {
|
||||
std::vector<TaskPtr> tasks;
|
||||
for (auto &index_file : job->index_files()) {
|
||||
auto task = std::make_shared<XSearchTask>(index_file.second);
|
||||
tasks.emplace_back(task);
|
||||
}
|
||||
|
||||
return tasks;
|
||||
}
|
||||
|
||||
std::vector<TaskPtr>
|
||||
TaskCreator::Create(const DeleteJobPtr &job) {
|
||||
std::vector<TaskPtr> tasks;
|
||||
// auto task = std::make_shared<XDeleteTask>(job);
|
||||
// tasks.emplace_back(task);
|
||||
|
||||
return tasks;
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,61 @@
|
|||
// Licensed to the Apache Software Foundation (ASF) under one
|
||||
// or more contributor license agreements. See the NOTICE file
|
||||
// distributed with this work for additional information
|
||||
// regarding copyright ownership. The ASF licenses this file
|
||||
// to you under the Apache License, Version 2.0 (the
|
||||
// "License"); you may not use this file except in compliance
|
||||
// with the License. You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing,
|
||||
// software distributed under the License is distributed on an
|
||||
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
|
||||
// KIND, either express or implied. See the License for the
|
||||
// specific language governing permissions and limitations
|
||||
// under the License.
|
||||
#pragma once
|
||||
|
||||
#include <string>
|
||||
#include <vector>
|
||||
#include <list>
|
||||
#include <queue>
|
||||
#include <deque>
|
||||
#include <unordered_map>
|
||||
#include <thread>
|
||||
#include <mutex>
|
||||
#include <condition_variable>
|
||||
#include <memory>
|
||||
|
||||
#include "job/Job.h"
|
||||
#include "job/SearchJob.h"
|
||||
#include "job/DeleteJob.h"
|
||||
#include "task/Task.h"
|
||||
#include "task/SearchTask.h"
|
||||
#include "task/DeleteTask.h"
|
||||
|
||||
|
||||
namespace zilliz {
|
||||
namespace milvus {
|
||||
namespace scheduler {
|
||||
|
||||
using engine::TaskPtr;
|
||||
using engine::XSearchTask;
|
||||
using engine::XDeleteTask;
|
||||
|
||||
class TaskCreator {
|
||||
public:
|
||||
static std::vector<TaskPtr>
|
||||
Create(const JobPtr &job);
|
||||
|
||||
public:
|
||||
static std::vector<TaskPtr>
|
||||
Create(const SearchJobPtr &job);
|
||||
|
||||
static std::vector<TaskPtr>
|
||||
Create(const DeleteJobPtr &job);
|
||||
};
|
||||
|
||||
}
|
||||
}
|
||||
}
|
|
@ -14,6 +14,7 @@
|
|||
// KIND, either express or implied. See the License for the
|
||||
// specific language governing permissions and limitations
|
||||
// under the License.
|
||||
#pragma once
|
||||
|
||||
#include <string>
|
||||
#include <vector>
|
||||
|
|
|
@ -14,6 +14,7 @@
|
|||
// KIND, either express or implied. See the License for the
|
||||
// specific language governing permissions and limitations
|
||||
// under the License.
|
||||
#pragma once
|
||||
|
||||
#include <string>
|
||||
#include <vector>
|
||||
|
|
|
@ -14,6 +14,7 @@
|
|||
// KIND, either express or implied. See the License for the
|
||||
// specific language governing permissions and limitations
|
||||
// under the License.
|
||||
#pragma once
|
||||
|
||||
#include <string>
|
||||
#include <vector>
|
||||
|
@ -78,6 +79,11 @@ public:
|
|||
return vectors_;
|
||||
}
|
||||
|
||||
Id2IndexMap &
|
||||
index_files() {
|
||||
return index_files_;
|
||||
}
|
||||
|
||||
private:
|
||||
uint64_t topk_ = 0;
|
||||
uint64_t nq_ = 0;
|
||||
|
|
Loading…
Reference in New Issue