Go ทำให้การเริ่มต้น concurrency มีต้นทุนต่ำมาก แค่ใส่ go ไว้หน้าการเรียกฟังก์ชัน คุณก็ได้ goroutine ที่มี stack เพียงไม่กี่กิโลไบต์ โดยมี runtime scheduler คอยจัดการให้ แต่ต้นทุนที่ต่ำนี่เองคือจุดเริ่มต้นของปัญหา เพราะ goroutine สร้างง่าย มันจึงรั่ว (leak) ได้ง่าย ถูกสร้างเกินจำเป็นได้ง่าย และถูกลืมตอน shutdown ได้ง่ายเช่นกัน
บทความนี้พูดถึง pattern ที่ใช้งานได้จริงใน service บน production แนวคิดหลักที่อยู่เบื้องหลังทุก pattern คือ ทุก goroutine ที่คุณสร้างต้องมีเจ้าของที่ชัดเจน และมีวิธีจบการทำงานที่ชัดเจน
สามคำถามที่ต้องตอบก่อนเขียน go
ก่อนจะสั่งรัน goroutine ให้ตอบคำถามสามข้อนี้ก่อน:
- ใครเป็นคนรอมัน? ต้องมีบางอย่างที่รู้ว่ามันทำงานเสร็จเมื่อไร ไม่ว่าจะเป็น
sync.WaitGroup,errgroup.Groupหรือ channel สำหรับรับผลลัพธ์ - มันจะหยุดก่อนกำหนดได้อย่างไร? ถ้า request ถูกยกเลิกหรือ process กำลัง shutdown ตัว goroutine ต้องมีสัญญาณที่มันคอยเช็กอยู่จริง ๆ
- มีได้พร้อมกันสูงสุดกี่ตัว? แนวคิดแบบ "หนึ่งตัวต่อหนึ่ง item ที่เข้ามา" ไม่มีขีดจำกัดบน และอะไรก็ตามที่ไม่มีขีดจำกัด สุดท้ายจะก่อปัญหาเสมอ
ถ้าตอบไม่ได้ครบทั้งสามข้อ goroutine ตัวนั้นก็น่าจะกลายเป็น leak ในอนาคต
การยกเลิกงานควรอยู่ใน context.Context
ให้ส่ง context.Context เป็น argument ตัวแรกให้กับทุกอย่างที่ block หรือทำ I/O นี่คือวิธีมาตรฐานในการส่งต่อ deadline และการยกเลิก (cancellation) ข้ามขอบเขตของ API
func fetchPrice(ctx context.Context, client *http.Client, sku string) (Price, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, priceURL(sku), nil)
if err != nil {
return Price{}, err
}
resp, err := client.Do(req)
if err != nil {
return Price{}, fmt.Errorf("fetch price %s: %w", sku, err)
}
defer resp.Body.Close()
var p Price
return p, json.NewDecoder(resp.Body).Decode(&p)
}เมื่อ goroutine ทำงานเป็น loop มันควรเช็ก ctx.Done() ใน select เดียวกับที่ใช้รับงาน:
for {
select {
case <-ctx.Done():
return ctx.Err()
case job, ok := <-jobs:
if !ok {
return nil
}
process(job)
}
}loop ที่ block รอรับค่าจาก channel โดยไม่ได้เฝ้าดู ctx.Done() ไปด้วย จะยังคงทำงานต่อไปแม้ว่าจะไม่มีใครต้องการผลลัพธ์ของมันแล้ว
จำกัดความขนานด้วย worker pool
สมมติว่าคุณต้อง enrich ข้อมูล 50,000 record ด้วยการเรียก API ภายนอก ถ้าสร้าง goroutine 50,000 ตัว connection pool จะหมด โดน rate limit และ latency จะคาดเดาไม่ได้ การใช้ worker จำนวนคงที่ที่อ่านงานจาก channel ร่วมกันจะช่วยให้โหลดคงที่:
func enrichAll(ctx context.Context, records []Record, workers int) error {
g, ctx := errgroup.WithContext(ctx)
jobs := make(chan Record)
// Producer: feeds work and stops early if a worker fails.
g.Go(func() error {
defer close(jobs)
for _, r := range records {
select {
case jobs <- r:
case <-ctx.Done():
return ctx.Err()
}
}
return nil
})
for i := 0; i < workers; i++ {
g.Go(func() error {
for r := range jobs {
if err := enrich(ctx, r); err != nil {
return err // cancels ctx for everyone else
}
}
return nil
})
}
return g.Wait()
}errgroup จาก golang.org/x/sync จัดการให้สามอย่างในโค้ดนี้ คือรอ goroutine ทุกตัว คืนค่า error ตัวแรก และยกเลิก context ที่ใช้ร่วมกันเพื่อให้ worker ตัวอื่นหยุดได้เร็ว ถ้าคุณต้องการแค่จำกัดจำนวน concurrency โดยไม่ต้องมี producer แยก การใช้ g.SetLimit(n) จะง่ายกว่านี้อีก
เลือกจำนวน worker อย่างไร
- งานที่ใช้ CPU เป็นหลัก (CPU-bound): เริ่มจากค่าใกล้ ๆ
runtime.GOMAXPROCS(0)การมี worker มากกว่าจำนวน core มีแต่จะเพิ่ม overhead ในการ schedule - งานที่รอ I/O เป็นหลัก (I/O-bound): ข้อจำกัดอยู่ที่ระบบปลายทาง ไม่ใช่ CPU ของคุณ ให้กำหนดขนาด pool ตามที่ database หรือ API รับไหว ซึ่งมักหมายถึงให้เท่ากับ
MaxConnsPerHostของ HTTP client หรือขนาดของ DB pool
ต้องวัดผลจริงเสมอ ตัวเลขที่เหมาะสมคือค่าที่ทำให้ p99 latency คงที่ในขณะที่ throughput เพิ่มขึ้น
Channel: กฎความเป็นเจ้าของที่ช่วยป้องกัน panic
bug ส่วนใหญ่ที่เกี่ยวกับ channel มาจากความไม่ชัดเจนว่าใครเป็นเจ้าของ กฎสองข้อนี้ครอบคลุมเกือบทุกกรณี:
- เฉพาะฝั่งผู้ส่งเท่านั้นที่ปิด channel การปิดจากฝั่งผู้รับจะทำให้เกิด panic เมื่อผู้ส่งเขียนค่าเข้าไปอีกครั้ง
- ถ้ามีผู้ส่งหลายตัว ห้ามให้ตัวใดตัวหนึ่งปิด channel เอง ให้ใช้
WaitGroupร่วมกับ goroutine ประสานงานหนึ่งตัวที่คอยปิด channel หลังจากผู้ส่งทุกตัวทำงานเสร็จแล้ว
var wg sync.WaitGroup
out := make(chan Result)
for _, src := range sources {
wg.Add(1)
go func(s Source) {
defer wg.Done()
for r := range s.Stream(ctx) {
select {
case out <- r:
case <-ctx.Done():
return
}
}
}(src)
}
go func() { wg.Wait(); close(out) }()buffered channel มีไว้รองรับ burst สั้น ๆ ไม่ได้มีไว้แก้ปัญหา consumer ที่ช้า ถ้าคุณพบว่าตัวเองต้องเพิ่มขนาด buffer เพื่อไม่ให้ระบบค้าง แปลว่า consumer คือคอขวด และ buffer ก็แค่เลื่อนปัญหาออกไป
ตรวจจับ goroutine leak
goroutine ที่รั่วมักจะค้างอยู่ตลอดไปกับการส่งหรือรับค่าที่ไม่มีวันเสร็จ ใน service ที่รันต่อเนื่องเป็นเวลานาน มันจะสะสมเพิ่มขึ้นเรื่อย ๆ อย่างเงียบ ๆ จนหน่วยความจำหรือ file descriptor หมด
- export ค่า
runtime.NumGoroutine()เป็น metric ถ้าเห็นแนวโน้มค่อย ๆ เพิ่มขึ้นทั้งที่ traffic คงที่ แทบจะแน่นอนว่ามี leak - ใช้ endpoint
/debug/pprof/goroutine?debug=2เพื่อดู stack ของ goroutine ทุกตัว leak จะปรากฏเป็น stack หน้าตาเหมือนกันหลายร้อยตัวที่ค้างอยู่ในฟังก์ชันเดียวกัน - ในเทสต์
go.uber.org/goleakสามารถทำให้เทสต์ fail ได้ถ้ามี goroutine หลงเหลืออยู่:
func TestMain(m *testing.M) {
goleak.VerifyTestMain(m)
}ปกป้อง shared state
channel ไม่ใช่คำตอบของปัญหาการประสานงานทุกแบบ counter หรือ cache map ที่ถูกเข้าถึงจากหลาย goroutine มักจะอ่านเข้าใจง่ายกว่าถ้าใช้ sync.Mutex แทนการให้ goroutine ตัวหนึ่งถือครองมันผ่าน channel ให้ใช้ channel เพื่อ ส่งต่อความเป็นเจ้าของข้อมูล และใช้ mutex เพื่อ ควบคุมการเข้าถึงข้อมูลที่ใช้ร่วมกัน
สำหรับ counter ตัวเลขธรรมดา type ใน sync/atomic อย่าง atomic.Int64 ช่วยเลี่ยง overhead ของการ lock ได้ สำหรับ map ที่ถูกอ่านบ่อยกว่าเขียน sync.RWMutex จะช่วยให้ฝั่งอ่านทำงานขนานกันได้ และควรรัน test suite ด้วย -race เสมอ race detector หา bug จริงได้ โดยแทบไม่มีต้นทุนเพิ่มใน CI
Graceful shutdown
service บน production ควรทำงานที่ค้างอยู่ให้เสร็จก่อนปิดตัว วิธีที่นิยมคือใช้ root context ที่จะถูกยกเลิกเมื่อได้รับสัญญาณจาก OS:
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
srv := &http.Server{Addr: ":8080", Handler: router}
go func() {
if err := srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
log.Fatal(err)
}
}()
<-ctx.Done()
shutdownCtx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
defer cancel()
_ = srv.Shutdown(shutdownCtx)ถ้ารันบน Kubernetes ให้ตั้ง shutdown timeout ให้สั้นกว่า terminationGracePeriodSeconds เพื่อให้ process ปิดตัวเองได้ก่อนจะถูก kill
Checklist
- ทุก goroutine มีเจ้าของที่คอยรอมัน
- ทุก operation ที่ block เคารพ
context - ทุกแหล่งที่สร้าง concurrency มีขีดจำกัดบน
- เฉพาะผู้ส่งเท่านั้นที่ปิด channel
- รัน
-raceใน CI และติดตามจำนวน goroutine บน production
การยึดตามกฎเหล่านี้จะทำให้ concurrency ใน Go คาดเดาได้และเข้าใจง่าย
