Merge pull request #13265 from swordqiu/hotfix/qj-batch-migrate-in-parrallel

fix: batch migrate in parrallel
This commit is contained in:
Zexi Li
2022-01-21 10:12:40 +08:00
committed by GitHub
6 changed files with 76 additions and 46 deletions
+15 -8
View File
@@ -4448,7 +4448,8 @@ func (manager *SGuestManager) PerformBatchMigrate(ctx context.Context, userCred
guests[i] = *guest
}
var hostGuests = map[string][]*api.GuestBatchMigrateParams{}
var hostGuests = map[string][]*SGuest{}
var hostGuestParams = map[string][]*api.GuestBatchMigrateParams{}
for i := 0; i < len(guests); i++ {
bmp := &api.GuestBatchMigrateParams{
Id: guests[i].Id,
@@ -4459,28 +4460,34 @@ func (manager *SGuestManager) PerformBatchMigrate(ctx context.Context, userCred
}
guests[i].SetStatus(userCred, api.VM_START_MIGRATE, "batch migrate")
if _, ok := hostGuests[guests[i].HostId]; ok {
hostGuests[guests[i].HostId] = append(hostGuests[guests[i].HostId], bmp)
hostGuests[guests[i].HostId] = append(hostGuests[guests[i].HostId], &guests[i])
hostGuestParams[guests[i].HostId] = append(hostGuestParams[guests[i].HostId], bmp)
} else {
hostGuests[guests[i].HostId] = []*api.GuestBatchMigrateParams{bmp}
hostGuests[guests[i].HostId] = []*SGuest{&guests[i]}
hostGuestParams[guests[i].HostId] = []*api.GuestBatchMigrateParams{bmp}
}
}
for hostId, params := range hostGuests {
for hostId, guests := range hostGuests {
params := hostGuestParams[hostId]
kwargs := jsonutils.NewDict()
kwargs.Set("guests", jsonutils.Marshal(params))
if len(preferHostId) > 0 {
kwargs.Set("prefer_host_id", jsonutils.NewString(preferHostId))
}
host := HostManager.FetchHostById(hostId)
manager.StartHostGuestsMigrateTask(ctx, userCred, host, kwargs, "")
manager.StartHostGuestsMigrateTask(ctx, userCred, guests, kwargs, "")
}
return nil, nil
}
func (manager *SGuestManager) StartHostGuestsMigrateTask(
ctx context.Context, userCred mcclient.TokenCredential,
host *SHost, kwargs *jsonutils.JSONDict, parentTaskId string,
guests []*SGuest, kwargs *jsonutils.JSONDict, parentTaskId string,
) error {
task, err := taskman.TaskManager.NewTask(ctx, "HostGuestsMigrateTask", host, userCred, kwargs, parentTaskId, "", nil)
taskItems := make([]db.IStandaloneModel, len(guests))
for i := range guests {
taskItems[i] = guests[i]
}
task, err := taskman.TaskManager.NewParallelTask(ctx, "HostGuestsMigrateTask", taskItems, userCred, kwargs, parentTaskId, "", nil)
if err != nil {
log.Errorln(err)
return err
+3 -1
View File
@@ -5316,6 +5316,7 @@ func (host *SHost) MigrateSharedStorageServers(ctx context.Context, userCred mcc
if err != nil {
return errors.Wrapf(err, "host %s(%s) get guests", host.Name, host.Id)
}
migGuests := []*SGuest{}
hostGuests := []*api.GuestBatchMigrateParams{}
for i := 0; i < len(guests); i++ {
@@ -5333,11 +5334,12 @@ func (host *SHost) MigrateSharedStorageServers(ctx context.Context, userCred mcc
}
guests[i].SetStatus(userCred, api.VM_START_MIGRATE, "host down")
hostGuests = append(hostGuests, bmp)
migGuests = append(migGuests, &guests[i])
}
}
kwargs := jsonutils.NewDict()
kwargs.Set("guests", jsonutils.Marshal(hostGuests))
return GuestManager.StartHostGuestsMigrateTask(ctx, userCred, host, kwargs, "")
return GuestManager.StartHostGuestsMigrateTask(ctx, userCred, migGuests, kwargs, "")
}
func (host *SHost) SetStatus(userCred mcclient.TokenCredential, status string, reason string) error {
+17 -1
View File
@@ -245,10 +245,26 @@ func (self *GuestMigrateTask) OnSrcPrepareComplete(ctx context.Context, guest *m
func (self *GuestMigrateTask) OnMigrateConfAndDiskCompleteFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) {
targetHostId, _ := self.Params.GetString("target_host_id")
guest.StartUndeployGuestTask(ctx, self.UserCred, "", targetHostId)
err := jsonutils.NewDict()
err.Set("MigrateConfAndDiskFailedReason", data)
self.SetStage("OnUndeployTargetGuestSucc", err)
guest.StartUndeployGuestTask(ctx, self.UserCred, self.GetTaskId(), targetHostId)
self.TaskFailed(ctx, guest, data)
}
func (self *GuestMigrateTask) OnUndeployTargetGuestSucc(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) {
err, _ := self.Params.Get("MigrateConfAndDiskFailedReason")
self.TaskFailed(ctx, guest, err)
}
func (self *GuestMigrateTask) OnUndeployTargetGuestSuccFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) {
prevErr, _ := self.Params.Get("MigrateConfAndDiskFailedReason")
err := jsonutils.NewDict()
err.Set("MigrateConfAndDiskFailedReason", prevErr)
err.Set("UndeployTargetGuestFailedReason", data)
self.TaskFailed(ctx, guest, err)
}
func (self *GuestMigrateTask) OnMigrateConfAndDiskComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) {
guestStatus, _ := self.Params.GetString("guest_status")
if !jsonutils.QueryBoolean(self.Params, "is_rescue_mode", false) && (guestStatus == api.VM_RUNNING || guestStatus == api.VM_SUSPEND) {
+11 -32
View File
@@ -34,7 +34,7 @@ type HostGuestsMigrateTask struct {
taskman.STask
}
func (self *HostGuestsMigrateTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
func (self *HostGuestsMigrateTask) OnInit(ctx context.Context, objs []db.IStandaloneModel, data jsonutils.JSONObject) {
guests := make([]*api.GuestBatchMigrateParams, 0)
err := self.Params.Unmarshal(&guests, "guests")
if err != nil {
@@ -43,51 +43,30 @@ func (self *HostGuestsMigrateTask) OnInit(ctx context.Context, obj db.IStandalon
}
preferHostId, _ := self.Params.GetString("prefer_host_id")
var guestMigrating bool
var migrateIndex int
for i := 0; i < len(guests); i++ {
guest := models.GuestManager.FetchGuestById(guests[i].Id)
self.SetStage("OnMigrateComplete", nil)
for i := range objs {
guest := objs[i].(*models.SGuest)
if guests[i].LiveMigrate {
err := guest.StartGuestLiveMigrateTask(
ctx, self.UserCred, guests[i].OldStatus, preferHostId, &guests[i].SkipCpuCheck, self.Id)
if err != nil {
log.Errorln(err)
continue
} else {
guestMigrating = true
migrateIndex = i
break
}
} else {
err := guest.StartMigrateTask(ctx, self.UserCred, guests[i].RescueMode,
false, guests[i].OldStatus, preferHostId, self.Id)
if err != nil {
log.Errorln(err)
continue
} else {
guestMigrating = true
migrateIndex = i
break
}
}
}
if !guestMigrating {
if jsonutils.QueryBoolean(self.Params, "some_guest_migrate_failed", false) {
self.SetStageFailed(ctx, jsonutils.NewString("some guest migrate failed"))
} else {
self.SetStageComplete(ctx, nil)
}
} else {
guests := append(guests[:migrateIndex], guests[migrateIndex+1:]...)
params := jsonutils.NewDict()
params.Set("guests", jsonutils.Marshal(guests))
self.SaveParams(params)
}
}
func (self *HostGuestsMigrateTask) OnInitFailed(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
kwargs := jsonutils.NewDict()
kwargs.Set("some_guest_migrate_failed", jsonutils.JSONTrue)
self.SaveParams(kwargs)
self.OnInit(ctx, obj, data)
func (self *HostGuestsMigrateTask) OnMigrateComplete(ctx context.Context, objs []db.IStandaloneModel, data jsonutils.JSONObject) {
self.SetStageComplete(ctx, nil)
}
func (self *HostGuestsMigrateTask) OnMigrateCompleteFailed(ctx context.Context, objs []db.IStandaloneModel, data jsonutils.JSONObject) {
self.SetStageFailed(ctx, data)
}
+25 -3
View File
@@ -37,14 +37,36 @@ func init() {
func (self *HostMaintainTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
host := obj.(*models.SHost)
guests, _ := self.Params.Get("guests")
preferHostId, _ := self.Params.Get("prefer_host_id")
var hostGuests = []*api.GuestBatchMigrateParams{}
err := self.Params.Unmarshal(&hostGuests, "guests")
if err != nil {
self.TaskFailed(ctx, host, jsonutils.NewString(err.Error()))
return
}
guests := make([]*models.SGuest, 0)
hostGuestParams := make([]*api.GuestBatchMigrateParams, 0)
for i := range hostGuests {
guest := models.GuestManager.FetchGuestById(hostGuests[i].Id)
if guest != nil {
guests = append(guests, guest)
hostGuestParams = append(hostGuestParams, hostGuests[i])
}
}
if len(guests) == 0 {
// no guest to migrate
self.SetStageComplete(ctx, nil)
return
}
kwargs := jsonutils.NewDict()
kwargs.Set("guests", guests)
kwargs.Set("guests", jsonutils.Marshal(hostGuestParams))
kwargs.Set("prefer_host_id", preferHostId)
self.SetStage("OnGuestsMigrate", nil)
err := models.GuestManager.StartHostGuestsMigrateTask(ctx, self.UserCred, host, self.Params, self.Id)
err = models.GuestManager.StartHostGuestsMigrateTask(ctx, self.UserCred, guests, kwargs, self.Id)
if err != nil {
self.TaskFailed(ctx, host, jsonutils.NewString(err.Error()))
return
+5 -1
View File
@@ -292,7 +292,11 @@ func (s *SKVMGuestInstance) asyncScriptStart(ctx context.Context, params interfa
if ctx != nil && len(appctx.AppContextTaskId(ctx)) >= 0 {
hostutils.TaskFailed(ctx, fmt.Sprintf("Async start server failed: %s", err))
}
s.SyncStatus("")
needMigrate := jsonutils.QueryBoolean(data, "need_migrate", false)
// do not syncstatus if need_migrate
if !needMigrate {
s.SyncStatus("")
}
return nil, err
}