如何使用频道广播消息

Jin*_*ang 20 concurrency channel go goroutine

我是新手,我正在尝试创建一个简单的聊天服务器,客户端可以向所有连接的客户端广播消息.

在我的服务器中,我有一个goroutine(无限循环)接受连接,所有连接都由一个通道接收.

go func() {
    for {
        conn, _ := listener.Accept()
        ch <- conn
        }
}()
Run Code Online (Sandbox Code Playgroud)

然后,我为每个连接的客户端启动一个处理程序(goroutine).在处理程序内部,我尝试通过迭代通道来广播所有连接.

for c := range ch {
    conn.Write(msg)
}
Run Code Online (Sandbox Code Playgroud)

但是,我不能播放因为(我认为从阅读文档)通道需要在迭代之前关闭.我不确定何时应关闭频道,因为我想继续接受新的连接,关闭频道不会让我这样做.如果有人可以帮助我,或提供更好的方式向所有连接的客户广播消息,我们将不胜感激.

nev*_*ets 36

你正在做的是扇出模式,也就是说,多个端点正在监听单个输入源.这种模式的结果是,只要输入源中有消息,这些侦听器中只有一个能够获取消息.唯一的例外是close渠道.这close将被所有听众识别,因此是"广播".

但你想要做的是广播从连接读取的消息,所以我们可以这样做:

当听众的数量已知时

让每个工作人员收听专用广播信道,并将消息从主信道发送到每个专用广播信道.

type worker struct {
    source chan interface{}
    quit chan struct{}
}

func (w *worker) Start() {
    w.source = make(chan interface{}, 10) // some buffer size to avoid blocking
    go func() {
        for {
            select {
            case msg := <-w.source
                // do something with msg
            case <-quit: // will explain this in the last section
                return
            }
        }
    }()
}
Run Code Online (Sandbox Code Playgroud)

然后我们可以有一堆工人:

workers := []*worker{&worker{}, &worker{}}
for _, worker := range workers { worker.Start() }
Run Code Online (Sandbox Code Playgroud)

然后开始我们的听众:

go func() {
for {
    conn, _ := listener.Accept()
    ch <- conn
    }
}()
Run Code Online (Sandbox Code Playgroud)

和调度员:

go func() {
    for {
        msg := <- ch
        for _, worker := workers {
            worker.source <- msg
        }
    }
}()
Run Code Online (Sandbox Code Playgroud)

当听众的数量未知时

在这种情况下,上面给出的解决方案仍然有效.唯一的区别是,无论何时需要新工作人员,您都需要创建一个新工作人员,启动它,然后将其推入workers切片.但是这种方法需要一个线程安全的切片,需要锁定它.其中一个实现可能如下所示:

type threadSafeSlice struct {
    sync.Mutex
    workers []*worker
}

func (slice *threadSafeSlice) Push(w *worker) {
    slice.Lock()
    defer slice.Unlock()

    workers = append(workers, w)
}

func (slice *threadSafeSlice) Iter(routine func(*worker)) {
    slice.Lock()
    defer slice.Unlock()

    for _, worker := range workers {
        routine(worker)
    }
}
Run Code Online (Sandbox Code Playgroud)

每当你想要开始一个工人时:

w := &worker{}
w.Start()
threadSafeSlice.Push(w)
Run Code Online (Sandbox Code Playgroud)

您的调度员将更改为:

go func() {
    for {
        msg := <- ch
        threadSafeSlice.Iter(func(w *worker) { w.source <- msg })
    }
}()
Run Code Online (Sandbox Code Playgroud)

最后一句话:永远不要留下悬垂的goroutine

其中一个好的做法是:永远不要留下悬垂的goroutine.所以当你听完之后,你需要关闭你开枪的所有goroutine.这将通过以下quit渠道完成worker:

首先,我们需要创建一个全局quit信令通道:

globalQuit := make(chan struct{})
Run Code Online (Sandbox Code Playgroud)

每当我们创建一个worker时,我们globalQuit就将它作为退出信号分配给它:

worker.quit = globalQuit
Run Code Online (Sandbox Code Playgroud)

然后,当我们要关闭所有工人时,我们只需:

close(globalQuit)
Run Code Online (Sandbox Code Playgroud)

由于close将被所有听力goroutines识别(这是你理解的点),所有goroutines将被返回.记得关闭你的调度程序例程,但我会留给你:)


icz*_*cza 13

更优雅的解决方案是"经纪人",客户可以订阅和取消订阅消息.

为了优雅地处理订阅和取消订阅,我们可以为此使用通道,因此接收和分发消息的代理的主循环可以使用单个select语句合并所有这些,并且从解决方案的性质给出同步.

另一个技巧是将订户存储在地图中,从我们用于向其分发消息的频道进行映射.因此,使用通道作为地图中的键,然后添加和删除客户端是"死"简单.这是可能的,因为通道值是可比较的,并且它们的比较非常有效,因为通道值是到通道描述符的简单指针.

不用多说,这是一个简单的经纪人实现:

type Broker struct {
    stopCh    chan struct{}
    publishCh chan interface{}
    subCh     chan chan interface{}
    unsubCh   chan chan interface{}
}

func NewBroker() *Broker {
    return &Broker{
        stopCh:    make(chan struct{}),
        publishCh: make(chan interface{}, 1),
        subCh:     make(chan chan interface{}, 1),
        unsubCh:   make(chan chan interface{}, 1),
    }
}

func (b *Broker) Start() {
    subs := map[chan interface{}]struct{}{}
    for {
        select {
        case <-b.stopCh:
            return
        case msgCh := <-b.subCh:
            subs[msgCh] = struct{}{}
        case msgCh := <-b.unsubCh:
            delete(subs, msgCh)
        case msg := <-b.publishCh:
            for msgCh := range subs {
                // msgCh is buffered, use non-blocking send to protect the broker:
                select {
                case msgCh <- msg:
                default:
                }
            }
        }
    }
}

func (b *Broker) Stop() {
    close(b.stopCh)
}

func (b *Broker) Subscribe() chan interface{} {
    msgCh := make(chan interface{}, 5)
    b.subCh <- msgCh
    return msgCh
}

func (b *Broker) Unsubscribe(msgCh chan interface{}) {
    b.unsubCh <- msgCh
}

func (b *Broker) Publish(msg interface{}) {
    b.publishCh <- msg
}
Run Code Online (Sandbox Code Playgroud)

使用它的示例:

func main() {
    // Create and start a broker:
    b := NewBroker()
    go b.Start()

    // Create and subscribe 3 clients:
    clientFunc := func(id int) {
        msgCh := b.Subscribe()
        for {
            fmt.Printf("Client %d got message: %v\n", id, <-msgCh)
        }
    }
    for i := 0; i < 3; i++ {
        go clientFunc(i)
    }

    // Start publishing messages:
    go func() {
        for msgId := 0; ; msgId++ {
            b.Publish(fmt.Sprintf("msg#%d", msgId))
            time.Sleep(300 * time.Millisecond)
        }
    }()

    time.Sleep(time.Second)
}
Run Code Online (Sandbox Code Playgroud)

上面的输出将是(在Go Playground上试试):

Client 2 got message: msg#0
Client 0 got message: msg#0
Client 1 got message: msg#0
Client 2 got message: msg#1
Client 0 got message: msg#1
Client 1 got message: msg#1
Client 1 got message: msg#2
Client 2 got message: msg#2
Client 0 got message: msg#2
Client 2 got message: msg#3
Client 0 got message: msg#3
Client 1 got message: msg#3
Run Code Online (Sandbox Code Playgroud)

改进

您可以考虑以下改进.根据您使用经纪人的方式/方式,这些可能有用也可能没用.

Broker.Unsubscribe() 可以关闭消息通道,表示不再发送消息:

func (b *Broker) Unsubscribe(msgCh chan interface{}) {
    b.unsubCh <- msgCh
    close(msgCh)
}
Run Code Online (Sandbox Code Playgroud)

这将允许客户端range通过消息通道,如下所示:

msgCh := b.Subscribe()
for msg := range msgCh {
    fmt.Printf("Client %d got message: %v\n", id, msg)
}
Run Code Online (Sandbox Code Playgroud)

然后,如果有人取消订阅msgCh这样:

b.Unsubscribe(msgCh)
Run Code Online (Sandbox Code Playgroud)

上述范围循环将在处理调用之前发送的所有消息后终止Unsubscribe().

如果您希望客户依赖于关闭的消息通道,并且代理的生命周期比应用程序的生命周期更窄,那么您也可以在代理停止时关闭所有订阅的客户端,Start()方法如下:

case <-b.stopCh:
    for msgCh := range subs {
        close(msgCh)
    }
    return
Run Code Online (Sandbox Code Playgroud)


bro*_*man 5

广播到频道切片并使用sync.Mutex 来管理频道添加和删除可能是您的情况下最简单的方法。

您可以在 golang 中执行以下操作broadcast:

  • 您可以使用sync.Cond广播共享状态更改。这种方式不需要任何分配一次设置,但您不能添加超时功能或与其他通道一起使用。
  • 您可以通过关闭旧频道广播共享状态更改并创建新频道和sync.Mutex。这样,每次状态更改都会有一个分配,但您可以添加超时功能并使用另一个通道。
  • 您可以广播到函数回调片段并使用sync.Mutex 来管理它们。调用者可以执行频道操作。这种方式为每个调用者分配多个分配,并与另一个通道一起使用。
  • 您可以广播到一部分频道并使用sync.Mutex 来管理它们。这种方式为每个调用者分配多个分配,并与另一个通道一起使用。
  • 您可以广播到sync.WaitGroup 的一部分并使用sync.Mutex 来管理它们。