提交 67a69fea authored 作者: mooncake9527's avatar mooncake9527

update

上级 7815ac1e
package xcron
// func AddHandler(k string, h controller.Fn) {
// controller.Handlers.Add(k, h)
// }
import (
"context"
"fmt"
// func AddHandlerMap(handlers map[string]controller.Fn) {
// for k, h := range handlers {
// AddHandler(k, h)
// }
// }
"gitlab.wanzhuangkj.com/tush/xpkg/config"
"gitlab.wanzhuangkj.com/tush/xpkg/pkg/eventbus"
"gitlab.wanzhuangkj.com/tush/xpkg/pkg/gocron"
"gitlab.wanzhuangkj.com/tush/xpkg/pkg/logger"
"gitlab.wanzhuangkj.com/tush/xpkg/utils/ctxUtils"
"gitlab.wanzhuangkj.com/tush/xpkg/utils/sliceUtils"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/biz"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/controller"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/enums"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/service"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/types"
)
// func initCron() error {
// cronCfg := config.GetCron()
// if !cronCfg.Enable {
// return nil
// }
// ctx := context.TODO()
// cronJobs, err := service.CronJobService.GetSliceByEnable(ctx, enums.CronJob_Enable_OPEN, &biz.CronJobOpts{})
// if err != nil {
// return err
// }
// if len(cronJobs) == 0 {
// return nil
// }
// cronJobIDs := sliceUtils.GetIDs(cronJobs)
// _ = service.CronJobService.UpdateByIDs(ctx, &types.CronJobUpdateByIDsReq{IDs: cronJobIDs, State: enums.CronJob_State_NOT_RUNNING}, nil)
// var tasks []*gocron.Task
// for _, job := range cronJobs {
// job.ParseJobInfo()
// _, ok := controller.Handlers.Get(job.JobCode)
// if ok {
// tsk := gocron.Task{
// TimeSpec: job.TimeSpec,
// Name: job.JobCode,
// Fn: func() {
// ctx := context.WithValue(ctx, ctxUtils.HeaderXRequestIDKey, ctxUtils.GenerateTid())
// if err := controller.ExecJob(ctx, job); err != nil {
// logger.Error(fmt.Sprintf("[cron]exec err: %s", err.Error()), logger.Any("name", job.JobName), logger.Any("timeSpec", job.TimeSpec), logger.Any("args", job.Info), ctxUtils.CtxTraceIDField(ctx))
// }
// },
// }
// logger.Info("[cron]add cron job", logger.Any("name", job.JobName), logger.Any("timeSpec", job.TimeSpec), logger.Any("args", job.Info), ctxUtils.CtxTraceIDField(ctx))
// tasks = append(tasks, &tsk)
// }
// }
// if len(tasks) == 0 {
// return nil
// }
// if err = gocron.Init(gocron.WithLog(logger.Get(), false), gocron.WithGranularity(gocron.SecondType)); err != nil {
// return err
// }
// if err = gocron.Add(tasks...); err != nil {
// return err
// }
// eventbus.Eb.Publish(ctx, eventbus.TopicCronInitFinish)
// logger.Info("[cron] initialized", ctxUtils.CtxTraceIDField(ctx))
// return nil
// }
func AddHandler(k string, h controller.Fn) {
controller.Handlers.Add(k, h)
}
func AddHandlerMap(handlers map[string]controller.Fn) {
for k, h := range handlers {
AddHandler(k, h)
}
}
func initCron() error {
cronCfg := config.GetCron()
if !cronCfg.Enable {
return nil
}
ctx := context.TODO()
cronJobs, err := service.CronJobService.GetSliceByEnable(ctx, enums.CronJob_Enable_OPEN, &biz.CronJobOpts{})
if err != nil {
return err
}
if len(cronJobs) == 0 {
return nil
}
cronJobIDs := sliceUtils.GetIDs(cronJobs)
_ = service.CronJobService.UpdateByIDs(ctx, &types.CronJobUpdateByIDsReq{IDs: cronJobIDs, State: enums.CronJob_State_NOT_RUNNING}, nil)
var tasks []*gocron.Task
for _, job := range cronJobs {
job.ParseJobInfo()
_, ok := controller.Handlers.Get(job.JobCode)
if ok {
tsk := gocron.Task{
TimeSpec: job.TimeSpec,
Name: job.JobCode,
Fn: func() {
ctx := context.WithValue(ctx, ctxUtils.HeaderXRequestIDKey, ctxUtils.GenerateTid())
if err := controller.ExecJob(ctx, job); err != nil {
logger.Error(fmt.Sprintf("[cron]exec err: %s", err.Error()), logger.Any("name", job.JobName), logger.Any("timeSpec", job.TimeSpec), logger.Any("args", job.Info), ctxUtils.CtxTraceIDField(ctx))
}
},
}
logger.Info("[cron]add cron job", logger.Any("name", job.JobName), logger.Any("timeSpec", job.TimeSpec), logger.Any("args", job.Info), ctxUtils.CtxTraceIDField(ctx))
tasks = append(tasks, &tsk)
}
}
if len(tasks) == 0 {
return nil
}
if err = gocron.Init(gocron.WithLog(logger.Get(), false), gocron.WithGranularity(gocron.SecondType)); err != nil {
return err
}
if err = gocron.Add(tasks...); err != nil {
return err
}
eventbus.Eb.Publish(ctx, eventbus.TopicCronInitFinish)
logger.Info("[cron] initialized", ctxUtils.CtxTraceIDField(ctx))
return nil
}
......@@ -3,27 +3,36 @@ package xcron
import (
"context"
"github.com/jinzhu/copier"
"gitlab.wanzhuangkj.com/tush/xpkg/config"
"gitlab.wanzhuangkj.com/tush/xpkg/database"
"gitlab.wanzhuangkj.com/tush/xpkg/pkg/eventbus"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/cache"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/dao"
)
// 使用xxl
// func init() {
// _ = eventbus.Eb.Subscribe(eventbus.TopicCoreInitFinish, func(ctx context.Context) {
// cronCfg := config.GetCron()
// if !cronCfg.Enable {
// return
// }
// schema := cronCfg.Schema
// dao.CronJobDao.ODao.XDB = &database.XDB{Schema: schema}
// dao.CronJobLogDao.ODao.XDB = &database.XDB{Schema: schema}
// dao.CronJobRecordDao.ODao.XDB = &database.XDB{Schema: schema}
// dao.CronJobDao.ODao.Cache = cache.NewCronJobCache(database.GetCacheType())
// initCron()
// })
// }
func init() {
eventbus.Eb.Subscribe(eventbus.TopicGinEngineCreated, func(ctx context.Context) {
initXCron()
initLocalCron()
})
}
func initLocalCron() {
var cron config.Cron
config.Read(func(c *config.Config) {
_ = copier.Copy(&cron, c.Cron)
})
if !cron.Enable {
return
}
if cron.Type != "db" {
return
}
schema := cron.Schema
dao.CronJobDao.ODao.XDB = &database.XDB{Schema: schema}
dao.CronJobLogDao.ODao.XDB = &database.XDB{Schema: schema}
dao.CronJobRecordDao.ODao.XDB = &database.XDB{Schema: schema}
dao.CronJobDao.ODao.Cache = cache.NewCronJobCache(database.GetCacheType())
initCron()
}
......@@ -33,12 +33,6 @@ func initXCron() {
if !cron.Enable {
return
}
// if cron.Type == "" {
// cron.Type = "xxl"
// config.Write(func(c *config.Config) {
// c.Cron.Type = "xxl"
// })
// }
if cron.Type != "xxl" {
return
}
......
Markdown 格式
0%
您添加了 0 到此讨论。请谨慎行事。
请先完成此评论的编辑!
注册 或者 后发表评论