Go: أنماط «Go Select» والتزامن
選択هو بمثابة «سكين الجيش السويسري» للبرمجة المتزامنة في لغة Go — فهو يستقبل البيانات من قنوات متعددة في آن واحد ويعالج أي قناة تصبح جاهزة أولاً. وهو يعمل بمثابة «محول ذكي» يربط بين goroutines.
إذا تصورنا القناة على أنها خط هاتفي يربط بين «الغوروتينات»، فإن 選択 هي بمثابة عاملة المقسم الهاتفي — فهي تستمع إلى جميع المكالمات في آن واحد وترد على أي مكالمة ترن. في هذا الدرس، ستتقن جميع الأنماط الأساسية لـ 選択.
1. ستتعلم
- قواعد بناء الجملة وسلوك تقنية التعدد
選択 - الاختيار العشوائي (توافر عدة حالات في الوقت نفسه)
time.Afterالتحكم في مهلة الانتظارdefault: عملية غير معطلةfor-選択الخروج من الحلقة (القناة انتهت)- توزيع المهام باستخدام وضع التوزيع المتفرع
- نتائج التجميع في وضع «fan-in»
- إنشاء مسار لمعالجة البيانات باستخدام نموذج المسار
2. القصة الحقيقية لمهندس تداول كمي
(1) المشكلة: يؤدي استعلام ثلاثة مصادر بيانات إلى تشغيل وحدة المعالجة المركزية (CPU) بكامل طاقتها
بوب هو مهندس «الخلفية» في فريق التداول الكمي. ويحتاج إلى مراقبة ثلاثة مصادر لبيانات الأسهم في آن واحد:
"تتطلب استراتيجيتنا استرداد أسعار الأسهم في الوقت الفعلي من ثلاث بورصات في آن واحد: بورصة نيويورك (NYSE)، وناسداك (NASDAQ)، وبورصة لندن (LSE). وكان التطبيق السابق المكتوب بلغة جافا يستخدم خيطًا واحدًا لاستقصاء ثلاثة مآخذ WebSocket، مع انتظار نشط مدته 10 مللي ثانية في كل مرة، مما كان يستهلك 30% من طاقة وحدة المعالجة المركزية (CPU) — وقد قال مديري إنني كنت أستهلك طاقة أكثر من نظام التداول نفسه."
فتح ملف الكود الحالي:
// Bad code: busy polling
func pollDataSources() {
for {
// Poll every 10ms, wasting CPU
if data1 := pollNYSE(); data1 != nil {
process(data1)
}
if data2 := pollNASDAQ(); data2 != nil {
process(data2)
}
if data3 := pollLSE(); data3 != nil {
process(data3)
}
time.Sleep(10 * time.Millisecond)
}
}
ثلاث مشكلات: (1) يؤدي فاصل الاستقصاء البالغ 10 مللي ثانية إلى إهدار موارد وحدة المعالجة المركزية؛ (2) عدم تزامن وصول البيانات ومعالجتها؛ (3) عدم إمكانية التنبؤ بتأخيرات الاستقصاء التي تتراوح بين 0 و10 مللي ثانية.
(2) حل Go: استخدم select للاستماع إلى ثلاث قنوات
// market_data.go
package main
import (
"fmt"
"math/rand"
"time"
)
type Quote struct {
Source string
Symbol string
Price float64
}
func simulateExchange(name string, out chan<- Quote) {
symbols := []string{"AAPL", "GOOGL", "MSFT", "AMZN"}
for {
time.Sleep(time.Duration(50+rand.Intn(200)) * time.Millisecond)
quote := Quote{
Source: name,
Symbol: symbols[rand.Intn(len(symbols))],
Price: 100 + rand.Float64()*200,
}
out <- quote
}
}
func main() {
nyse := make(chan Quote)
nasdaq := make(chan Quote)
lse := make(chan Quote)
go simulateExchange("NYSE", nyse)
go simulateExchange("NASDAQ", nasdaq)
go simulateExchange("LSE", lse)
// الاختيار listens to all three channels simultaneously
timeout := time.After(2 * time.Second)
for {
الاختيار {
case q := <-nyse:
fmt.Printf("[NYSE] %s: $%.2f\n", q.Symbol, q.Price)
case q := <-nasdaq:
fmt.Printf("[NASDAQ] %s: $%.2f\n", q.Symbol, q.Price)
case q := <-lse:
fmt.Printf("[LSE] %s: $%.2f\n", q.Symbol, q.Price)
case <-timeout:
fmt.Println("Demo ended")
return
}
}
}
الناتج:
[NASDAQ] AMZN: $198.32
[NYSE] AAPL: $150.45
[LSE] MSFT: $287.10
...
(3) الأداء: SELECT Listen مقابل الاستقصاء
| الطريقة | استخدام وحدة المعالجة المركزية | زمن الاستجابة | تعقيد الكود |
|---|---|---|---|
| معدل الاستقصاء 10 مللي ثانية | 30% | 0–10 مللي ثانية | منخفض |
| معدل الاستقصاء 100 مللي ثانية | 3% | 0–100 مللي ثانية | منخفض |
| الاختيار: الاستماع | 0% | 0 مللي ثانية (في الوقت الفعلي) | معتدل |
選択 + channel تعمل بناءً على الأحداث — حيث تدخل الجوروتين في حالة سكون (دون استهلاك موارد وحدة المعالجة المركزية) في حالة عدم وجود بيانات، ويتم إيقاظها بواسطة بيئة التشغيل عند وصول البيانات. وهذا ما يجعل نموذج التزامن في لغة Go فعالاً للغاية.
3. أساسيات SELECT
(1) بناء الجملة الانتقائية
select {
case v := <-ch1:
// ch1 is ready
case v := <-ch2:
// ch2 is ready
case ch3 <- value:
// Can send to ch3 (ch3 has space or has a receiver)
default:
// All channels are not ready
}
▶ مثال: الاختيار (الاختيار العشوائي)
package main
import (
"fmt"
)
func main() {
ch1 := make(chan string, 1)
ch2 := make(chan string, 1)
ch1 <- "from ch1"
ch2 <- "from ch2"
// Both channels are ready; الاختيار picks one at random
for i := 0; i < 2; i++ {
الاختيار {
case msg := <-ch1:
fmt.Println(msg)
case msg := <-ch2:
fmt.Println(msg)
}
}
}
النتيجة (قد تختلف في كل مرة):
from ch1
from ch2
4. التحكم في مهلة الانتظار
(1) وضع انتهاء المهلة
package main
import (
"fmt"
"time"
)
func slowOperation() string {
time.Sleep(2 * time.Second)
return "result"
}
func main() {
ch := make(chan string)
go func() {
ch <- slowOperation()
}()
select {
case result := <-ch:
fmt.Println("Success:", result)
case <-time.After(1 * time.Second):
fmt.Println("Timeout! Operation exceeded 1 second")
}
}
الناتج:
Timeout! Operation exceeded 1 second
▶ مثال: فترات الانتظار المتدرجة (الانتظار المتدرج)
package main
import (
"fmt"
"time"
)
func fetchFromCache() string {
time.Sleep(50 * time.Millisecond)
return "cache-data"
}
func fetchFromDB() string {
time.Sleep(200 * time.Millisecond)
return "db-data"
}
func fetchFromAPI() string {
time.Sleep(500 * time.Millisecond)
return "api-data"
}
func main() {
cache := make(chan string)
db := make(chan string)
api := make(chan string)
go func() { cache <- fetchFromCache() }()
go func() { db <- fetchFromDB() }()
go func() { api <- fetchFromAPI() }()
الاختيار {
case r := <-cache:
fmt.Println("Cache hit:", r)
case <-time.After(100 * time.Millisecond):
الاختيار {
case r := <-db:
fmt.Println("DB returned:", r)
case <-time.After(300 * time.Millisecond):
الاختيار {
case r := <-api:
fmt.Println("API returned:", r)
case <-time.After(600 * time.Millisecond):
fmt.Println("All data sources timed out!")
}
}
}
}
time.After(d) القيمة <-chan time.Time وترسل الوقت الحالي بعد مرور d ثانية. في كل مرة تستدعي فيها الدالة select الدالة time.After، فإنها تنشئ مؤقتًا جديدًا — إذا كنت تستدعيها بشكل متكرر في حلقة، فتأكد من استخدام الدالة time.NewTimer وإيقافها لتجنب تسرب الموارد.
5. الإعداد الافتراضي: عملية غير معطلة
(1) الإرسال/الاستلام غير المعطل
package main
import (
"fmt"
)
func main() {
ch := make(chan int, 1)
// Non-blocking receive
select {
case v := <-ch:
fmt.Println("Received:", v)
default:
fmt.Println("No data (non-blocking)")
}
// Non-blocking send
ch <- 1
select {
case ch <- 2:
fmt.Println("Send successful")
default:
fmt.Println("Buffer full (non-blocking)")
}
}
الناتج:
No data (non-blocking)
Buffer full (non-blocking)
▶ مثال: قناة غير معطلة + قاطع دائرة الاستقصاء
package main
import (
"fmt"
"time"
)
func main() {
ch := make(chan int, 3)
go func() {
for i := 1; i <= 10; i++ {
select {
case ch <- i:
// Send successful
default:
fmt.Printf("Buffer full, dropped %d\n", i)
}
time.Sleep(10 * time.Millisecond)
}
close(ch)
}()
// Slow consumer
for v := range ch {
fmt.Printf("Processing: %d\n", v)
time.Sleep(50 * time.Millisecond)
}
}
الناتج:
Processing: 1
Processing: 2
Processing: 3
Buffer full, dropped 4
Buffer full, dropped 5
Processing: 6
...
default إلى العودة فورًا من select — ويتم تنفيذ default إذا لم تكن أي من القنوات جاهزة. وهذا مثالي لـعمليات القنوات غير المعطلة، ولكن كن حذرًا: استخدام default في حلقة for قد يتسبب في حدوث حلقة مشغولة.
6. حلقات for-select وقناة done
(1) وضع الخروج «القناة المنتهية»
package main
import (
"fmt"
"time"
)
func worker(done <-chan struct{}) {
for {
select {
case <-done:
fmt.Println("worker exiting")
return
default:
fmt.Println("worker working...")
time.Sleep(200 * time.Millisecond)
}
}
}
func main() {
done := make(chan struct{})
go worker(done)
time.Sleep(1 * time.Second)
close(done)
time.Sleep(100 * time.Millisecond)
fmt.Println("main exiting")
}
▶ مثال: ثلاث طرق للخروج من حلقة for-اختيار
package main
import (
"fmt"
"time"
)
func main() {
ch := make(chan int)
done := make(chan struct{})
// Producer
go func() {
for i := 1; i <= 5; i++ {
ch <- i
time.Sleep(100 * time.Millisecond)
}
close(ch)
}()
// Consumer (for-select loop)
go func() {
for {
select {
case v, ok := <-ch:
if !ok {
fmt.Println("Style 1: channel closed exit")
close(done)
return
}
fmt.Printf("Processing: %d\n", v)
case <-time.After(1 * time.Second):
fmt.Println("Style 2: timeout exit")
close(done)
return
}
}
}()
<-done
fmt.Println("main exiting")
}
(3) مقارنة بين أوضاع الخروج من «for-select»
| الوضع | شروط التشغيل | المزايا | العيوب |
|---|---|---|---|
| إغلاق (ch) | يقوم المرسل بإغلاق القناة | إنهاء طبيعي | لا يمكن استخدامه إلا من قِبل المستقبل |
| قناة "done" | إشارة close(done) |
مرنة؛ يمكن تشغيلها من مصدر خارجي | تتطلب قناة إضافية |
| المهلة | time.After المهلة |
تمنع حدوث حالة التعطل | المهلات الثابتة ليست مرنة بما يكفي |
| السياق | ctx.Done() |
يمكن إرسال إشارة إلغاء | الدرس 17: تحليل متعمق |
7. وضع التوزيع / التجميع
(1) التوزيع: توزيع المهام
package main
import (
"fmt"
"sync"
)
func fanOut(jobs <-chan int, workers int) []<-chan int {
channels := make([]<-chan int, workers)
for w := 0; w < workers; w++ {
ch := make(chan int, 10)
channels[w] = ch
go func(id int, out chan<- int) {
defer close(out)
for job := range jobs {
result := job * job
fmt.Printf("Worker %d: %d^2 = %d\n", id, job, result)
out <- result
}
}(w, ch)
}
return channels
}
▶ مثال: fan-in: تجميع النتائج
package main
import (
"fmt"
"sync"
)
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() {
jobs := make(chan int, 10)
for i := 1; i <= 6; i++ {
jobs <- i
}
close(jobs)
// fan-out: distribute to 3 workers
workers := fanOut(jobs, 3)
// fan-in: aggregate all worker results
results := fanIn(workers...)
// Collect results
sum := 0
count := 0
for r := range results {
sum += r
count++
}
fmt.Printf("Total %d results, sum = %d\n", count, sum)
}
الناتج:
Worker 2: 1^2 = 1
Worker 0: 3^2 = 9
Worker 1: 2^2 = 4
Worker 1: 5^2 = 25
Worker 0: 4^2 = 16
Worker 2: 6^2 = 36
Total 6 results, sum = 91
(3) التوزيع (fan-out) مقابل التجميع (fan-in)
| الوضع | الاتجاه | الغرض |
|---|---|---|
| التوزيع | قناة واحدة → N قنوات | توزيع المهام، الحساب المتوازي |
| تجميع المدخلات | N قنوات → قناة واحدة | تجميع النتائج، وجمع السجلات |
flowchart LR
subgraph FanOut [Fan-Out]
J[Jobs Channel] --> W1[Worker 1]
J --> W2[Worker 2]
J --> W3[Worker 3]
end
subgraph FanIn [Fan-In]
W1 --> R[Results Channel]
W2 --> R
W3 --> R
end
style J fill:#e1f5fe
style R fill:#fff3e0
8. مثال كامل: تجميع أسعار الأسهم من ثلاثة مصادر للبيانات
// market_aggregator.go
package main
import (
"fmt"
"math/rand"
"sync"
"time"
)
// ---------- Data Types ----------
type Trade struct {
Source string
Symbol string
Price float64
Volume int
Time time.Time
}
type AggregatedQuote struct {
Symbol string
AvgPrice float64
TotalVolume int
Sources int
High float64
Low float64
Time time.Time
}
// ---------- Simulated Exchange Data Source ----------
var symbols = []string{"AAPL", "GOOGL", "MSFT", "AMZN", "TSLA", "META"}
func simulateExchange(name string, out chan<- Trade, done <-chan struct{}) {
for {
select {
case <-done:
return
default:
time.Sleep(time.Duration(100+rand.Intn(300)) * time.Millisecond)
trade := Trade{
Source: name,
Symbol: symbols[rand.Intn(len(symbols))],
Price: 100 + rand.Float64()*200,
Volume: rand.Intn(1000) + 100,
Time: time.Now(),
}
select {
case out <- trade:
case <-done:
return
}
}
}
}
// ---------- Aggregator ----------
type Aggregator struct {
trades chan Trade
done chan struct{}
wg sync.WaitGroup
}
func NewAggregator() *Aggregator {
return &Aggregator{
trades: make(chan Trade, 100),
done: make(chan struct{}),
}
}
func (a *Aggregator) AddSource(name string) {
a.wg.Add(1)
go func() {
defer a.wg.Done()
simulateExchange(name, a.trades, a.done)
}()
fmt.Printf("Added data source: %s\n", name)
}
func (a *Aggregator) Start(interval time.Duration, callback func(map[string]*AggregatedQuote)) {
// Aggregation window
window := make(map[string][]Trade)
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case trade := <-a.trades:
window[trade.Symbol] = append(window[trade.Symbol], trade)
case <-ticker.C:
// Time window reached; compute aggregated quotes
quotes := make(map[string]*AggregatedQuote)
for symbol, trades := range window {
if len(trades) == 0 {
continue
}
var sumPrice float64
var totalVol int
high := trades[0].Price
low := trades[0].Price
for _, t := range trades {
sumPrice += t.Price * float64(t.Volume)
totalVol += t.Volume
if t.Price > high {
high = t.Price
}
if t.Price < low {
low = t.Price
}
}
sources := make(map[string]bool)
for _, t := range trades {
sources[t.Source] = true
}
quotes[symbol] = &AggregatedQuote{
Symbol: symbol,
AvgPrice: sumPrice / float64(totalVol),
TotalVolume: totalVol,
Sources: len(sources),
High: high,
Low: low,
Time: time.Now(),
}
}
callback(quotes)
window = make(map[string][]Trade) // Reset window
case <-a.done:
return
}
}
}
func (a *Aggregator) Stop() {
close(a.done)
a.wg.Wait()
}
// ---------- Main Function ----------
func main() {
aggregator := NewAggregator()
// Add 3 data sources
aggregator.AddSource("NYSE")
aggregator.AddSource("NASDAQ")
aggregator.AddSource("LSE")
fmt.Println("\nStarting aggregation (output every 2 seconds)...\n")
// Start aggregation with 2-second time window
done := make(chan struct{})
go func() {
aggregator.Start(2*time.Second, func(quotes map[string]*AggregatedQuote) {
fmt.Printf("=== Aggregation Report %s ===\n", time.Now().Format("15:04:05"))
for _, q := range quotes {
fmt.Printf("%-6s | $%.2f (avg) | vol=%d | src=%d | H=%.2f L=%.2f\n",
q.Symbol, q.AvgPrice, q.TotalVolume, q.Sources, q.High, q.Low)
}
fmt.Println()
})
close(done)
}()
// Run for 10 seconds then stop
time.Sleep(10 * time.Second)
aggregator.Stop()
<-done
fmt.Println("Quote aggregator stopped")
}
النتيجة المتوقعة:
Added data source: NYSE
Added data source: NASDAQ
Added data source: LSE
Starting aggregation (output every 2 seconds)...
=== Aggregation Report 10:00:02 ===
AAPL | $152.34 (avg) | vol=2341 | src=3 | H=165.20 L=142.10
GOOGL | $178.90 (avg) | vol=1567 | src=2 | H=185.00 L=172.30
MSFT | $295.40 (avg) | vol=890 | src=3 | H=301.20 L=288.50
=== Aggregation Report 10:00:04 ===
...
Quote aggregator stopped
select، يجب الاشتراك في كل من case <-done: وcase trade := <-a.trades:. إذا اشتركت في قناة trades فقط، فلن يتم توصيل إشارة Stop. يجب أن تتضمن كل حلقة for-select شرط خروج.
❓ أسئلة شائعة
select أو جملة default مقترنة بفحوصات يدوية.time.After مع حالات انتهاء المهلة؟time.After(d) قناة تتلقى قيمة بعد مرور d وحدة زمنية. في عبارة select، استخدم case <- time.After(d): كفرع انتهاء المهلة. ملاحظة: كل استدعاء لـ select لـ time.After ينشئ مؤقتًا جديدًا؛ وفي الحلقة، يُنصح بإعادة استخدامه مع time.NewTimer.default في عبارة select؟default تجعل عبارة select غير معطلة — فهي تنفذ كتلة default فورًا عندما لا تكون أي من القنوات جاهزة. مناسب لـ: محاولة الإرسال/الاستقبال دون حجب، وحساب معدلات الفشل، وقطع الدائرة/تدهور الأداء. ملاحظة: يؤدي استخدام default في حلقة for إلى إنشاء حلقة مشغولة؛ وعادةً ما تحتاج إلى إضافة time.Sleep أو الحد من التكرار.done := make(chan struct{})، واستمع إلى case <-done: داخل كتلة select، واستدعِ close(done) عندما تريد الخروج. وستتلقى جميع goroutines التي تستمع إلى done إشارة الخروج في وقت واحد. يُفضل استخدام close بدلاً من إرسال قيمة — حيث يمكن استلامها عدة مرات.select للإرسال؟case ch <- value: هو أيضًا select case صالح. ويصبح جاهزًا عندما تتوفر مساحة في القناة (مع التخزين المؤقت) أو عندما يكون هناك مستقبل (بدون تخزين مؤقت). تُستخدم هذه الميزة في أوضاع تحديد معدل الإرسال وأوضاع الإشارات.select {} فارغ؟select {} معطلة إلى أجل غير مسمى — نظرًا لعدم وجود حالات ولا قيمة افتراضية. ويمكن أن يُعتبر هذا إشارة إلى «الانتظار إلى أجل غير مسمى». ومع ذلك، في كود الإنتاج، يجب عليك التأكد من وجود مسار للخروج؛ وإلا، فسوف تتسرب الجوروتينات.break الموجودة داخل حلقة for-select أن تخرج من الحلقة؟break الموجودة داخل كتلة select لا تخرج إلا من كتلة select؛ وهي لا تخرج من حلقة for. للخروج من حلقة for-select، استخدم تسمية متبوعة بـ break label أو return أو القناة done.📖 ملخص
- تقوم ميزة «select» بمراقبة عدة قنوات في آن واحد وتقوم بمعالجة القناة التي تصبح جاهزة أولاً
- عندما تكون هناك عدة حالات جاهزة في الوقت نفسه، اختر واحدة منها عشوائيًا
- يُستخدم
time.Afterللتحكم في مهلة الانتظار، بينما يُعاد استخدامtime.NewTimerفي الحلقات defaultيجعل عملية الاختيار غير معطلة، وهو ما يجعلها مناسبة لعمليات قطع الدائرة أو إعادة المحاولةdone channelهي الطريقة القياسية للخروج من حلقة for-select- تعمل عملية «التوزيع» (fan-out) على توزيع المهام، بينما تعمل عملية «التجميع» (fan-in) على تجميع النتائج، ويتحد الاثنان معًا لتشكيل مسار متوازي
- يجب أن يتضمن كل تعبير «for-select» شرطًا للخروج
📝 تمارين
-
المشكلة الأساسية (الصعوبة ⭐): اكتب برنامجًا يبدأ ثلاث «جوروتينات»، ترسل كل منها عددًا صحيحًا إلى قناتها الخاصة. استخدم
selectلاستلام هذه الأعداد وطباعتها في آن واحد. أضف فرعًا لانتهاء المهلة للخروج من البرنامج بعد 3 ثوانٍ. -
مشكلة متقدمة (درجة الصعوبة ⭐⭐): قم بتنفيذ ذاكرة تخزين مؤقتة ذات مهلة:
func FetchWithCache(key string, cache map[string]string) string. أولاً، تحقق من قناة ذاكرة التخزين المؤقتة (مهلة 10 مللي ثانية)؛ وإذا لم يتم العثور على نتيجة مطابقة، فقم بالاستعلام عن «قاعدة البيانات المحاكاة» (500 مللي ثانية) وقم بتعبئة ذاكرة التخزين المؤقتة بالنتيجة. يجب عليك استخدامselectمعtime.Afterلتنفيذ فترات انتظار متدرجة. -
التحدي (الصعوبة: ⭐⭐⭐): تنفيذ نظام تجميع السجلات: ثلاثة مصادر للسجلات (goroutines) تولد كل منها سجلات بمستويات مختلفة (INFO/WARN/ERROR)، والتي يتم تجميعها في قناة واحدة باستخدام عملية fan-in، ثم تتم معالجتها حسب المستوى باستخدام
select—تؤدي سجلات ERROR إلى إصدار تنبيه فوري (طباعة)، ويتم تخزين سجلات WARN مؤقتًا (يتم تفريغها على دفعات من 5)، ويتم كتابة سجلات INFO على دفعات (يتم تفريغها على دفعات مكونة من 10). المتطلبات: إنهاء العمل بشكل سلس عبر التفريغ عند انتهاء المهلة وقناةdone.