M Receivers N Senders

  • 多個接收, 多個發送. 引用 moderator 角色來關閉額外的訊號通到close(stopCh)
  • 能讓任何接收方和發送方關閉數據通道。
  • 並且我們不能讓任何接收者關閉額外的信號通道來通知所有發送者和接收者退出遊戲
  • 引入一個主持人角色(moderator)來關閉額外的信號通道
    // moderator
    go func() {
      stoppedBy = <-toStop
      close(stopCh)
    }()
    
  • toStop := make(chan string, 1) 此channel 的size為1, 避免第一個notification 錯過, 當此notification被發送是在 moderator goroutine 啟動之前, 還沒準備好接收 toStop
  • 發送端/接收端 發送值到 toStop channel -> moderator 收到 toStop 有直, close(stopCh) -> 發送端/接收端 接收到 stopCh就 return -> 等待接收端 WaitGroup跑完

Exercises

https://go.dev/play/p/8g1u2D1-p2l

package main

import (
    "log"
    "math/rand"
    "strconv"
    "sync"
    "time"
)

func main() {
    rand.Seed(time.Now().UnixNano())
    log.SetFlags(0)

    // ...
    const Max = 100000
    const NumReceivers = 10
    const NumSenders = 1000

    wgReceivers := sync.WaitGroup{}
    wgReceivers.Add(NumReceivers)

    // ...
    dataCh := make(chan int)
    stopCh := make(chan struct{}) // for sender, receiver to stop
    // stopCh is an additional signal channel.
    // Its sender is the moderator goroutine shown
    // below, and its receivers are all senders
    // and receivers of dataCh.
    toStop := make(chan string, 1)
    // The channel toStop is used to notify the
    // moderator to close the additional signal
    // channel (stopCh). Its senders are any senders
    // and receivers of dataCh, and its receiver is
    // the moderator goroutine shown below.
    // It must be a buffered channel.

    var stoppedBy string

    // moderator
    go func() {
        stoppedBy = <-toStop
        close(stopCh)
    }()

    // senders
    for i := 0; i < NumSenders; i++ {
        go func(id string) {
            for {
                value := rand.Intn(Max)
                if value == 0 {
                    // Here, the try-send operation is
                    // to notify the moderator to close
                    // the additional signal channel.
                    select {
                    case toStop <- "sender#" + id:
                    default:
                    }
                    return
                }

                // The try-receive operation here is to
                // try to exit the sender goroutine as
                // early as possible. Try-receive and
                // try-send select blocks are specially
                // optimized by the standard Go
                // compiler, so they are very efficient.
                select {
                case <-stopCh:
                    return
                default:
                }

                // Even if stopCh is closed, the first
                // branch in this select block might be
                // still not selected for some loops
                // (and for ever in theory) if the send
                // to dataCh is also non-blocking. If
                // this is unacceptable, then the above
                // try-receive operation is essential.
                select {
                case <-stopCh:
                    return
                case dataCh <- value:
                }
            }
        }(strconv.Itoa(i))
    }

    // receivers
    for i := 0; i < NumReceivers; i++ {
        go func(id string) {
            defer wgReceivers.Done()
            for {
                // Same as the sender goroutine, the
                // try-receive operation here is to
                // try to exit the receiver goroutine
                // as early as possible.
                select {
                case <-stopCh:
                    return
                default:
                }

                // Even if stopCh is closed, the first
                // branch in this select block might be
                // still not selected for some loops
                // (and forever in theory) if the receive
                // from dataCh is also non-blocking. If
                // this is not acceptable, then the above
                // try-receive operation is essential.
                select {
                case <-stopCh:
                    return
                case value := <-dataCh:
                    if value == Max-1 {
                        // Here, the same trick is
                        // used to notify the moderator
                        // to close the additional
                        // signal channel.
                        select {
                        case toStop <- "receiver#" + id:
                        default:
                        }
                        return
                    }
                    log.Println(value)
                }
            }

        }(strconv.Itoa(i))
    }
    // ...
    wgReceivers.Wait()
    log.Println("stopped by", stoppedBy)
}
© Kimi Tsai all right reserved.            Updated : 2023-07-12 09:04:54

results matching ""

    No results matching ""

    results matching ""

      No results matching ""