Middle
Czym jest workshopping?
sobes.tech AI
Odpowiedź od AI
Work stealing (kradzież pracy) to mechanizm równoważenia obciążenia w systemach wielowątkowych, w którym nieaktywny wątek (worker) "kradnie" zadania z kolejki innego, zajętego wątku.
Kluczowe cechy:
- Decentralizowany: Każdy wątek zarządza własną kolejką zadań.
- Aktywny: Nieaktywny wątek aktywnie szuka pracy u innych.
- Wysoka wydajność: Efektywny dla zadań z nierównomiernym obciążeniem.
W Golang, jest używany w schedulerze do rozdzielania goroutine pomiędzy wątki systemu operacyjnego (M). Gdy M kończy swoją bieżącą pracę (pakiet goroutine na P), szuka goroutine u innych P przez ich lokalne kolejki.
// Pseudokod demonstrujący koncepcję work stealing
package main
import "fmt"
import "sync"
import "time"
type Task struct {
id int
}
// Symulacja kolejki zadań
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óba "kradzieży" zadania z innej kolejki
func (q *TaskQueue) StealTask() (Task, bool) {
q.mu.Lock()
defer q.mu.Unlock()
if len(q.tasks) < 2 { // Nie "kradniemy", jeśli mało zadań
return Task{}, false
}
// Kradniemy z środka lub końca, aby uniknąć konfliktów z lokalnym pobraniem
index := len(q.tasks) / 2
task := q.tasks[index]
q.tasks = append(q.tasks[:index], q.tasks[index+1:]...)
return task, true
}
// Symulacja worker'a
func worker(id int, localQueue *TaskQueue, otherQueues []*TaskQueue, wg *sync.WaitGroup) {
defer wg.Done()
for {
// Próba pobrania zadania z lokalnej kolejki
task, ok := localQueue.GetLocalTask()
if ok {
fmt.Printf("Worker %d wykonuje zadanie %d lokalnie\n", id, task.id)
time.Sleep(100 * time.Millisecond) // Symulacja pracy
continue
}
// Jeśli lokalna kolejka jest pusta, próbujemy ukraść
stolen := false
for _, queue := range otherQueues {
if queue == localQueue {
continue // Nie kradniemy od siebie
}
task, ok := queue.StealTask()
if ok {
fmt.Printf("Worker %d ukradł zadanie %d\n", id, task.id)
time.Sleep(100 * time.Millisecond) // Symulacja pracy
stolen = true
break // Ukończyliśmy po udanym kradzieży
}
}
if !stolen {
// Jeśli nie udało się ukraść, może nie ma więcej zadań
// W rzeczywistym schedulerze są bardziej złożone mechanizmy kończenia
// Dla demo, kończymy po pewnym czasie oczekiwania
fmt.Printf("Worker %d jest bezczynny, czeka...\n", id)
time.Sleep(50 * time.Millisecond)
// W scenariuszu rzeczywistym, tutaj byłby mechanizm kończenia lub parkingu
// return // Dla demo pozwalamy mu wyjść
break // Dla uproszczenia symulacji
}
}
}
func main() {
numWorkers := 4
numTasks := 20
taskQueues := make([]*TaskQueue, numWorkers)
for i := range taskQueues {
taskQueues[i] = &TaskQueue{}
}
// Rozkład zadań nierównomierny, aby pokazać work stealing
for i := 0; i < numTasks; i++ {
queueIndex := i % 2 // Więcej zadań dla pierwszych dwóch workerów
taskQueues[queueIndex].AddTask(Task{id: i})
}
var wg sync.WaitGroup
for i := 0; i < numWorkers; i++ {
wg.Add(1)
// Przekazujemy referencje do wszystkich kolejek, aby worker mógł "kradnąć"
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("Wszystkie zadania wykonane")
}
```}}}}}]}},{