From f47f57ef3c7da3925b9b05596b275b390ddc0fee Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=96=B0=E4=BA=AE?= Date: Sun, 5 Sep 2021 17:53:50 +0800 Subject: [PATCH] =?UTF-8?q?feature(1.2.7):=20=E6=96=B0=E5=A2=9E=20cron=5Fs?= =?UTF-8?q?erver=20-=20=E5=90=8E=E5=8F=B0=E4=BB=BB=E5=8A=A1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 优化代码; --- internal/cron/cron_server/server.go | 12 +++++++++++- internal/cron/cron_server/service_add_job.go | 5 ++++- 2 files changed, 15 insertions(+), 2 deletions(-) diff --git a/internal/cron/cron_server/server.go b/internal/cron/cron_server/server.go index 5cce0f5..06eaaee 100644 --- a/internal/cron/cron_server/server.go +++ b/internal/cron/cron_server/server.go @@ -50,11 +50,21 @@ type server struct { type Server interface { i() + + // Start 启动 cron 服务 Start() + + // Stop 停止 cron 服务 Stop() + + // AddTask 增加定时任务 AddTask(task *cron_task_repo.CronTask) - AddJob(task *cron_task_repo.CronTask) cron.FuncJob + + // RemoveTask 删除定时任务 RemoveTask(taskId int) + + // AddJob 增加定时任务执行的工作内容 + AddJob(task *cron_task_repo.CronTask) cron.FuncJob } func New(logger *zap.Logger, db db.Repo, cache cache.Repo) (Server, error) { diff --git a/internal/cron/cron_server/service_add_job.go b/internal/cron/cron_server/service_add_job.go index 414e953..9d4c649 100644 --- a/internal/cron/cron_server/service_add_job.go +++ b/internal/cron/cron_server/service_add_job.go @@ -13,7 +13,10 @@ func (s *server) AddJob(task *cron_task_repo.CronTask) cron.FuncJob { s.taskCount.Add() defer s.taskCount.Done() - msg := fmt.Sprintf("开始执行任务:(%d)%s [%s]", task.Id, task.Name, task.Spec) + // 将 task 信息写入到 Kafka Topic 中,任务执行器订阅 Topic 如果为符合条件的任务并进行执行,反之不执行 + // 为了便于演示,不写入到 Kafka 中,仅记录日志 + + msg := fmt.Sprintf("执行任务:(%d)%s [%s]", task.Id, task.Name, task.Spec) s.logger.Info(msg) } }