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.这将通过以下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)
广播到频道切片并使用sync.Mutex 来管理频道添加和删除可能是您的情况下最简单的方法。
您可以在 golang 中执行以下操作broadcast:
| 归档时间: |
|
| 查看次数: |
18028 次 |
| 最近记录: |