fix: image cache use specific workers (#22047)

Co-authored-by: Qiu Jian <qiujian@yunionyun.com>
This commit is contained in:
Jian Qiu
2025-02-05 11:33:41 +08:00
committed by GitHub
parent 8ca9fd5640
commit 2677a402fc
9 changed files with 31 additions and 12 deletions
+4
View File
@@ -70,6 +70,10 @@ func init() {
cmd.Get("ipmi", &options.BaseIdOptions{})
cmd.Get("vnc", &options.BaseIdOptions{})
cmd.Get("app-options", &options.BaseIdOptions{})
cmd.GetWithCustomShow("worker-stats", func(data jsonutils.JSONObject) {
stats, _ := data.GetArray("workers")
printList(&printutils.ListResult{Data: stats}, nil)
}, &options.BaseIdOptions{})
cmd.Get("tap-config", &options.BaseIdOptions{})
cmd.GetWithCustomShow("nics", func(data jsonutils.JSONObject) {
results := printutils.ListResult{}
@@ -232,5 +232,5 @@ func taskCompleted(ctx context.Context, data jsonutils.JSONObject) {
}
func InitPlaybookWorker() {
PlaybookWorker = workmanager.NewWorkManger(taskFailed, taskCompleted, options.Options.PlaybookWorkerCount)
PlaybookWorker = workmanager.NewWorkManger("PlaybookWorker", taskFailed, taskCompleted, options.Options.PlaybookWorkerCount)
}
+1 -1
View File
@@ -54,6 +54,6 @@ func performImageCache(
return
}
hostutils.DelayTask(ctx, performTask, disk)
hostutils.DelayImageCacheTask(ctx, performTask, disk)
hostutils.ResponseOk(ctx, w)
}
+5 -2
View File
@@ -189,7 +189,10 @@ func (w *SWorkManager) Stop() {
}
}
func NewWorkManger(onFailed OnTaskFailed, onCompleted OnTaskCompleted, workerCount int) *SWorkManager {
func NewWorkManger(name string, onFailed OnTaskFailed, onCompleted OnTaskCompleted, workerCount int) *SWorkManager {
if len(name) == 0 {
name = "RequestWorker"
}
if workerCount <= 0 {
workerCount = 1
}
@@ -197,6 +200,6 @@ func NewWorkManger(onFailed OnTaskFailed, onCompleted OnTaskCompleted, workerCou
onFailed: onFailed,
onCompleted: onCompleted,
worker: appsrv.NewWorkerManager(
"RequestWorker", workerCount, appsrv.DEFAULT_BACKLOG, false),
name, workerCount, appsrv.DEFAULT_BACKLOG, false),
}
}
+4
View File
@@ -7206,6 +7206,10 @@ func (h *SHost) GetDetailsAppOptions(ctx context.Context, userCred mcclient.Toke
return h.Request(ctx, userCred, httputils.GET, "/app-options", nil, nil)
}
func (h *SHost) GetDetailsWorkerStats(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) (jsonutils.JSONObject, error) {
return h.Request(ctx, userCred, httputils.GET, "/worker_stats", nil, nil)
}
func (hh *SHost) IsAttach2Wire(wireId string) bool {
netifs := hh.getNetifsOnWire(wireId)
return len(netifs) > 0
-4
View File
@@ -2408,16 +2408,12 @@ func (h *SHostInfo) OnCatalogChanged(catalog mcclient.KeystoneServiceCatalogV3)
}
}
if !reflect.DeepEqual(telegraf.GetConf(), conf) || (!strings.Contains(svcs, "telegraf") && !telegraf.IsActive()) {
log.Infof("telegraf configuration change, to reload ...")
log.Debugf("telegraf config: %s", conf)
telegraf.SetConf(conf)
if !strings.Contains(svcs, "telegraf") {
telegraf.BgReload(conf)
} else {
telegraf.BgReloadConf(conf)
}
} else {
log.Debugf("telegraf configuration no change")
}
/*urls, _ = catalog.GetServiceURLs("elasticsearch",
+12 -2
View File
@@ -203,6 +203,7 @@ func Response(ctx context.Context, w http.ResponseWriter, res interface{}) {
var (
wm *workmanager.SWorkManager
imageCacheW *workmanager.SWorkManager
k8sWm *workmanager.SWorkManager
ParamsError = fmt.Errorf("Delay task parse params error")
)
@@ -215,6 +216,10 @@ func DelayTask(ctx context.Context, task workmanager.DelayTaskFunc, params inter
wm.DelayTask(ctx, task, params)
}
func DelayImageCacheTask(ctx context.Context, task workmanager.DelayTaskFunc, params interface{}) {
imageCacheW.DelayTask(ctx, task, params)
}
func DelayKubeTask(ctx context.Context, task workmanager.DelayTaskFunc, params interface{}) {
k8sWm.DelayTask(ctx, task, params)
}
@@ -232,14 +237,19 @@ func DelayTaskWithWorker(
func InitWorkerManager() {
InitWorkerManagerWithCount(options.HostOptions.DefaultRequestWorkerCount)
initImageCacheWorkerManager()
}
func initImageCacheWorkerManager() {
imageCacheW = workmanager.NewWorkManger("ImageCacheDelayTaskWorkers", TaskFailed, TaskComplete, options.HostOptions.DefaultRequestWorkerCount)
}
func InitWorkerManagerWithCount(count int) {
wm = workmanager.NewWorkManger(TaskFailed, TaskComplete, count)
wm = workmanager.NewWorkManger("GeneralDelayedTaskWorkers", TaskFailed, TaskComplete, count)
}
func InitK8sWorkerManager() {
k8sWm = workmanager.NewWorkManger(K8sTaskFailed, K8sTaskComplete, options.HostOptions.DefaultRequestWorkerCount)
k8sWm = workmanager.NewWorkManger("K8sDelayedTaskWorkers", K8sTaskFailed, K8sTaskComplete, options.HostOptions.DefaultRequestWorkerCount)
}
func Init() {
@@ -130,7 +130,7 @@ func performImageCache(
}
}
hostutils.DelayTask(ctx, performTask, disk)
hostutils.DelayImageCacheTask(ctx, performTask, disk)
hostutils.ResponseOk(ctx, w)
}
+3 -1
View File
@@ -168,6 +168,8 @@ func (s *STelegraf) GetConfig(kwargs map[string]interface{}) string {
}
ignorePathSegments := []string{
"/run/k3s/containerd/",
"/run/onecloud/containerd/",
"/var/lib/",
}
ignorePathSegments = append(ignorePathSegments, kwargs["server_path"].(string))
for i := range ignorePathSegments {
@@ -175,7 +177,7 @@ func (s *STelegraf) GetConfig(kwargs map[string]interface{}) string {
}
conf += " ignore_mount_points = [" + strings.Join(ignoreMountPoints, ", ") + "]\n"
conf += " ignore_path_segments = [" + strings.Join(ignorePathSegments, ", ") + "]\n"
conf += " ignore_fs = [\"tmpfs\", \"devtmpfs\", \"overlay\", \"squashfs\", \"iso9660\", \"rootfs\", \"hugetlbfs\", \"autofs\"]\n"
conf += " ignore_fs = [\"devtmpfs\", \"devfs\", \"overlayfs\", \"overlay\", \"squashfs\", \"iso9660\", \"rootfs\", \"hugetlbfs\", \"autofs\", \"aufs\"]\n"
conf += "\n"
conf += "[[inputs.diskio]]\n"
conf += " skip_serial_number = false\n"