Middle
Ce este workshopping-ul?
sobes.tech AI
Răspuns de la AI
Work stealing este un mecanism de echilibrare a încărcării în sistemele multi-thread, unde un thread inactiv "fură" sarcini din coada altui thread ocupat.
Caracteristici cheie:
- Decentralizat: Fiecare thread gestionează propria coadă locală de sarcini.
- Activ: Thread-ul inactiv caută activ muncă de la alții.
- Performanță ridicată: Eficient pentru sarcini cu încărcare neuniformă.
În Golang, este folosit în scheduler pentru distribuirea gorutinelor pe firele de sistem (M). Când un M își termină munca curentă (pachet de gorutine pe P), caută gorutine în alte P prin cozi locale.
// Pseudo-cod care demonstrează conceptul de work stealing
package main
import "fmt"
import "sync"
import "time"
type Task struct {
id int
}
// Simulare a unei cozi de sarcini
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
}
// Încercare de "furare" a unei sarcini din altă coadă
func (q *TaskQueue) StealTask() (Task, bool) {
q.mu.Lock()
defer q.mu.Unlock()
if len(q.tasks) < 2 { // Nu "furăm" dacă sunt puține sarcini
return Task{}, false
}
// Furează din mijloc sau de la sfârșit pentru a evita conflictele cu extragerea locală
index := len(q.tasks) / 2
task := q.tasks[index]
q.tasks = append(q.tasks[:index], q.tasks[index+1:]...)
return task, true
}
// Simulare a unui thread de lucru
func worker(id int, localQueue *TaskQueue, otherQueues []*TaskQueue, wg *sync.WaitGroup) {
defer wg.Done()
for {
// Încearcă să ia o sarcină din coada locală
task, ok := localQueue.GetLocalTask()
if ok {
fmt.Printf("Worker %d execută sarcina %d\n", id, task.id)
time.Sleep(100 * time.Millisecond) // Simulare de muncă
continue
}
// Dacă coada locală este goală, încearcă să fure
stolen := false
for _, queue := range otherQueues {
if queue == localQueue {
continue // Nu furăm de la noi înșine
}
task, ok := queue.StealTask()
if ok {
fmt.Printf("Worker %d a furat sarcina %d\n", id, task.id)
time.Sleep(100 * time.Millisecond) // Simulare de muncă
stolen = true
break // Am furat cu succes, ieșim din bucla de furt
}
}
if !stolen {
// Dacă nu am furat, poate nu mai sunt sarcini
// Într-un scheduler real, există mecanisme mai complexe de încheiere
// Pentru demonstrație, pur și simplu ieșim după o așteptare
fmt.Printf("Worker %d este inactiv, așteaptă...\n", id)
time.Sleep(50 * time.Millisecond)
// În scenariul real, aici ar fi un mecanism de încheiere sau de parcarea thread-ului
// return // Pentru demonstrație, vom permite ieșirea
break // Simplificare a simulării
}
}
}
func main() {
numWorkers := 4
numTasks := 20
taskQueues := make([]*TaskQueue, numWorkers)
for i := range taskQueues {
taskQueues[i] = &TaskQueue{}
}
// Distribuim sarcinile in mod neuniform pentru demonstrație de work stealing
for i := 0; i < numTasks; i++ {
queueIndex := i % 2 // Mai multe sarcini pentru primele două thread-uri
taskQueues[queueIndex].AddTask(Task{id: i})
}
var wg sync.WaitGroup
for i := 0; i < numWorkers; i++ {
wg.Add(1)
// Transmitem referințe către toate cozii pentru ca worker-ul să poată "fure"
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("Toate sarcinile au fost finalizate")
}