Middle
Mis on workshopping?
sobes.tech AI
Vastus AI-lt
Töö varastamine (work stealing) on mehhanism koormuse tasakaalustamiseks mitme-kiudses süsteemis, kus mitteaktiivne kiud (tööline) "varastab" ülesandeid teise, hõivatud kiu järjekorrast.
Peamised omadused:
- Dekentraliseeritud: Iga kiud haldab oma kohaliku ülesannete järjekorda.
- Aktiivne: Mitteaktiivne kiud otsib aktiivselt tööd teistelt.
- Kõrge jõudlus: Efektiivne ebaühtlase koormuse korral.
Golangis kasutatakse seda ajastajas (scheduler) goroutine'ide jaotamiseks operatsioonisüsteemi kiudude (M) vahel. Kui üks M lõpetab oma praeguse töö (goroutine'ide kogum P), siis ta otsib goroutine'e teiste P-de kohalike järjekordade kaudu.
// Pseudo-kood, mis demonstreerib töö varastamise kontseptsiooni
package main
import "fmt"
import "sync"
import "time"
type Task struct {
id int
}
// Tööde järjekord
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
}
// Katse "varastada" ülesanne teise järjekorrast
func (q *TaskQueue) StealTask() (Task, bool) {
q.mu.Lock()
defer q.mu.Unlock()
if len(q.tasks) < 2 { // Ei varasta, kui vähe ülesandeid
return Task{}, false
}
// Varastame keskmisest või lõpust, et vältida konflikte kohaliku eemaldamisega
index := len(q.tasks) / 2
task := q.tasks[index]
q.tasks = append(q.tasks[:index], q.tasks[index+1:]...)
return task, true
}
// Töötava kiu funktsioon
func worker(id int, localQueue *TaskQueue, otherQueues []*TaskQueue, wg *sync.WaitGroup) {
defer wg.Done()
for {
// Püüdke võtta ülesanne kohalikust järjekorrast
task, ok := localQueue.GetLocalTask()
if ok {
fmt.Printf("Tööline %d täidab ülesannet %d kohapeal\n", id, task.id)
time.Sleep(100 * time.Millisecond) // Imiteerimine
continue
}
// Kui kohalik järjekord on tühi, proovime varastada
stolen := false
for _, queue := range otherQueues {
if queue == localQueue {
continue // Ei varasta iseendast
}
task, ok := queue.StealTask()
if ok {
fmt.Printf("Tööline %d varastas ülesande %d\n", id, task.id)
time.Sleep(100 * time.Millisecond) // Imiteerimine
stolen = true
break // Edukas varastamine, väljumine
}
}
if !stolen {
// Kui ei õnnestunud varastada, võib olla ülesandeid enam pole
// Reaalses ajastajas on keerukamad mehhanismid
// Demonstratsiooniks lihtsalt väljumine pärast teatud ootamist
fmt.Printf("Tööline %d ootab...\n", id)
time.Sleep(50 * time.Millisecond)
// Reaalses stsenaariumis oleks siin lõpetamise või pargimise mehhanism
// return // Demonstratsiooniks lubame väljuda
break // Lihtsustamise eesmärgil
}
}
}
func main() {
numWorkers := 4
numTasks := 20
taskQueues := make([]*TaskQueue, numWorkers)
for i := range taskQueues {
taskQueues[i] = &TaskQueue{}
}
// Jaotame ülesanded ebaühtlaselt, et demonstreerida töö varastamist
for i := 0; i < numTasks; i++ {
queueIndex := i % 2 // Esimesed kaks töötajat saavad rohkem ülesandeid
taskQueues[queueIndex].AddTask(Task{id: i})
}
var wg sync.WaitGroup
for i := 0; i < numWorkers; i++ {
wg.Add(1)
// Anname viited kõigile järjekordadele, et töötaja saaks "varastada"
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("Kõik ülesanded on täidetud")
}