Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
25 commits
Select commit Hold shift + click to select a range
1e42d4b
Add scheduler lease table
sitaowang1998 Apr 3, 2025
7dc9c87
Remove unused scheduler state
sitaowang1998 Apr 3, 2025
37fbc5f
Remove unnecessary scheudle state change
sitaowang1998 Apr 3, 2025
3d198d1
Format code
sitaowang1998 Apr 3, 2025
c1dc18f
Fix clang tidy
sitaowang1998 Apr 3, 2025
053a1bb
Merge branch 'main' into scheduler_state
sitaowang1998 Apr 5, 2025
6fcbbd3
Merge branch 'main' into scheduler_lease
sitaowang1998 Apr 7, 2025
d4e7f18
Merge branch 'scheduler_state' into scheduler_lease
sitaowang1998 Apr 7, 2025
648cd38
Add scheduler lease in get_ready_tasks and unlease in create_task_ins…
sitaowang1998 Apr 7, 2025
a818ec1
Remove scheduler state when creating scheduler
sitaowang1998 Apr 7, 2025
8062659
Remove scheduler state in get active scheduler
sitaowang1998 Apr 7, 2025
806cd55
Merge branch 'scheduler_state' into scheduler_lease
sitaowang1998 Apr 7, 2025
f426d97
Add lease expire in get_ready_tasks
sitaowang1998 Apr 7, 2025
f240a8c
Reformat file
sitaowang1998 Apr 7, 2025
da06e14
Fix typo
sitaowang1998 Apr 7, 2025
cc6742e
Remove scheduler lease table creation
sitaowang1998 Apr 7, 2025
a444a78
Merge branch 'scheduler_state' of github.com:sitaowang1998/spider int…
sitaowang1998 Apr 7, 2025
4642922
Merge branch 'scheduler_state' into scheduler_lease
sitaowang1998 Apr 7, 2025
16976b7
Revert "Remove scheduler lease table creation"
sitaowang1998 Apr 7, 2025
aa71a60
Add unit test for lease timeout
sitaowang1998 Apr 7, 2025
56a04ba
Fix clang tidy and improve REQUIRE checks
sitaowang1998 Apr 7, 2025
b728fed
Merge branch 'main' into scheduler_lease
sitaowang1998 Apr 7, 2025
4bca9e9
Fix the comment on job id in map
sitaowang1998 Apr 8, 2025
8a795ff
Add primary key in scheduler lease table
sitaowang1998 Apr 9, 2025
c81d25a
Merge branch 'scheduler_lease' of github.com:sitaowang1998/spider int…
sitaowang1998 Apr 9, 2025
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 4 additions & 2 deletions src/spider/scheduler/FifoPolicy.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -17,11 +17,13 @@

namespace spider::scheduler {
FifoPolicy::FifoPolicy(
boost::uuids::uuid const scheduler_id,
std::shared_ptr<core::MetadataStorage> const& metadata_store,
std::shared_ptr<core::DataStorage> const& data_store,
std::shared_ptr<core::StorageConnection> const& conn
)
: m_metadata_store{metadata_store},
: m_scheduler_id{scheduler_id},
m_metadata_store{metadata_store},
m_data_store{data_store},
m_conn{conn} {}

Expand Down Expand Up @@ -63,7 +65,7 @@ auto FifoPolicy::pop_next_task(std::string const& worker_addr)
}

auto FifoPolicy::fetch_tasks() -> void {
m_metadata_store->get_ready_tasks(*m_conn, &m_tasks);
m_metadata_store->get_ready_tasks(*m_conn, m_scheduler_id, &m_tasks);
m_metadata_store->get_task_timeout(*m_conn, &m_tasks);

// Sort tasks based on job creation time in descending order.
Expand Down
3 changes: 3 additions & 0 deletions src/spider/scheduler/FifoPolicy.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ namespace spider::scheduler {
class FifoPolicy final : public SchedulerPolicy {
public:
FifoPolicy(
boost::uuids::uuid scheduler_id,
std::shared_ptr<core::MetadataStorage> const& metadata_store,
std::shared_ptr<core::DataStorage> const& data_store,
std::shared_ptr<core::StorageConnection> const& conn
Expand All @@ -31,6 +32,8 @@ class FifoPolicy final : public SchedulerPolicy {

auto pop_next_task(std::string const& worker_addr) -> std::optional<boost::uuids::uuid>;

boost::uuids::uuid m_scheduler_id;

std::shared_ptr<core::MetadataStorage> m_metadata_store;
std::shared_ptr<core::DataStorage> m_data_store;
std::shared_ptr<core::StorageConnection> m_conn;
Expand Down
19 changes: 12 additions & 7 deletions src/spider/scheduler/scheduler.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -206,13 +206,6 @@ auto main(int argc, char** argv) -> int {
boost::uuids::random_generator gen;
boost::uuids::uuid const scheduler_id = gen();

// Start scheduler server
spider::core::StopToken stop_token;
std::shared_ptr<spider::scheduler::SchedulerPolicy> const policy
= std::make_shared<spider::scheduler::FifoPolicy>(metadata_store, data_store, conn);
spider::scheduler::SchedulerServer
server{port, policy, metadata_store, data_store, conn, stop_token};

// Register scheduler with storage
spider::core::Scheduler const scheduler{scheduler_id, scheduler_addr, port};
err = metadata_store->add_scheduler(*conn, scheduler);
Expand All @@ -221,6 +214,18 @@ auto main(int argc, char** argv) -> int {
return cStorageErr;
}

// Start scheduler server
spider::core::StopToken stop_token;
std::shared_ptr<spider::scheduler::SchedulerPolicy> const policy
= std::make_shared<spider::scheduler::FifoPolicy>(
scheduler_id,
metadata_store,
data_store,
conn
);
spider::scheduler::SchedulerServer
server{port, policy, metadata_store, data_store, conn, stop_token};

try {
// Start a thread that periodically updates the scheduler's heartbeat
std::thread heartbeat_thread{
Expand Down
7 changes: 5 additions & 2 deletions src/spider/storage/MetadataStorage.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -83,8 +83,11 @@ class MetadataStorage {
get_task_job_id(StorageConnection& conn, boost::uuids::uuid id, boost::uuids::uuid* job_id)
-> StorageErr
= 0;
virtual auto get_ready_tasks(StorageConnection& conn, std::vector<ScheduleTaskMetadata>* tasks)
-> StorageErr
virtual auto get_ready_tasks(
StorageConnection& conn,
boost::uuids::uuid scheduler_id,
std::vector<ScheduleTaskMetadata>* tasks
) -> StorageErr
= 0;
virtual auto set_task_state(StorageConnection& conn, boost::uuids::uuid id, TaskState state)
-> StorageErr
Expand Down
46 changes: 44 additions & 2 deletions src/spider/storage/mysql/MySqlStorage.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -1226,19 +1226,32 @@ auto MySqlMetadataStorage::get_task_job_id(
return StorageErr{};
}

constexpr int cLeaseExpireTime = 1000 * 10; // 10 ms
Comment thread
davidlion marked this conversation as resolved.

auto MySqlMetadataStorage::get_ready_tasks(
StorageConnection& conn,
boost::uuids::uuid scheduler_id,
std::vector<ScheduleTaskMetadata>* tasks
) -> StorageErr {
try {
// Remove timeout scheduler leases
std::unique_ptr<sql::PreparedStatement> lease_timeout_statement(
static_cast<MySqlConnection&>(conn)->prepareStatement(
"DELETE FROM `scheduler_leases` WHERE TIMESTAMPDIFF(MICROSECOND, "
"`lease_time`, CURRENT_TIMESTAMP()) > ?"
)
);
lease_timeout_statement->setInt(1, cLeaseExpireTime);
lease_timeout_statement->executeUpdate();

// Get all ready tasks from job that has not failed or cancelled
std::unique_ptr<sql::Statement> task_statement(
static_cast<MySqlConnection&>(conn)->createStatement()
);
std::unique_ptr<sql::ResultSet> const res{task_statement->executeQuery(
"SELECT `id`, `func_name`, `job_id` FROM `tasks` WHERE `state` = 'ready' "
"AND `job_id` NOT IN (SELECT `job_id` FROM `tasks` WHERE `state` = 'fail' OR "
"`state` = 'cancel')"
"`state` = 'cancel') AND `id` NOT IN (SELECT `task_id` FROM `scheduler_leases`)"
)};

if (res->rowsCount() == 0) {
Expand Down Expand Up @@ -1270,6 +1283,7 @@ auto MySqlMetadataStorage::get_ready_tasks(
"(SELECT `job_id` FROM `tasks` WHERE `state` = 'fail' OR `state` = 'cancel'))"
)};

// Get job metadata
while (job_res->next()) {
boost::uuids::uuid const job_id = read_id(job_res->getBinaryStream("id"));
boost::uuids::uuid const client_id = read_id(job_res->getBinaryStream("client_id"));
Expand All @@ -1285,6 +1299,10 @@ auto MySqlMetadataStorage::get_ready_tasks(
)
};
}
// Job id will not be in job_id_to_task_ids if the job's tasks are leased
if (job_id_to_task_ids.find(job_id) == job_id_to_task_ids.end()) {
continue;
}
for (boost::uuids::uuid const& task_id : job_id_to_task_ids[job_id]) {
new_tasks[task_id].set_client_id(client_id);
new_tasks[task_id].set_job_creation_time(optional_creation_time.value());
Expand All @@ -1301,7 +1319,8 @@ auto MySqlMetadataStorage::get_ready_tasks(
"`task_inputs`.`data_id` = `data`.`id` JOIN `data_locality` ON `data`.`id` "
"= `data_locality`.`id` WHERE `task_inputs`.`task_id` IN (SELECT `id` "
"FROM `tasks` WHERE `state` = 'ready' AND `job_id` NOT IN (SELECT `job_id` "
"FROM `tasks` WHERE `state` = 'fail' OR `state` = 'cancel'))"
"FROM `tasks` WHERE `state` = 'fail' OR `state` = 'cancel')) AND "
"`task_inputs`.`task_id` NOT IN (SELECT `task_id` FROM `scheduler_leases`)"
)};

while (locality_res->next()) {
Expand All @@ -1315,6 +1334,21 @@ auto MySqlMetadataStorage::get_ready_tasks(
}
}

// Add scheduler lease
std::unique_ptr<sql::PreparedStatement> lease_statement(
static_cast<MySqlConnection&>(conn)->prepareStatement(
"INSERT INTO `scheduler_leases` (`scheduler_id`, `task_id`) VALUES (?, ?)"
)
);
sql::bytes scheduler_id_bytes = uuid_get_bytes(scheduler_id);
for (auto const& [task_id, task] : new_tasks) {
sql::bytes task_id_bytes = uuid_get_bytes(task_id);
lease_statement->setBytes(1, &scheduler_id_bytes);
lease_statement->setBytes(2, &task_id_bytes);
lease_statement->addBatch();
}
lease_statement->executeBatch();

Comment thread
sitaowang1998 marked this conversation as resolved.
// Add all tasks to the output
absl::flat_hash_set<boost::uuids::uuid> task_ids;
for (ScheduleTaskMetadata const& task : *tasks) {
Expand Down Expand Up @@ -1458,6 +1492,14 @@ MySqlMetadataStorage::create_task_instance(StorageConnection& conn, TaskInstance
instance_statement->setBytes(1, &instance_id_bytes);
instance_statement->setBytes(2, &id_bytes);
instance_statement->executeUpdate();
// Remove task from scheduler leases
std::unique_ptr<sql::PreparedStatement> const lease_statement(
static_cast<MySqlConnection&>(conn)->prepareStatement(
"DELETE FROM `scheduler_leases` WHERE `task_id` = ?"
)
);
lease_statement->setBytes(1, &id_bytes);
lease_statement->executeUpdate();
} catch (sql::SQLException& e) {
static_cast<MySqlConnection&>(conn)->rollback();
return StorageErr{StorageErrType::OtherErr, e.what()};
Expand Down
7 changes: 5 additions & 2 deletions src/spider/storage/mysql/MySqlStorage.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -79,8 +79,11 @@ class MySqlMetadataStorage : public MetadataStorage {
-> StorageErr override;
auto get_task_job_id(StorageConnection& conn, boost::uuids::uuid id, boost::uuids::uuid* job_id)
-> StorageErr override;
auto get_ready_tasks(StorageConnection& conn, std::vector<ScheduleTaskMetadata>* tasks)
-> StorageErr override;
auto get_ready_tasks(
StorageConnection& conn,
boost::uuids::uuid scheduler_id,
std::vector<ScheduleTaskMetadata>* tasks
) -> StorageErr override;
auto set_task_state(StorageConnection& conn, boost::uuids::uuid id, TaskState state)
-> StorageErr override;
auto set_task_running(StorageConnection& conn, boost::uuids::uuid id) -> StorageErr override;
Expand Down
14 changes: 13 additions & 1 deletion src/spider/storage/mysql/mysql_stmt.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,16 @@ std::string const cCreateTaskInstanceTable = R"(CREATE TABLE IF NOT EXISTS `task
PRIMARY KEY (`id`)
))";

std::string const cCreateSchedulerLeaseTable = R"(CREATE TABLE IF NOT EXISTS `scheduler_leases` (
`scheduler_id` BINARY(16) NOT NULL,
`task_id` BINARY(16) NOT NULL,
`lease_time` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
CONSTRAINT `lease_scheduler_id` FOREIGN KEY (`scheduler_id`) REFERENCES `schedulers` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE,
CONSTRAINT `lease_task_id` FOREIGN KEY (`task_id`) REFERENCES `tasks` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE,
INDEX (`scheduler_id`),
PRIMARY KEY (`scheduler_id`, `task_id`)
))";

std::string const cCreateDataTable = R"(CREATE TABLE IF NOT EXISTS `data` (
`id` BINARY(16) NOT NULL,
`value` VARBINARY(999) NOT NULL,
Expand Down Expand Up @@ -153,7 +163,7 @@ std::string const cCreateTaskKVDataTable = R"(CREATE TABLE IF NOT EXISTS `task_k
CONSTRAINT `kv_data_task_id` FOREIGN KEY (`task_id`) REFERENCES `tasks` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE
))";

std::array<std::string const, 16> const cCreateStorage = {
std::array<std::string const, 17> const cCreateStorage = {
cCreateDriverTable, // drivers table must be created before data_ref_driver
cCreateSchedulerTable,
cCreateJobTable, // jobs table must be created before task
Expand All @@ -170,6 +180,8 @@ std::array<std::string const, 16> const cCreateStorage = {
cCreateTaskInputTable,
cCreateTaskDependencyTable,
cCreateTaskInstanceTable,
cCreateSchedulerLeaseTable // scheduler_lease table must be created after scheduler and
// task
};

std::string const cInsertJob = R"(INSERT INTO `jobs` (`id`, `client_id`) VALUES (?, ?))";
Expand Down
30 changes: 26 additions & 4 deletions tests/scheduler/test-SchedulerPolicy.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,13 @@ TEMPLATE_LIST_TEST_CASE(
= std::move(std::get<std::unique_ptr<spider::core::StorageConnection>>(conn_result));

boost::uuids::random_generator gen;

// Add scheduler
boost::uuids::uuid const scheduler_id = gen();
REQUIRE(metadata_store
->add_scheduler(*conn, spider::core::Scheduler{scheduler_id, "127.0.0.1", 8080})
.success());

boost::uuids::uuid const client_id = gen();
// Submit tasks
spider::core::Task const task_1{"task_1"};
Expand All @@ -62,7 +69,7 @@ TEMPLATE_LIST_TEST_CASE(
boost::uuids::uuid const job_id_2 = gen();
REQUIRE(metadata_store->add_job(*conn, job_id_2, client_id, graph_2).success());

spider::scheduler::FifoPolicy policy{metadata_store, data_store, conn};
spider::scheduler::FifoPolicy policy{scheduler_id, metadata_store, data_store, conn};

// Schedule the earlier task
std::optional<boost::uuids::uuid> optional_task_id = policy.schedule_next(gen(), "");
Expand Down Expand Up @@ -107,6 +114,13 @@ TEMPLATE_LIST_TEST_CASE(
= std::move(std::get<std::unique_ptr<spider::core::StorageConnection>>(conn_result));

boost::uuids::random_generator gen;

// Add scheduler
boost::uuids::uuid const scheduler_id = gen();
REQUIRE(metadata_store
->add_scheduler(*conn, spider::core::Scheduler{scheduler_id, "127.0.0.1", 8080})
.success());

boost::uuids::uuid const job_id = gen();
boost::uuids::uuid const client_id = gen();
// Submit task with hard locality
Expand All @@ -123,7 +137,7 @@ TEMPLATE_LIST_TEST_CASE(
graph.add_output_task(task.get_id());
REQUIRE(metadata_store->add_job(*conn, job_id, client_id, graph).success());

spider::scheduler::FifoPolicy policy{metadata_store, data_store, conn};
spider::scheduler::FifoPolicy policy{scheduler_id, metadata_store, data_store, conn};
// Schedule with wrong address
REQUIRE_FALSE(policy.schedule_next(gen(), "").has_value());
// Schedule with correct address
Expand Down Expand Up @@ -156,8 +170,15 @@ TEMPLATE_LIST_TEST_CASE(
std::shared_ptr<spider::core::StorageConnection> const conn
= std::move(std::get<std::unique_ptr<spider::core::StorageConnection>>(conn_result));

// Add task
boost::uuids::random_generator gen;

// Add scheduler
boost::uuids::uuid const scheduler_id = gen();
REQUIRE(metadata_store
->add_scheduler(*conn, spider::core::Scheduler{scheduler_id, "127.0.0.1", 8080})
.success());

// Add task
boost::uuids::uuid const job_id = gen();
boost::uuids::uuid const client_id = gen();
spider::core::Task task{"task"};
Expand All @@ -173,7 +194,8 @@ TEMPLATE_LIST_TEST_CASE(
graph.add_output_task(task.get_id());
REQUIRE(metadata_store->add_job(*conn, job_id, client_id, graph).success());

spider::scheduler::FifoPolicy policy{metadata_store, data_store, conn};
spider::scheduler::FifoPolicy policy{scheduler_id, metadata_store, data_store, conn};

// Schedule with wrong address
std::optional<boost::uuids::uuid> const optional_task_id = policy.schedule_next(gen(), "");
REQUIRE(optional_task_id.has_value());
Expand Down
15 changes: 13 additions & 2 deletions tests/scheduler/test-SchedulerServer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -50,8 +50,20 @@ TEMPLATE_LIST_TEST_CASE(
std::shared_ptr<spider::core::StorageConnection> const conn
= std::move(std::get<std::unique_ptr<spider::core::StorageConnection>>(conn_result));

// Add scheduler
boost::uuids::random_generator gen;
boost::uuids::uuid const scheduler_id = gen();
REQUIRE(metadata_store
->add_scheduler(*conn, spider::core::Scheduler{scheduler_id, "127.0.0.1", 8080})
.success());

std::shared_ptr<spider::scheduler::SchedulerPolicy> const policy
= std::make_shared<spider::scheduler::FifoPolicy>(metadata_store, data_store, conn);
= std::make_shared<spider::scheduler::FifoPolicy>(
scheduler_id,
metadata_store,
data_store,
conn
);

constexpr unsigned short cPort = 6021;
spider::core::StopToken stop_token;
Expand Down Expand Up @@ -79,7 +91,6 @@ TEMPLATE_LIST_TEST_CASE(
graph.add_dependency(parent_task.get_id(), child_task.get_id());
graph.add_input_task(parent_task.get_id());
graph.add_output_task(child_task.get_id());
boost::uuids::random_generator gen;
boost::uuids::uuid const job_id = gen();
REQUIRE(metadata_store->add_job(*conn, job_id, gen(), graph).success());

Expand Down
Loading