Go: Go Goroutines و WaitGroups
تُعدّ «الغوروتينات» حجر الزاوية في البرمجة المتزامنة بلغة «Go» — فهي ليست «خيوطًا» ولا «كوروتينات»، بل وحدات تزامن خفيفة الوزن يديرها بيئة تشغيل «Go». وبفضل مساحة المكدس التي تبلغ 2 كيلوبايت، فإنها تتيح تحقيق التزامن بمعدل يصل إلى الملايين.
لا يضاهي نموذج التزامن في لغة Go أي نموذج آخر: فإطلاق «غوروتين» لا يتطلب سوى الكلمة الرئيسية go، في حين أن بيئة تشغيل Go تقوم بكفاءة بتخصيص آلاف وآلاف من «الغوروتينات» لعدد قليل من خيوط نظام التشغيل خلف الكواليس.
1. ستتعلم
- مفهوم «الغوروتينات» وكلمة
go sync.WaitGroupينتظر انتهاء جميع الغوروتينات- الجوروتين مقابل مؤشر ترابط نظام التشغيل (حجم المكدس/التكلفة الإضافية للإنشاء)
GOMAXPROCSونموذج جدولة GMP- أسباب تسربات الجوروتينات وكيفية معالجتها
runtime.NumGoroutineالمراقبة- استخدام «غوروتينات» لمعالجة 100,000 إدخال سجل في وقت واحد
2. قصة حقيقية لمهندس بيانات
(1) المشكلة: تستغرق معالجة 100,000 إدخال سجل 5 ساعات في مؤشر ترابط واحد
تشارلي هو مهندس بيانات يقوم بمعالجة 100,000 سجل خادم يوميًا:
«كل يوم عند الفجر، تُنفَّذ مهمة تنظيف السجلات، حيث تتم معالجة 100,000 سطر واحدًا تلو الآخر — وتستغرق العملية أكثر من خمس ساعات. وخلال الاجتماع الصباحي، يسأل مدير المشروع دائمًا: «متى ستكون بيانات الأمس جاهزة؟» فأجيب: «في الساعة 3:00 مساءً»، فيردون: «لماذا لا يمكن أن تكون جاهزة بحلول الساعة 9:00 صباحًا؟»»
فتح ملف الكود الحالي:
// Serial processing: 100k entries × 200ms/entry = 5.5 hours
func processLogs(logs []LogEntry) []Result {
results := make([]Result, 0, len(logs))
for _, log := range logs {
result := processSingleLog(log) // 200ms each, including network IO
results = append(results, result)
}
return results
}
تتطلب معالجة كل إدخال في السجل استدعاءً واحدًا لواجهة برمجة التطبيقات الخارجية (مع وقت انتظار يبلغ حوالي 200 مللي ثانية)، لكن الحساب الفعلي الذي تقوم به وحدة المعالجة المركزية يستغرق أقل من 1 مللي ثانية — حيث يُقضى 99.5% من الوقت في انتظار الشبكة.
(2) حل Go: التوازي في Goroutine
// log_processor.go
package main
import (
"fmt"
"runtime"
"sync"
"time"
)
type LogEntry struct {
ID int
Message string
Level string
}
type Result struct {
ID int
Status string
}
func processSingleLog(log LogEntry) Result {
time.Sleep(200 * time.Millisecond)
return Result{ID: log.ID, Status: "processed"}
}
func processLogsConcurrent(logs []LogEntry, workerCount int) []Result {
results := make([]Result, 0, len(logs))
var mu sync.Mutex
var wg sync.WaitGroup
// Use channels for task distribution
jobs := make(chan LogEntry, len(logs))
for _, log := range logs {
jobs <- log
}
close(jobs)
// Start workerCount workers
for w := 0; w < workerCount; w++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
for log := range jobs {
result := processSingleLog(log)
mu.Lock()
results = append(results, result)
mu.Unlock()
}
}(w)
}
wg.Wait()
return results
}
func main() {
// Simulate 100 log entries (demo; can scale to 100k)
logs := make([]LogEntry, 100)
for i := 0; i < 100; i++ {
logs[i] = LogEntry{ID: i, Message: fmt.Sprintf("log-%d", i), Level: "INFO"}
}
start := time.Now()
results := processLogsConcurrent(logs, 10)
elapsed := time.Since(start)
fmt.Printf("Processed %d logs in %v\n", len(results), elapsed)
fmt.Printf("Current goroutine count: %d\n", runtime.NumGoroutine())
}
الناتج:
Processed 100 logs in 2.05s
Current goroutine count: 1
(3) الأداء: التسلسلي مقابل المتزامن
| طريقة المعالجة | 100 سجل | 100,000 سجل | عدد الغوروتينات |
|---|---|---|---|
| مسلسل | 20 حلقة | 5.5 ساعة | 1 |
| 10 عمال متزامنين | 2 ثانية | 33 دقيقة | ~12 |
| 100 عامل متزامن | 0.2 ثانية | 3.3 دقيقة | ~102 |
worker count = GOMAXPROCS * 2–10؛ أما بالنسبة للمهام التي تتطلب استخدامًا مكثفًا لوحدة المعالجة المركزية، فنوصي باستخدام worker count = GOMAXPROCS.
3. أساسيات الجوروتين
(1) الكلمة الرئيسية go
package main
import (
"fmt"
"time"
)
func printNumbers() {
for i := 1; i <= 5; i++ {
time.Sleep(100 * time.Millisecond)
fmt.Printf("%d ", i)
}
}
func printLetters() {
for _, c := range "ABCDE" {
time.Sleep(150 * time.Millisecond)
fmt.Printf("%c ", c)
}
}
func main() {
go printNumbers() // Execute concurrently in a goroutine
go printLetters()
time.Sleep(1 * time.Second)
fmt.Println()
}
النتيجة (قد تختلف في كل مرة):
1 A 2 B 3 C 4 D 5 E
(2) غوروتين مجهول
package main
import (
"fmt"
"time"
)
func main() {
go func() {
fmt.Println("Anonymous goroutine running")
}()
go func(msg string) {
fmt.Println("Anonymous goroutine with parameter:", msg)
}("hello")
time.Sleep(100 * time.Millisecond)
fmt.Println("main done")
}
time.Sleep أعلاه هو انتظار انتهاء الغوروتين — وفي كود الإنتاج، يجب عليك استخدام sync.WaitGroup.
4. sync.WaitGroup: الانتظار حتى تتم المزامنة
(1) ثلاث طرق لاستخدام WaitGroup
package main
import (
"fmt"
"sync"
)
func worker(id int, wg *sync.WaitGroup) {
defer wg.Done() // 3. Tell WaitGroup this worker is done
fmt.Printf("Worker %d started\n", id)
// Simulate work...
fmt.Printf("Worker %d done\n", id)
}
func main() {
var wg sync.WaitGroup
for i := 1; i <= 3; i++ {
wg.Add(1) // 1. Increment counter
go worker(i, &wg)
}
wg.Wait() // 2. Block until all workers finish
fmt.Println("All workers done!")
}
| الطريقة | الدالة |
|---|---|
wg.Add(delta int) |
يزيد قيمة العداد (يُستدعى عادةً قبل بدء goroutine) |
wg.Done() |
يقلل قيمة العداد (عادةً ما يكون استدعاءً لـ defer داخل جورووتين) |
wg.Wait() |
يوقف التشغيل حتى يصل العداد إلى الصفر |
▶ مثال: WaitGroup + مخاطر الإغلاق
package main
import (
"fmt"
"sync"
)
func main() {
var wg sync.WaitGroup
// ❌ Wrong: Closure captures loop variable
for i := 1; i <= 3; i++ {
wg.Add(1)
go func() {
defer wg.Done()
fmt.Printf("Wrong i=%d\n", i) // All print 3 or 4
}()
}
wg.Wait()
// ✅ Correct: Pass parameter (creates a copy)
for i := 1; i <= 3; i++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
fmt.Printf("Correct i=%d\n", id)
}(i)
}
wg.Wait()
}
i مباشرةً — وبحلول الوقت الذي تبدأ فيه الجوروتين، قد تكون الحلقة قد انتهت بالفعل، وقد تكون قيمة i قد تم استبدالها. يجب عليك تمرير نسخة من المعلمة.
5. الجوروتينات مقابل خيوط نظام التشغيل
(1) الاختلافات الرئيسية
package main
import (
"fmt"
"runtime"
)
func main() {
// View current GOMAXPROCS
fmt.Printf("GOMAXPROCS: %d\n", runtime.GOMAXPROCS(0))
fmt.Printf("CPU cores: %d\n", runtime.NumCPU())
fmt.Printf("Current goroutine count: %d\n", runtime.NumGoroutine())
}
| البعد | مؤشر ترابط نظام التشغيل | غوروتين |
|---|---|---|
| الحجم الأولي للمكدس | 1–8 ميغابايت | 2 كيلوبايت |
| الحد الأقصى للمكدس | ثابت عند 1–8 ميغابايت | يتوسع ديناميكيًا إلى 1 غيغابايت |
| الوقت الإضافي اللازم للإنشاء | ~1 µs | ~0.1 µs |
| تبديل السياق | وضع النواة (~1 ميكروثانية) | وضع المستخدم (~0.1 ميكروثانية) |
| العدد الأقصى | ~10,000 | ملايين |
| المجدول | جدولة نواة نظام التشغيل | GMP في بيئة تشغيل Go |
(2) نموذج جدولة GMP
graph TB
G[G: Goroutine Queue] --> P[P: Logical Processor<br/>GOMAXPROCS count]
P --> M[M: OS Thread<br/>Scheduled by kernel]
M --> CPU[CPU Core]
subgraph Global Queue
GQ[(Global Goroutine Queue)]
end
GQ --> P
style P fill:#e1f5fe
style M fill:#fff3e0
GMP: G (غوروتين) — P (معالج، المعالج المنطقي؛ العدد = GOMAXPROCS) — M (الجهاز، مؤشر ترابط نظام التشغيل). يقوم بيئة تشغيل Go بجدولة الغوروتينات (G) في قائمة الانتظار المحلية للمعالج (P)، ثم يرتبط المعالج (P) بالجهاز (M) للتنفيذ.
6. GOMAXPROCS والتحكم في التزامن
(1) إعدادات GOMAXPROCS
package main
import (
"fmt"
"runtime"
)
func main() {
fmt.Printf("Default GOMAXPROCS: %d\n", runtime.GOMAXPROCS(0))
// CPU-intensive: GOMAXPROCS = NumCPU
// IO-intensive: GOMAXPROCS = NumCPU * 2~10
runtime.GOMAXPROCS(4) // Set to 4
fmt.Printf("After setting GOMAXPROCS: %d\n", runtime.GOMAXPROCS(0))
}
▶ مثال: تأثير GOMAXPROCS على الأداء
package main
import (
"fmt"
"runtime"
"sync"
"time"
)
func cpuIntensiveTask() {
sum := 0
for i := 0; i < 10000000; i++ {
sum += i
}
}
func benchmarkGOMAXPROCS(n int) time.Duration {
runtime.GOMAXPROCS(n)
var wg sync.WaitGroup
start := time.Now()
for i := 0; i < 8; i++ {
wg.Add(1)
go func() {
defer wg.Done()
cpuIntensiveTask()
}()
}
wg.Wait()
return time.Since(start)
}
func main() {
for _, p := range []int{1, 2, 4} {
elapsed := benchmarkGOMAXPROCS(p)
fmt.Printf("GOMAXPROCS=%d: %v\n", p, elapsed)
}
}
7. تسربات الجوروتينات واستكشاف الأخطاء وإصلاحها
(1) سيناريوهات التسرب
package main
import (
"fmt"
"runtime"
"time"
)
// Leaky goroutine: reads from a channel that is never closed
func leakyGoroutine() {
ch := make(chan int)
go func() {
<-ch // Blocks forever—no one writes to ch
}()
}
func main() {
for i := 0; i < 10; i++ {
leakyGoroutine()
}
time.Sleep(100 * time.Millisecond)
fmt.Printf("Goroutine count after leak: %d\n", runtime.NumGoroutine())
// Output: Goroutine count after leak: 11 (10 leaked + 1 main)
}
▶ مثال: 4 طرق لبدء «غوروتين»
package main
import (
"fmt"
"sync"
"time"
)
func say(msg string) {
fmt.Println(msg)
}
func main() {
var wg sync.WaitGroup
// Style 1: Named function
wg.Add(1)
go func() {
defer wg.Done()
say("Style 1: Named function")
}()
// Style 2: Anonymous function
wg.Add(1)
go func() {
defer wg.Done()
fmt.Println("Style 2: Anonymous function")
}()
// Style 3: Anonymous function with parameter (recommended—avoids closure trap)
wg.Add(1)
go func(msg string) {
defer wg.Done()
fmt.Println(msg)
}("Style 3: With parameter")
// Style 4: Function as value
wg.Add(1)
fn := func() {
defer wg.Done()
fmt.Println("Style 4: Function variable")
}
go fn()
wg.Wait()
}
▶ مثال: مقارنة بين مكدسات «غوروتين» ومكدسات خيوط نظام التشغيل
package main
import (
"fmt"
"runtime"
"sync"
)
func main() {
var memStats runtime.MemStats
var wg sync.WaitGroup
num := 100000 // 100k goroutines
runtime.GC()
runtime.ReadMemStats(&memStats)
before := memStats.HeapAlloc
for i := 0; i < num; i++ {
wg.Add(1)
go func(n int) {
defer wg.Done()
_ = n
}(i)
}
wg.Wait()
runtime.ReadMemStats(&memStats)
after := memStats.HeapAlloc
perGoroutine := float64(after-before) / float64(num)
fmt.Printf("Started %d goroutines\n", num)
fmt.Printf("Memory increase: %.2f MB\n", float64(after-before)/1024/1024)
fmt.Printf("Per goroutine: ~%.2f KB\n", perGoroutine/1024)
fmt.Printf("Total goroutine count: %d\n", runtime.NumGoroutine())
}
الناتج:
Started 100000 goroutines
Memory increase: 21.45 MB
Per goroutine: ~0.22 KB
Total goroutine count: 1
▶ مثال: مراقبة دورة حياة الجوروتين
package main
import (
"fmt"
"runtime"
"time"
)
func safeWorker(done chan struct{}) {
<-done // Wait for exit signal
}
func main() {
done := make(chan struct{})
// Start 5 workers
for i := 0; i < 5; i++ {
go safeWorker(done)
}
fmt.Printf("Goroutine count after start: %d\n", runtime.NumGoroutine())
// Send exit signal
close(done)
time.Sleep(10 * time.Millisecond)
fmt.Printf("Goroutine count after close: %d\n", runtime.NumGoroutine())
}
(5) قائمة مراجعة لمنع التسرب
| السيناريو | طرق الوقاية |
|---|---|
| القراءة من قناة دون أن يقوم أحد بالكتابة إليها | استخدم قناة ذات مخزن مؤقت أو select مع الخيار default |
| الكتابة إلى قناة دون أن يقرأها أحد | تأكد من وجود مستهلك أو استخدم select مع القيمة الافتراضية |
| الحلقات اللانهائية في Goroutine | التحكم في الخروج عبر قناة done أو السياق |
يتسبب time.After في حدوث تسرب في الذاكرة داخل حلقة |
قم بإنشاء مؤقت جديد في كل تكرار باستخدام time.NewTimer + Stop |
select {} فارغ |
تأكد من وجود مسار للخروج |
8. مثال كامل: المعالجة المتزامنة لـ 100,000 إدخال في السجل
// log_processor.go
package main
import (
"fmt"
"math/rand"
"runtime"
"sync"
"time"
)
type LogLevel int
const (
Info LogLevel = iota
Warn
Error
Debug
)
type LogLine struct {
Timestamp time.Time
Level LogLevel
Message string
Source string
}
type ProcessedLog struct {
Original LogLine
Severity string
Alert bool
ProcessedAt time.Time
}
func (l LogLine) Process() ProcessedLog {
// Simulate processing time (IO wait)
time.Sleep(time.Duration(50+rand.Intn(50)) * time.Millisecond)
severity := "low"
alert := false
switch l.Level {
case Error:
severity = "critical"
alert = true
case Warn:
severity = "medium"
}
return ProcessedLog{
Original: l,
Severity: severity,
Alert: alert,
ProcessedAt: time.Now(),
}
}
// Serial processing
func processSerial(logs []LogLine) []ProcessedLog {
results := make([]ProcessedLog, 0, len(logs))
for _, log := range logs {
results = append(results, log.Process())
}
return results
}
// Concurrent processing: worker pool pattern
func processConcurrent(logs []LogLine, workers int) []ProcessedLog {
jobs := make(chan LogLine, len(logs))
results := make(chan ProcessedLog, len(logs))
// Fill jobs
for _, log := range logs {
jobs <- log
}
close(jobs)
var wg sync.WaitGroup
for w := 0; w < workers; w++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
for log := range jobs {
results <- log.Process()
}
}(w)
}
// Wait for all workers to finish, then close results
go func() {
wg.Wait()
close(results)
}()
// Collect results
processed := make([]ProcessedLog, 0, len(logs))
for r := range results {
processed = append(processed, r)
}
return processed
}
func main() {
const logCount = 10000 // 10k for demo; can scale to 100k
levels := []LogLevel{Info, Warn, Error, Debug}
sources := []string{"api-gateway", "user-service", "payment", "database"}
logs := make([]LogLine, logCount)
for i := 0; i < logCount; i++ {
logs[i] = LogLine{
Timestamp: time.Now(),
Level: levels[rand.Intn(len(levels))],
Message: fmt.Sprintf("log-entry-%d", i),
Source: sources[rand.Intn(len(sources))],
}
}
fmt.Printf("Log count: %d\n", logCount)
fmt.Printf("CPU cores: %d, GOMAXPROCS: %d\n", runtime.NumCPU(), runtime.GOMAXPROCS(0))
// Serial processing
start := time.Now()
serialResults := processSerial(logs)
serialTime := time.Since(start)
fmt.Printf("\nSerial: %v (%d entries/sec)\n", serialTime,
int(float64(logCount)/serialTime.Seconds()))
// Concurrent processing (10 workers)
runtime.GC()
start = time.Now()
concurrentResults := processConcurrent(logs, 10)
concurrentTime := time.Since(start)
fmt.Printf("Concurrent(10): %v (%d entries/sec)\n", concurrentTime,
int(float64(logCount)/concurrentTime.Seconds()))
// Concurrent processing (50 workers)
runtime.GC()
start = time.Now()
concurrentResults = processConcurrent(logs, 50)
concurrentTime = time.Since(start)
fmt.Printf("Concurrent(50): %v (%d entries/sec)\n", concurrentTime,
int(float64(logCount)/concurrentTime.Seconds()))
fmt.Printf("\nResult validation: serial=%d, concurrent=%d\n",
len(serialResults), len(concurrentResults))
fmt.Printf("Current goroutine count: %d\n", runtime.NumGoroutine())
// Count alerts
alertCount := 0
for _, r := range concurrentResults {
if r.Alert {
alertCount++
}
}
fmt.Printf("Alert count: %d\n", alertCount)
}
النتيجة المتوقعة:
Log count: 10000
CPU cores: 8, GOMAXPROCS: 8
Serial: 5.2s (1923 entries/sec)
Concurrent(10): 520ms (19230 entries/sec)
Concurrent(50): 110ms (90909 entries/sec)
Result validation: serial=10000, concurrent=10000
Current goroutine count: 1
Alert count: 2453
recover الخارجية. يجب حماية نقطة الدخول لكل «goroutine» باستخدام defer recover(). سيتناول الدرس 14 المتعلق بالقنوات هذا النمط بمزيد من التفصيل.
❓ أسئلة شائعة
sync.WaitGroup—تقوم wg.Add(n) بزيادة العدد، وwg.Done() بتقليله، بينما تقوم wg.Wait() بحجب التنفيذ حتى يصل العدد إلى الصفر. لاحظ أنه يجب استدعاء Add في النطاق الخارجي أو قبل تشغيل goroutine، ويجب تغليف Done بـ defer لضمان تنفيذه.runtime.NumGoroutine() للتحقق مما إذا كان العدد في تزايد مستمر؛ واستخدم net/http/pprof لفحص مسارات مكدس الجوروتينات. الأسباب الشائعة للتسربات: عمليات قراءة أو كتابة محجوبة في القنوات، وعبارات select التي لا تحتوي على جملة default، وحلقات for التي تطلق goroutines دون التحكم في خروجها.runtime.Gosched؟time.Sleep.defer func() { if r := recover(); r != nil { ... } }().go خالية من المعلمات؟go myFunc(arg1, arg2). احذر من مشاكل المتغيرات الدورية عند التقاط المتغيرات في الإغلاقات — استخدم تمرير المعلمات (الذي ينشئ نسخًا) بدلاً من التقاطها مباشرةً.📖 ملخص
- تعمل الكلمة الرئيسية
goعلى تشغيل «جوروتين»، مما يتيح التزامن الخفيف في وضع المستخدم - تبدأ مكدسات الغوروتينات من 2 كيلوبايت وتتوسع ديناميكيًا؛ أما خيوط نظام التشغيل فهي ثابتة عند 1–8 ميغابايت
sync.WaitGroup: الطرق الثلاث —AddوDoneوWait— تتولى إدارة دورة حياة الجوروتين- عندما تلتقط دالة الإغلاق الخاصة بـ goroutine متغيرًا في حلقة، يجب تمرير المعلمة عن طريق النسخ
GOMAXPROCSيتحكم في درجة التوازي؛ القيمة الافتراضية = NumCPU- نموذج جدولة GMP: غوروتين → معالج → جهاز
- تسربات الجوروتينات: القنوات المحجوبة، والحلقات اللانهائية التي لا تحتوي على شروط خروج
runtime.NumGoroutine()يراقب العدد الحالي لـ goroutines
📝 تمارين
-
المسألة الأساسية (الصعوبة ⭐): اكتب برنامجًا يبدأ 5 goroutines، حيث تطبع كل واحدة منها رقمها الخاص وكلمة "Hello". استخدم
sync.WaitGroupللانتظار حتى تنتهي جميعها، وتأكد من أن ترتيب تنفيذ الـ goroutines عشوائي. -
مشكلة متقدمة (درجة الصعوبة ⭐⭐): قم بتنفيذ مرشح للأعداد الأولية يعمل بشكل متزامن: عند إعطاء عدد n، استخدم 4 goroutines للتحقق بشكل متزامن من الأعداد بين 1 و n التي تعتبر أعدادًا أولية، ثم أعد قائمة بالأعداد الأولية. يجب عليك استخدام WaitGroup للتزامن ومقارنة الفرق في الأداء بين النهج المتزامن والنهج التسلسلي.
-
التحدي (الصعوبة ⭐⭐⭐): قم ببناء مجدول مهام متزامن: عند وجود مجموعة من المهام (
[]func() Result)، قم بتنفيذها بشكل متزامن باستخدام نموذج تجمع العمال. يجب أن يدعم المجدول ما يلي: (1) عدد قابل للتكوين من العمال؛ (2) تحمل الأعطال (لا يؤثر حدوث حالة ذعر في goroutine واحد على البقية)؛ (3) استدعاءات التقدم (طباعة رسالة كلما تم إكمال 10% من المهام)؛ (4) الكشف عن تسرباتruntime.NumGoroutine.