M Receivers N Senders Closed By 3Party
- 多個發送端, 但是要關閉data channel, 讓接收端知道發送已結束
- 為了保持 The Channel Closing Principle. 需要藉由 middle layer 將 N 個 Sender, 轉換成 1 個 sender. 這樣可以在 middle layer 關閉 data channel (M receivers and 1 Sender, 由sender close
dataChchannel) - third party 呼叫
stopfunction -> stop function 傳值到closingchannel -> sender 收到closedchannel 後 return -> middle 收到 closing後 呼叫exit()function 並 return -> exit function 會close(closed)和close(dataCh)
Exercises
https://go.dev/play/p/svDdJMumD-_n
package main
import (
"log"
"math/rand"
"strconv"
"sync"
"time"
)
func main() {
rand.Seed(time.Now().UnixNano())
log.SetFlags(0)
// ...
const Max = 1000000
const NumReceivers = 10
const NumSenders = 1000
const NumThirdParties = 15
wgReceivers := sync.WaitGroup{}
wgReceivers.Add(NumReceivers)
// ...
dataCh := make(chan int) // will be closed
middleCh := make(chan int) // will never be closed
closing := make(chan string) // signal channel
closed := make(chan struct{})
var stoppedBy string
// the middle layer : 接受middleCh 並轉傳至 dataCh.等待 closing 訊號, 然後關閉 closed, dataCh channel
go func() {
exit := func(v int, needSend bool) {
close(closed)
if needSend {
dataCh <- v
}
close(dataCh)
}
for {
select {
case stoppedBy = <-closing:
exit(0, false)
return
case v := <-middleCh:
select {
case stoppedBy = <-closing:
exit(v, true)
return
case dataCh <- v:
}
}
}
}()
// The stop function can be called
// multiple times safely.
stop := func(by string) {
select {
case closing <- by:
<-closed
case <-closed:
}
}
// some third-party goroutines
for i := 0; i < NumThirdParties; i++ {
go func(id string) {
r := 1 + rand.Intn(3)
time.Sleep(time.Duration(r) * time.Second)
stop("3rd-party#" + id)
}(strconv.Itoa(i))
}
// Senders 丟資料到 middleCh
for i := 0; i < NumSenders; i++ {
go func(id string) {
value := rand.Intn(Max)
if value == 0 {
stop("sender#" + id)
return
}
for {
select {
case <-closed:
return
default:
}
select {
case <-closed:
return
case middleCh <- value:
}
}
}(strconv.Itoa(i))
}
// receivers
for range [NumReceivers]struct{}{} {
go func() {
defer wgReceivers.Done()
for value := range dataCh {
log.Println(value)
}
}()
}
// ...
wgReceivers.Wait()
log.Println("stopped by", stoppedBy)
}