Go: أنماط التوازي في Go
آخر تحديث: 2026-08-26
لا تعتمد المعالجة المتزامنة في لغة Go على المكتبات، بل على الأنماط — Fan-in وFan-out وPipeline. ويمكن لهذه الأنماط مجتمعةً حل 90% من مشكلات المعالجة المتزامنة.
عندما تحتاج إلى معالجة كميات كبيرة من البيانات — من خلال قراءتها من مصدر، وتكليف عدة عمال بمعالجتها بشكل متوازٍ، ثم دمج النتائج في النهاية — كيف يمكنك تصميم حل أنيق؟ في هذا الدرس، ستتقن نمط التزامن الأكثر كلاسيكية في مجتمع Go.
1. ستتعلم
- نمط خط الأنابيب
- التوزيع / التجميع
- مجموعة العمال
- الوضع المبسط لقناة Or-Done
- انتشار الإلغاء
- مجموعة الأخطاء
golang.org/x/sync/errgroup - تجربة عملية: مسار معالجة البيانات في الوقت الفعلي
2. قصة حقيقية لمهندس بيانات
(1) المشكلة: استغرقت معالجة 10 ملايين إدخال سجل في مؤشر ترابط واحد طوال الليل
فاطمة تعمل مهندسة منصات بيانات. تنتج الشركة يوميًا 10 ملايين سجل وصول يتعين معالجتها في الوقت الفعلي:
"استخدمت النسخة الأولى حلقة for: قراءة سجل → تحليله → كتابته في قاعدة البيانات. في اليوم الأول من إطلاقها، اكتشفنا أن سرعة المعالجة لم تكن قادرة على مواكبة معدل الإنتاج — فقد كانت قائمة انتظار السجلات تنمو بمعدل مليوني إدخال في الساعة. قال مديري: «تقارير الغد تعتمد على هذه البيانات». أضفت المزيد من الخوادم، لكن الكود كان لا يزال أحادي الخيط — حيث كانت وحدة المعالجة المركزية (CPU) تستخدم 5% فقط من سعتها."
الكود الذي كانت تستخدمه في ذلك الوقت:
// Bad code: single-threaded, not utilizing CPU
func processLogs(logs []string) {
for _, log := range logs {
parsed := parse(log) // CPU-intensive
enriched := enrich(parsed) // CPU-intensive
save(enriched) // I/O-intensive
}
// 10 million log entries: 3 hours!
}
(2) نهج لغة Go في التعامل مع التزامن: التسلسل + التوزيع
// pipeline.go
package main
import (
"fmt"
"sync"
"time"
)
type LogEntry struct {
Raw string
Level string
Message string
Time time.Time
}
func main() {
logs := make([]string, 100)
for i := range logs {
logs[i] = fmt.Sprintf("log-%d", i)
}
start := time.Now()
// Pipeline: generator → parse → enrich → save
rawCh := generator(logs)
parsedCh := parseStage(rawCh, 4) // 4 parse workers
enrichedCh := enrichStage(parsedCh, 4) // 4 enrich workers
saveStage(enrichedCh, 2) // 2 save workers
fmt.Printf("Processing complete, elapsed: %v\n", time.Since(start))
}
// Stage 1: Generator (data source)
func generator(logs []string) <-chan string {
out := make(chan string, 100)
go func() {
defer close(out)
for _, log := range logs {
out <- log
}
}()
return out
}
// Stage 2: Parse (Fan-out)
func parseStage(in <-chan string, workers int) <-chan LogEntry {
out := make(chan LogEntry, 100)
for i := 0; i < workers; i++ {
go func() {
for raw := range in {
time.Sleep(1 * time.Millisecond) // Simulate parsing
out <- LogEntry{Raw: raw, Level: "INFO", Message: raw, Time: time.Now()}
}
}()
}
return out
}
// Stage 3: Enrich (Fan-out)
func enrichStage(in <-chan LogEntry, workers int) <-chan LogEntry {
out := make(chan LogEntry, 100)
for i := 0; i < workers; i++ {
go func() {
for entry := range in {
time.Sleep(2 * time.Millisecond) // Simulate data enrichment
out <- entry
}
}()
}
return out
}
// Stage 4: Save (Fan-in to 2 workers)
func saveStage(in <-chan LogEntry, workers int) {
var wg sync.WaitGroup
for i := 0; i < workers; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for entry := range in {
_ = entry // Simulate writing to database
time.Sleep(500 * time.Microsecond)
}
}()
}
wg.Wait()
}
(3) الأداء: أحادي الخيط مقابل التسلسل
| المقياس | أحادي الخيط | خط الإنتاج (4+4+2 عمال) | التحسن |
|---|---|---|---|
| وقت معالجة 10 ملايين سجل | ~3 ساعات | ~12 دقيقة | 15 ضعفًا |
| استخدام وحدة المعالجة المركزية | 5% | 85% | 17x |
| تعقيد الكود | بسيط | متوسط | — |
| قابلية التوسع | إضافة خادم | إضافة عامل (تغيير العدد) | — |
3. خط الأنابيب
▶ مثال: مسار معالجة الأرقام
package main
import (
"fmt"
)
// Stage 1: Generate
func generate(nums ...int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, n := range nums {
out <- n
}
}()
return out
}
// Stage 2: Square
func square(in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for n := range in {
out <- n * n
}
}()
return out
}
// Stage 3: Double
func double(in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for n := range in {
out <- n * 2
}
}()
return out
}
func main() {
// Pipeline: generate → square → double
for result := range double(square(generate(1, 2, 3, 4, 5))) {
fmt.Println(result)
}
// Output: 2, 8, 18, 32, 50 (n²×2)
}
graph LR
A[generate] -->|chan int| B[square]
B -->|chan int| C[double]
C -->|chan int| D[main]
<-chan T (قناة للقراءة فقط) وتقبل قناة <-chan T كمدخل. تعمل كل مرحلة في goroutine خاص بها، وترتبط القنوات ببعضها البعض — وهذا هو نموذج CSP في لغة Go.
4. التوزيع / التجميع
▶ مثال: التوزيع + التجميع
package main
import (
"fmt"
"sync"
)
// Stage 1: Generate
func generate(nums ...int) <-chan int {
out := make(chan int, 10)
go func() {
defer close(out)
for _, n := range nums {
out <- n
}
}()
return out
}
// Stage 2: Worker (Fan-out target)
func worker(id int, in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for n := range in {
result := n * n
fmt.Printf("Worker %d: %d^2 = %d\n", id, n, result)
out <- result
}
}()
return out
}
// Fan-in: merge multiple channels
func fanIn(channels ...<-chan int) <-chan int {
out := make(chan int)
var wg sync.WaitGroup
for _, ch := range channels {
wg.Add(1)
go func(c <-chan int) {
defer wg.Done()
for v := range c {
out <- v
}
}(ch)
}
go func() {
wg.Wait()
close(out)
}()
return out
}
func main() {
// Fan-out: one generator distributes to 3 workers
in := generate(1, 2, 3, 4, 5, 6)
w1 := worker(1, in)
w2 := worker(2, in)
w3 := worker(3, in)
// Fan-in: merge results from 3 workers into one channel
results := fanIn(w1, w2, w3)
for result := range results {
fmt.Printf("Result: %d\n", result)
}
}
graph LR
G[Generate] -->|fan-out| W1[Worker 1]
G -->|fan-out| W2[Worker 2]
G -->|fan-out| W3[Worker 3]
W1 -->|fan-in| R[Results]
W2 -->|fan-in| R
W3 -->|fan-in| R
5. مجموعة العمال
▶ مثال: مجموعة العمال
package main
import (
"fmt"
"sync"
"time"
)
type Job struct {
ID int
Payload string
}
type Result struct {
JobID int
Output string
Err error
}
func worker(id int, jobs <-chan Job, results chan<- Result, wg *sync.WaitGroup) {
defer wg.Done()
for job := range jobs {
fmt.Printf("Worker %d processing job %d\n", id, job.ID)
time.Sleep(100 * time.Millisecond) // Simulate work
results <- Result{
JobID: job.ID,
Output: fmt.Sprintf("processed by worker %d", id),
}
}
}
func main() {
const numJobs = 20
const numWorkers = 5
jobs := make(chan Job, numJobs)
results := make(chan Result, numJobs)
// Start Worker Pool
var wg sync.WaitGroup
for i := 0; i < numWorkers; i++ {
wg.Add(1)
go worker(i, jobs, results, &wg)
}
// Send jobs
for i := 0; i < numJobs; i++ {
jobs <- Job{ID: i, Payload: fmt.Sprintf("data-%d", i)}
}
close(jobs)
// Wait for all workers to complete
wg.Wait()
close(results)
// Collect results
for result := range results {
fmt.Printf("Job %d → %s\n", result.JobID, result.Output)
}
}
(2) خط الأنابيب مقابل مجموعة العمال مقابل التوزيع/التجميع
| النمط | المفهوم الأساسي | السيناريوهات القابلة للتطبيق |
|---|---|---|
| مسار المعالجة | تتدفق البيانات عبر مراحل متعددة، حيث تؤدي كل مرحلة خطوة واحدة | توجد خطوات معالجة محددة بوضوح (التحليل → التحويل → الحفظ) |
| مجموعة العمال | يقوم عدد ثابت من العمال باسترداد المهام من قائمة انتظار المهام | معالجة المهام دفعة واحدة (ضغط الصور، إرسال البريد الإلكتروني) |
| التوزيع/التجميع | توزيع المهام على عدة عمال للمعالجة المتوازية، ثم دمج النتائج | الحساب المتوازي غير المرتبط بالحالة (العمليات الحسابية، تصفية البيانات) |
6. قناة «Or-Done»
▶ مثال: نمط «أو-تم»
package main
import (
"fmt"
"time"
)
// orDone wraps a channel, making it support cancellation via a done channel
func orDone(done <-chan struct{}, c <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for {
select {
case <-done:
return
case v, ok := <-c:
if !ok {
return
}
select {
case out <- v:
case <-done:
return
}
}
}
}()
return out
}
func generateWithCancel(done <-chan struct{}, nums ...int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, n := range nums {
select {
case <-done:
return
case out <- n:
}
}
}()
return out
}
func main() {
done := make(chan struct{})
// Start generator (will be canceled after reaching 5)
nums := generateWithCancel(done, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10)
// Consume safely via orDone
for n := range orDone(done, nums) {
fmt.Printf("Processing: %d\n", n)
if n >= 5 {
close(done) // Cancel all operations
}
}
fmt.Println("Canceled")
time.Sleep(100 * time.Millisecond) // Wait for goroutines to exit
}
7. انتشار الأخطاء في مجموعة الأخطاء (errgroup)
⚙️ المتطلبات المسبقة: قم بتشغيل
go get golang.org/x/sync/errgroupقبل استخدام هذه الحزمة.
▶ مثال: errgroup
package main
import (
"fmt"
"time"
"golang.org/x/sync/errgroup"
)
func main() {
// Create an errgroup
// The first goroutine to return an error will cancel the others
g := errgroup.Group{}
// Start 3 goroutines
for i := 0; i < 3; i++ {
id := i
g.Go(func() error {
return doWork(id)
})
}
// Wait for all goroutines to complete, get the first error
if err := g.Wait(); err != nil {
fmt.Printf("Task failed: %v\n", err)
} else {
fmt.Println("All succeeded")
}
}
func doWork(id int) error {
fmt.Printf("Worker %d starting\n", id)
time.Sleep(time.Duration(id+1) * 500 * time.Millisecond)
if id == 1 {
return fmt.Errorf("worker %d failed", id)
}
fmt.Printf("Worker %d complete\n", id)
return nil
}
▶ مثال: مجموعة الأخطاء + انتهاء مهلة السياق
package main
import (
"context"
"fmt"
"time"
"golang.org/x/sync/errgroup"
)
func main() {
// errgroup with Context
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()
g, ctx := errgroup.WithContext(ctx)
// Workers check if Context is canceled
for i := 0; i < 5; i++ {
id := i
g.Go(func() error {
select {
case <-time.After(time.Duration(id+1) * time.Second):
fmt.Printf("Worker %d complete\n", id)
return nil
case <-ctx.Done():
fmt.Printf("Worker %d canceled: %v\n", id, ctx.Err())
return ctx.Err()
}
})
}
if err := g.Wait(); err != nil {
fmt.Printf("errgroup stopped: %v\n", err)
}
}
(3) sync.WaitGroup مقابل errgroup
| الخاصية | sync.WaitGroup | errgroup |
|---|---|---|
| مجموعة الأخطاء | غير مدعوم | ✅ إرجاع الخطأ الأول |
| إلغاء الانتشار | غير مدعوم | ✅ الإلغاء التلقائي باستخدام WithContext |
| إدارة الجوروتينات | الإضافة/الإنهاء يدويًّا | ✅ تلقائيًّا عبر طريقة Go |
| التحكم في مهلة الانتظار | التنفيذ اليدوي | ✅ سهل التنفيذ باستخدام السياق |
8. مثال كامل: مسار معالجة البيانات في الوقت الفعلي
// data_pipeline.go
package main
import (
"context"
"fmt"
"math/rand"
"sync"
"time"
)
// ---------- Data types ----------
type DataPoint struct {
ID int
Value float64
Timestamp time.Time
}
type ProcessedData struct {
DataPoint
Normalized bool
Anomaly bool
Category string
}
// ---------- Pipeline Stages ----------
// Stage 1: Source (data source)
func source(ctx context.Context, count int) <-chan DataPoint {
out := make(chan DataPoint, 100)
go func() {
defer close(out)
for i := 0; i < count; i++ {
select {
case <-ctx.Done():
return
case out <- DataPoint{
ID: i,
Value: rand.Float64() * 100,
Timestamp: time.Now(),
}:
}
}
}()
return out
}
// Stage 2: Normalize (Fan-out)
func normalize(ctx context.Context, in <-chan DataPoint, workers int) <-chan ProcessedData {
out := make(chan ProcessedData, 100)
var wg sync.WaitGroup
for i := 0; i < workers; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for dp := range orDoneDataPoint(ctx.Done(), in) {
select {
case <-ctx.Done():
return
case out <- ProcessedData{
DataPoint: dp,
Normalized: true,
Category: categorize(dp.Value),
}:
}
}
}()
}
go func() {
wg.Wait()
close(out)
}()
return out
}
// Stage 3: Detect Anomaly (Fan-out)
func detectAnomaly(ctx context.Context, in <-chan ProcessedData, workers int) <-chan ProcessedData {
out := make(chan ProcessedData, 100)
var wg sync.WaitGroup
for i := 0; i < workers; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for pd := range orDoneProcessedData(ctx.Done(), in) {
pd.Anomaly = pd.Value > 90 || pd.Value < 10
select {
case <-ctx.Done():
return
case out <- pd:
}
}
}()
}
go func() {
wg.Wait()
close(out)
}()
return out
}
// Stage 4: Sink (output, Fan-in)
func sink(ctx context.Context, in <-chan ProcessedData) {
var normalCount, anomalyCount int
for pd := range orDoneProcessedData(ctx.Done(), in) {
if pd.Anomaly {
anomalyCount++
fmt.Printf("[ANOMALY] ID=%d, Value=%.2f, Category=%s\n",
pd.ID, pd.Value, pd.Category)
} else {
normalCount++
}
}
fmt.Printf("\nSummary: normal=%d, anomaly=%d\n", normalCount, anomalyCount)
}
// ---------- Helper functions ----------
func orDoneDataPoint(done <-chan struct{}, in <-chan DataPoint) <-chan DataPoint {
out := make(chan DataPoint)
go func() {
defer close(out)
for {
select {
case <-done:
return
case v, ok := <-in:
if !ok {
return
}
select {
case out <- v:
case <-done:
return
}
}
}
}()
return out
}
func orDoneProcessedData(done <-chan struct{}, in <-chan ProcessedData) <-chan ProcessedData {
out := make(chan ProcessedData)
go func() {
defer close(out)
for {
select {
case <-done:
return
case v, ok := <-in:
if !ok {
return
}
select {
case out <- v:
case <-done:
return
}
}
}
}()
return out
}
func categorize(value float64) string {
switch {
case value < 30:
return "low"
case value < 70:
return "medium"
default:
return "high"
}
}
func main() {
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
fmt.Println("Starting real-time data processing pipeline...")
// Build Pipeline
data := source(ctx, 1000)
normalized := normalize(ctx, data, 5) // 5 normalize workers
detected := detectAnomaly(ctx, normalized, 3) // 3 anomaly detection workers
sink(ctx, detected)
fmt.Println("Pipeline processing complete")
}
expvar أو prometheus لعرض سرعة المعالجة وطول قائمة الانتظار ومعدل الأخطاء لكل مرحلة.
❓ أسئلة شائعة
select مع قناة done.errgroup وsync.WaitGroup؟errgroup عندما تحتاج إلى تجميع الأخطاء أو نشر عمليات الإلغاء. واستخدم sync.WaitGroup (الذي يتميز بحجم أصغر) عندما تحتاج فقط إلى انتظار انتهاء goroutine. errgroup هو مجموعة شاملة لـ WaitGroup من حيث الوظائف، لكنه يضيف goroutine إضافية لإدارة عمليات الإلغاء.done)؛ (2) استخدم errgroup أو WaitGroup لإدارة دورة الحياة؛ (3) استخدم Or-Done لتغليف القنوات غير الخاضعة للتحكم. لا تنشئ أبدًا goroutine إذا كنت لا تعرف متى ستنتهي.📖 ملخص
- مسار المعالجة: تتدفق البيانات عبر مراحل معالجة متعددة، مع وجود «غوروتين» واحد لكل مرحلة
- التوزيع (Fan-out): يتم توزيع قناة واحدة على عدة عمال (لأغراض المعالجة المتوازية)
- تجميع القنوات: دمج قنوات متعددة في قناة واحدة (للاستهلاك الموحد)
- مجموعة العمال: يقوم عدد ثابت من العمال باسترداد المهام من قائمة انتظار المهام
- Or-Done: استهلاك القنوات التي لا يمكن التحكم فيها بأمان مع دعم الإلغاء
- errgroup: يدير مجموعات goroutine، ويدعم تجميع الأخطاء ونشر الإلغاء
- وضع الجمع: يستخدم «Pipeline» داخليًّا «مجموعة العمال» (Worker Pool) + «التوزيع/التجميع» (Fan-out/Fan-in)
📝 تمارين
-
المشكلة الأساسية (صعوبة ⭐): قم بتنفيذ مسار معالجة رقمي من 3 مراحل: التوليد (1–100) → التصفية (الاحتفاظ بالأعداد الزوجية فقط) → الجمع (التراكم). تعمل كل مرحلة في goroutine منفصل، وترتبط المراحل ببعضها عبر قناة.
-
مشكلة متقدمة (درجة الصعوبة ⭐⭐): قم بتنفيذ نظام معالجة السجلات المتزامنة: 10 عمال تحليل (توزيع إلى الخارج) + 3 عمال كتابة (تجميع من الخارج) + معالجة الأخطاء باستخدام errgroup. قم بمحاكاة 1,000 إدخال سجل وقم بإخراج الإحصائيات (إجمالي الوقت المنقضي، عدد الإدخالات التي تمت معالجتها، عدد الأخطاء).
-
التحدي (الصعوبة: ⭐⭐⭐): صمم وقم بتنفيذ مجموعة من العمال مزودة بميزة موازنة الحمل والتوسع الديناميكي. المتطلبات: (1) ابدأ بـ 5 عمال؛ (2) قم بزيادة عدد العمال تلقائيًا (حتى 20) عندما يتجاوز عدد المهام المتراكمة في قائمة الانتظار عتبة معينة؛ (3) قم بتقليل عدد العمال (حتى 2) عندما تكون قائمة الانتظار فارغة؛ (4) استخدم
Contextللتحكم في إنهاء عمل العمال؛ (5) استخدم-raceللتحقق من عدم وجود حالة تنافس.