Middle
Mi az a workshopping?
sobes.tech MI
Válasz az MI-től
A work stealing egy terheléselosztási mechanizmus több szálú rendszerekben, ahol egy inaktív szál "ellop" feladatokat egy másik elfoglalt szál sorából.
Fő jellemzők:
- Decentralizált: Minden szál kezeli a saját helyi feladatsorát.
- Aktív: Az inaktív szál aktívan keres munkát másoktól.
- Magas teljesítmény: Hatékony egyenetlen terhelés esetén.
A Golangban a schedulerben használják, hogy goroutine-okat oszthassanak szálak között (M). Amikor egy M befejezi a jelenlegi munkát (goroutine csomagokat P-n), más P-k helyi sorait vizsgálja.
// Pszeudokód a work stealing koncepciójának bemutatására
package main
import "fmt"
import "sync"
import "time"
type Task struct {
id int
}
// Feladatsor szimulációja
type TaskQueue struct {
tasks []Task
mu sync.Mutex
}
func (q *TaskQueue) AddTask(task Task) {
q.mu.Lock()
defer q.mu.Unlock()
q.tasks = append(q.tasks, task)
}
func (q *TaskQueue) GetLocalTask() (Task, bool) {
q.mu.Lock()
defer q.mu.Unlock()
if len(q.tasks) == 0 {
return Task{}, false
}
task := q.tasks[0]
q.tasks = q.tasks[1:]
return task, true
}
// Próbálkozás "lopásra" másik sorból
func (q *TaskQueue) StealTask() (Task, bool) {
q.mu.Lock()
defer q.mu.Unlock()
if len(q.tasks) < 2 { // Nem "lopunk" kevés feladattal
return Task{}, false
}
// Lopás középről vagy végéről, hogy elkerüljük a konfliktusokat a helyi kivétellel
index := len(q.tasks) / 2
task := q.tasks[index]
q.tasks = append(q.tasks[:index], q.tasks[index+1:]...)
return task, true
}
// Munkafolyamat szimulációja
func worker(id int, localQueue *TaskQueue, otherQueues []*TaskQueue, wg *sync.WaitGroup) {
defer wg.Done()
for {
// Próbálkozik feladatot szerezni a helyi sorból
task, ok := localQueue.GetLocalTask()
if ok {
fmt.Printf("Worker %d végrehajtja a %d feladatot\n", id, task.id)
time.Sleep(100 * time.Millisecond) // Munkaszimuláció
continue
}
// Ha a helyi sor üres, próbálkozik lopással
stolen := false
for _, queue := range otherQueues {
if queue == localQueue {
continue // Nem lopunk saját sorból
}
task, ok := queue.StealTask()
if ok {
fmt.Printf("Worker %d, lopott feladat %d\n", id, task.id)
time.Sleep(100 * time.Millisecond) // Munkaszimuláció
stolen = true
break // Sikeres lopás, kilép a lopási ciklusból
}
}
if !stolen {
// Ha nem sikerült lopni, lehet, hogy nincs több feladat
// Egy valódi ütemezőben bonyolultabb mechanizmusok vannak a befejezéshez
// Demonstrációként, egyszerűen kilépünk egy várakozás után
fmt.Printf("Worker %d inaktív, vár...\n", id)
time.Sleep(50 * time.Millisecond)
// Egy valódi szcenárióban itt lenne a befejező vagy a szál parkoltatásának mechanizmusa
// return // Demonstrációként engedjük a kilépést
break // Egyszerűsítés a szimulációban
}
}
}
func main() {
numWorkers := 4
numTasks := 20
taskQueues := make([]*TaskQueue, numWorkers)
for i := range taskQueues {
taskQueues[i] = &TaskQueue{}
}
// Feladatok elosztása nem egyenletes a work stealing demonstrációhoz
for i := 0; i < numTasks; i++ {
queueIndex := i % 2 // Több feladat az első két munkásnak
taskQueues[queueIndex].AddTask(Task{id: i})
}
var wg sync.WaitGroup
for i := 0; i < numWorkers; i++ {
wg.Add(1)
// Átadjuk az összes sor hivatkozását, hogy a munkás "lopni" tudjon
otherQueues := make([]*TaskQueue, 0, numWorkers-1)
for j := range taskQueues {
if i != j {
otherQueues = append(otherQueues, taskQueues[j])
}
}
go worker(i, taskQueues[i], otherQueues, &wg)
}
wg.Wait()
fmt.Println("Minden feladat végrehajtva")
}