【问题标题】:How to have multiple consumer from one io.Reader?如何从一个 io.Reader 拥有多个消费者?
【发布时间】:2020-01-29 01:05:06
【问题描述】:

我正在编写一个小脚本,它使用 bufio.Scannerhttp.Request 以及 go 例程来并行计算单词和行数。

package main

import (
    "bufio"
    "fmt"
    "io"
    "log"
    "net/http"
    "time"
)

func main() {
    err := request("http://www.google.com")

    if err != nil {
        log.Fatal(err)
    }

    // just keep main alive with sleep for now
    time.Sleep(2 * time.Second)
}

func request(url string) error {
    res, err := http.Get(url)

    if err != nil {
        return err
    }

    go scanLineWise(res.Body)
    go scanWordWise(res.Body)

    return err
}

func scanLineWise(r io.Reader) {
    s := bufio.NewScanner(r)
    s.Split(bufio.ScanLines)

    i := 0

    for s.Scan() {
        i++
    }

    fmt.Printf("Counted %d lines.\n", i)
}

func scanWordWise(r io.Reader) {
    s := bufio.NewScanner(r)
    s.Split(bufio.ScanWords)

    i := 0

    for s.Scan() {
        i++
    }

    fmt.Printf("Counted %d words.\n", i)
}

Source

正如scanLineWise 或多或少预期的那样,scalWordWise 将计算一个数字,而scalWordWise 将计算为零。这是因为scanLineWise 已经从req.Body 读取了所有内容。

我想知道:如何优雅地解决这个问题?

我的第一个想法是构建一个实现io.Readerio.Writer 的结构。我们可以使用io.Copy 读取req.Body 并将其写入writer。当扫描器从这个写入器读取数据时,写入器将复制数据而不是读取它。不幸的是,这只会随着时间的推移收集内存并打破流的整个想法......

【问题讨论】:

    标签: stream go


    【解决方案1】:

    选项非常简单——您要么维护数据“流”,要么缓冲正文。

    如果您确实需要多次读取正文而不是按顺序读取一次,则需要在某处缓冲它。没有办法。

    您可以通过多种方式流式传输数据,例如让行计数器输出行到字计数器(最好通过通道)。您还可以使用io.TeeReaderio.Pipe 构建管道,并为每个函数提供一个唯一的读取器。

    ...
    pipeReader, pipeWriter := io.Pipe()
    bodyReader := io.TeeReader(res.Body, pipeWriter)
    go scanLineWise(bodyReader)
    go scanWordWise(pipeReader)
    ...
    

    但是,如果有更多的消费者,这可能会变得笨拙,因此您可以使用 io.MultiWriter 多路复用到更多 io.Readers

    ...
    pipeOneR, pipeOneW := io.Pipe()
    pipeTwoR, pipeTwoW := io.Pipe()
    pipeThreeR, pipeThreeW := io.Pipe()
    
    go scanLineWise(pipeOneR)
    go scanWordWise(pipeTwoR)
    go scanSomething(pipeThreeR)
    
    // of course, this should probably have some error handling
    io.Copy(io.MultiWriter(pipeOneW, pipeTwoW, pipeThreeW), res.Body)
    ...
    

    【讨论】:

    • 我从来没有想过像这样使用io.Pipe(),+1 教我一些新东西!
    • @JimB 如果我有三个消费者,这会是什么样子?不会有一个读者作为开销吗?
    • 对于 3 个或更多,您将不得不制作更多管道并每次再次拆分 bodyReader,但这开始变得混乱。对于这种情况,我也会在答案中使用 MultiWriter 方法。不过,如果可能的话,我更愿意使用频道并逐步分解流。
    【解决方案2】:

    您可以使用频道,在您的scanLineWise 中进行实际阅读,然后将这些行传递给scanWordWise,对于example

    func countLines(r io.Reader) (ch chan string) {
        ch = make(chan string)
        go func() {
            s := bufio.NewScanner(r)
            s.Split(bufio.ScanLines)
    
            cnt := 0
    
            for s.Scan() {
                ch <- s.Text()
                cnt++
            }
            close(ch)
            fmt.Printf("Counted %d lines.\n", cnt)
        }()
    
        return
    }
    
    func countWords(ch <-chan string) {
        cnt := 0
        for line := range ch {
            s := bufio.NewScanner(strings.NewReader(line))
            s.Split(bufio.ScanWords)
            for s.Scan() {
                cnt++
            }
        }
        fmt.Printf("Counted %d words.\n", cnt)
    }
    
    func main() {
        r := strings.NewReader(body)
        ch := countLines(r)
        go countWords(ch)
        time.Sleep(1 * time.Second)
    }
    

    【讨论】:

    • 频道可能会在此处产生不必要的争用。 io 包可能会变得更有用。不过,设计方面的渠道可能会变得更有用。
    猜你喜欢
    • 2011-06-04
    • 1970-01-01
    • 2020-09-12
    • 1970-01-01
    • 1970-01-01
    • 2019-06-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多