【问题标题】:How to implement pipeline to goroutines?如何实现管道到 goroutines?
【发布时间】:2019-07-03 15:05:45
【问题描述】:

我需要一些帮助来了解如何使用管道获取数据以从一个 goroutine 传输到另一个。

我看了golang blogpost on pipeline,我明白了,但不能完全付诸行动,因此想向社区寻求帮助。

现在,我想出了这个丑陋的代码(Playground):

package main

import (
    "fmt"
    "sync"
    "time"
)

func main() {
    wg := sync.WaitGroup{}
    ch := make(chan int)
    for a := 0; a < 3; a++ {
        wg.Add(1)
        go func1(int(3-a), ch, &wg)
    }
    go func() {
        wg.Wait()
        close(ch)
    }()
    wg2 := sync.WaitGroup{}
    ch2 := make(chan string)
    for val := range ch {
        fmt.Println(val)
        wg2.Add(1)
        go func2(val, ch2, &wg2)
    }
    go func() {
        wg2.Wait()
        close(ch2)
    }()
    for val := range ch2 {
        fmt.Println(val)
    }
}

func func1(seconds int, ch chan<- int, wg *sync.WaitGroup) {
    defer wg.Done()
    time.Sleep(time.Duration(seconds) * time.Second)
    ch <- seconds
}

func func2(seconds int, ch chan<- string, wg *sync.WaitGroup) {
    defer wg.Done()
    ch <- "hello"
}

问题

我想以正确的方式使用 管道 或任何正确的方式。

另外,博文中显示的管道不适用于goroutines,因此我无法自己完成。

在现实生活中,func1func2 是从网络获取资源的函数,因此它们在自己的 goroutine 中启动。

谢谢。
临时工
(一个golang noobie)

P.S. 现实生活中的例子和使用 goroutine 的管道的用法也会有很大帮助。

【问题讨论】:

    标签: go


    【解决方案1】:

    该管道帖子的关键模式是,您可以将通道的内容视为数据流,并编写一组协作的 goroutine 来构建数据处理流图。这可能是一种在面向数据的应用程序中获得一些并发性的方法。

    在设计方面,您可能还会发现构建与 goroutine 结构无关的块并将它们包装在通道中会更有帮助。这使得测试低级代码变得更加容易,并且如果您改变主意是否在 goroutine 中运行事物,添加或删除包装器会更容易。

    因此,在您的示例中,我首先将最低级别的任务重构为它们自己的(同步)函数:

    func fetch(ms int) int {
        time.Sleep(time.Duration(ms) * time.Millisecond)
        return ms
    }
    
    func report(ms int) string {
        return fmt.Sprintf("Hello after %d ms", ms)
    }
    

    由于示例的后半部分是相当同步的,因此很容易适应管道模式。我们编写了一个函数,它消耗其所有输入流并产生一个完整的输出流,并在完成后关闭它。

    func reportAll(mss <-chan int, out chan<- string) {
        for ms := range mss {
            out <- report(ms)
        }
        close(out)
    }
    

    调用异步代码的函数有点小技巧。在函数的主循环中,每次读取一个值,都需要启动一个 goroutine 来处理它。然后,在您从输入通道中读取所有内容之后,您需要等待所有这些 goroutine 完成,然后再关闭输出通道。您可以在这里使用一个小的匿名函数来提供帮助。

    func fetchAll(mss <-chan int, out chan<- int) {
        var wg sync.WaitGroup
        for ms := range mss {
            wg.Add(1)
            go func(ms int) {
                out <- fetch(ms)
                wg.Done()
            }(ms)
        }
        wg.Wait()
        close(out)
    }
    

    这里也有帮助(因为通道写入是阻塞的)编写另一个函数来播种输入值。

    func produceInputs(mss chan<- int) {
        for ms := 1000; ms > 0; ms -= 300 {
            mss <- ms
        }
        close(mss)
    }
    

    现在您的 main 函数需要在它们之间创建通道并运行最终消费者。

    // main is the entry point to the program.
    //
    //                   mss        fetched       results
    //     produceInputs --> fetchAll --> reportAll --> main
    func main() {
        mss := make(chan int)
        fetched := make(chan int)
        results := make(chan string)
    
        go produceInputs(mss)
        go fetchAll(mss, fetched)
        go reportAll(fetched, results)
    
        for val := range results {
            fmt.Println(val)
        }
    }
    

    https://play.golang.org/p/V9Z7ECUVIJL 是一个完整的例子。

    我已经避免在这里手动传递sync.WaitGroups(并且通常倾向于这样做:除非您明确地将某些东西称为 goroutine 的顶级,否则您不会拥有 WaitGroup,因此推动 WaitGroup对调用者的管理使代码更加模块化;有关示例,请参见上面的 fetchAll 函数)。我怎么知道我所有的 goroutine 都完成了?我们可以追溯:

    • 如果我已到达main 的末尾,则results 频道将关闭。
    • results通道为reportAll的输出通道;如果它关闭,则该函数到达其执行结束;如果发生这种情况,则 fetched 频道将关闭。
    • fetched通道为fetchAll的输出通道; ...

    另一种看待这个问题的方式是,一旦管道的源 (produceInputs) 关闭其输出通道并完成,“我完成了”信号就会沿着管道向下流动并导致下游步骤关闭它们的输出频道和完成。

    博客文章提到了一个单独的明确关闭渠道。我根本没有在这里讨论。然而,自从它被编写后,标准库获得了context 包,它现在是管理这些包的标准习惯用法。您需要在主循环的主体中使用select 语句,这会使处理变得更加复杂。这可能看起来像:

    func reportAllCtx(ctx context.Context, mss <-chan int, out chan<- string) {
        for {
            select {
                case <-ctx.Done():
                    break
                case ms, ok := <-mss:
                    if ok {
                        out <- report(ms)
                    } else {
                        break
                    }
                }
            }
        }
        close(out)
    }
    

    【讨论】:

      【解决方案2】:

      这个article 涵盖了一个端口扫描器示例中的管道模式,推荐看看。

      端口扫描器旨在探测服务器或主机的开放端口

      上图显示了端口扫描器的整个管道。在下一节中,我们将逐个解释每个相关的功能。

      • 函数初始化
      • Func parsePortsToScan
      • 结构扫描操作
      • 功能基因
      • 功能扫描
      • 功能过滤器
      • 功能商店
      • 主函数

      init 函数定义了用户传入的参数。 ports 变量是要扫描的端口字符串,用破折号分隔。 outFile 变量是写入结果的文件。

      var ports string
      var outFile string
      
      func init() {
          flag.StringVar(&ports, "ports", "80", "Port(s) (e.g. 80, 22-100).")
          flag.StringVar(&outFile, "outfile", "scans.csv", 
          "Destination of scan results (defaults to scans.csv)")
      }
      

      主函数负责执行函数的管道。它从命令行参数中获取一片 int 端口和一个字符串输出文件。

      func main() {
          flag.Parse()
      
          portsToScan, err := parsePortsToScan(ports)
          if err != nil {
              fmt.Printf("Failed to parse ports to scan: %s\n", err)
              os.Exit(1)
          }
      
          dest, err := os.Create(outFile)
          if err != nil {
              fmt.Printf("Failed to create scan results destination: %s\n", err)
              os.Exit(2)
          }
      
          // pipeline
          // scanChan := store(dest, filter(scan(gen(portsToScan...))))
      
          // broken up for explainability
          var scanChan <-chan scanOp
          scanChan = gen(portsToScan...)
          scanChan = scan(scanChan)
          scanChan = filter(scanChan)
          scanChan = store(dest, scanChan)
      
          for s := range scanChan {
              if !s.open && s.scanErr != fmt.Sprintf("dial tcp 127.0.0.1:%d: connect: connection refused", s.port) {
                  fmt.Println(s.scanErr)
              }
          }
      }
      

      parsePortsToScan 函数从命令行参数解析要扫描的端口。如果参数无效,则返回错误。如果参数有效,则返回一个整数切片。

      func parsePortsToScan(portsFlag string) ([]int, error) {
          p, err := strconv.Atoi(portsFlag)
          if err == nil {
              return []int{p}, nil
          }
      
          ports := strings.Split(portsFlag, "-")
          if len(ports) != 2 {
              return nil, errors.New("unable to determine port(s) to scan")
          }
      
          minPort, err := strconv.Atoi(ports[0])
          if err != nil {
              return nil, fmt.Errorf("failed to convert %s to a valid port number", ports[0])
          }
      
          maxPort, err := strconv.Atoi(ports[1])
          if err != nil {
              return nil, fmt.Errorf("failed to convert %s to a valid port number", ports[1])
          }
      
          if minPort <= 0 || maxPort <= 0 {
              return nil, fmt.Errorf("port numbers must be greater than 0")
          }
      
          var results []int
          for p := minPort; p <= maxPort; p++ {
              results = append(results, p)
          }
          return results, nil
      }
      

      scanOp 表示单端口扫描操作及其结果(open、scanErr、scanDuration)。 open 是一个布尔值,指示端口是否打开。如果扫描失败,scanErr 是一条错误消息。scanDuration 是执行扫描所花费的时间。

      要将结果输出到 CSV 文件,CSV 编写器使用两种方法。 csvHeaders 返回字符串切片中的标题。 asSlice 将 scanOp 的值字段作为字符串切片返回。

      type scanOp struct {
          port         int
          open         bool
          scanErr      string
          scanDuration time.Duration
      }
      
      func (so scanOp) csvHeaders() []string {
          return []string{"port", "open", "scanError", "scanDuration"}
      }
      
      func (so scanOp) asSlice() []string {
          return []string{
              strconv.FormatInt(int64(so.port), 10),
              strconv.FormatBool(so.open),
              so.scanErr,
              so.scanDuration.String(),
          }
      }
      

      gen 函数是一个生成器函数,它从 int 端口切片返回一个 scanOps 结构值的缓冲通道。它用于创建将按顺序执行的函数管道,它是管道中的第一个函数。

      func gen(ports ...int) <-chan scanOp {
          out := make(chan scanOp, len(ports))
          go func() {
              defer close(out)
              for _, p := range ports {
                  out <- scanOp{port: p}
              }
          }()
          return out
      }
      

      扫描函数负责执行实际的端口扫描。它接受一个有缓冲的 scanOps 通道,并返回一个无缓冲的 scanOps 通道。

      func scan(in <-chan scanOp) <-chan scanOp {
          out := make(chan scanOp)
          go func() {
              defer close(out)
              for scan := range in {
                  address := fmt.Sprintf("127.0.0.1:%d", scan.port)
                  start := time.Now()
                  conn, err := net.Dial("tcp", address)
                  scan.scanDuration = time.Since(start)
                  if err != nil {
                      scan.scanErr = err.Error()
                  } else {
                      conn.Close()
                      scan.open = true
                  }
                  out <- scan
              }
          }()
          return out
      }
      

      filter函数负责过滤打开的scanOps。

      func filter(in <-chan scanOp) <-chan scanOp {
          out := make(chan scanOp)
          go func() {
              defer close(out)
              for scan := range in {
                  if scan.open {
                      out <- scan
                  }
              }
          }()
          return out
      }
      

      store 函数负责将 scanOps 存储在 CSV 文件中。它是管道中的最后一个函数。

      func store(file io.Writer, in <-chan scanOp) <-chan scanOp {
          csvWriter := csv.NewWriter(file)
          out := make(chan scanOp)
          go func() {
              defer csvWriter.Flush()
              defer close(out)
              var headerWritten bool
              for scan := range in {
                  if !headerWritten {
                      headers := scan.csvHeaders()
                      if err := csvWriter.Write(headers); err != nil {
                          fmt.Println(err)
                          break
                      }
                      headerWritten = true
                  }
                  values := scan.asSlice()
                  if err := csvWriter.Write(values); err != nil {
                      fmt.Println(err)
                      break
                  }
              }
          }()
      
          return out
      }
      

      通道可用于将 goroutine 连接在一起,以便一个的输出是另一个的输入。当您的管道中有许多功能并想要连接它们时,它非常有用。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2011-02-09
        • 2018-01-05
        • 2012-07-04
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多