【问题标题】:Asynchronous messages golang异步消息golang
【发布时间】:2015-04-12 09:35:08
【问题描述】:

我有一个 golang 服务器做这样的事情: 主包

func main() {
    for {
        c := listener.Accept()
        go handle(c)
    }
}

...
func handle(c net.Conn) {
    m := readMessage(c)    // func(net.Conn)Message
    r := processMessage(m) //func(Message)Result
    sendResult(c, r)       // func(net.Conn,Result)
}

同步读取和写入消息。我现在需要的是通过给定的开放连接异步发送消息,我知道我有点迷路了。

这是我的想法:

...
func someWhereElese(c chan Result) {
    // generate a message and a result
    r := createResultFromSomewhere()
    c <- r // send the result through the channel
}

并修改我的句柄以使用相同的频道

func handle(c net.Conn, rc chan Result) {
    m := readMessage(c)    // func(net.Conn)Message
    r := processMessage(m) //func(Message)Result
    //sendResult(c, r)       // func(net.Conn,Result)
    rc <- r
}

这就是我的困惑所在。

结果通道应该被创建并且它应该有一个连接来发送它收到的任何东西

func doSend(c net.Con, rc chan Result) {
    r := rc          // got a result channel
    sendResult(c, r) // send it through the wire
}

但是应该在哪里创建该频道?在主循环中?

func main() {
    ...
    for {
        c := l.Accept()
        rc := make(chan Result)
        go doSend(c, rc)
    }
}

阅读怎么样?它应该进入它自己的频道/gorutine吗? 如果我需要向 n 个客户广播,我应该保留一部分结果频道吗?一片连接?

我在这里有点困惑,但我觉得我很接近。

【问题讨论】:

  • 从这里开始:talks.golang.org/2012/concurrency.slide。您的问题的关键是使用select 来观察传入的连接通道、响应通道,也许还有退出通道。根据您的预期负载,为每个请求创建一个 goroutine 可能没问题;或者您可能需要创建一个池。但是,一旦您了解了 Go 并发模式,您就会更好地理解正确的问题。另请参阅blog.golang.org/advanced-go-concurrency-patterns,尤其是末尾的“相关文章”。理解select 是伟大的“哦,这就是 Go 的工作原理!”之一。时刻。
  • @RobNapier 嗯...我之前通读了幻灯片,我有点明白了,但不完全明白。稍后我将再次浏览它们。同时我设法做了一个小程序,如果您发现有什么特别危险的地方,您能评论一下吗?
  • 好的;我可能误解了您所说的“异步”是什么意思。我曾假设多个请求和响应会发生在同一个连接上(交错请求)。您似乎的意思是您想在读者获取数据时流式传输数据。这可能比在单个 goroutine 中读取两个字节并写入两个字节更并行一些。为此,您可能需要查看golang.org/x/text/transform 以及io.Copy()。 (对不起,我对这里的实际代码没有更多帮助;这是一个有趣的问题,但我现在没有时间提供广泛的帮助。)
  • 不,你是对的。虽然这是第一步。如果没有通道,我无法在读取之间进行另一次写入,并且流程将锁定,直到我有其他内容要读取。现在有了这个变化,我可以“广播”来自其他连接的消息,而无需等待读取发生。下一步将从单个 goroutine 读取和/或使用池以避免为每个打开的连接创建 goroutine。

标签: asynchronous go channel goroutine


【解决方案1】:

这个程序似乎解决了我的直接问题

package main

import (
    "bytes"
    "encoding/binary"
    "log"

    "net"
)

var rcs []chan int = make([]chan int,0)


func main() {
    a, e := net.ResolveTCPAddr("tcp", ":8082")
    if e != nil {
        log.Fatal(e)
    }
    l, e := net.ListenTCP("tcp", a)
    for {
        c, e := l.Accept()
        if e != nil {
            log.Fatal(e)
        }
        rc := make(chan int)
        go read(c, rc)
        go write(c, rc)
        rcs = append(rcs, rc)
        // simulate broacast
        log.Println(len(rcs))
        if len(rcs) > 5 {
            func() {
                for _, v := range rcs {
                    log.Println("sending")
                    select {
                    case v <- 34:
                        log.Println("done sending")
                    default:
                        log.Println("didn't send")
                    }
                }
            }()
        }
    }
}
func read(c net.Conn, rc chan int) {
    h := make([]byte, 2)
    for {
        _, err := c.Read(h)
        if err != nil {
            rc <- -1
        }
        var v int16
        binary.Read(bytes.NewReader(h[:2]), binary.BigEndian, &v)
        rc <- int(v)
    }
}
func write(c net.Conn, rc chan int) {
    for {
        r := <-rc
        o := []byte{byte(r * 2)}
        c.Write(o)
    }
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2022-06-10
    • 2015-11-25
    • 2011-04-05
    • 2010-12-12
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多