【问题标题】:Use multiple go routines to pull from a channel使用多个 goroutine 从通道中拉取
【发布时间】:2020-07-07 16:19:14
【问题描述】:

我有一个接收带有此代码的地图切片的频道:

func submitRecords(records []map[string]string) {
    batch := []map[string]string{}
    ch := make(chan []map[string]string)
    batchCt := 1
    go func() {
        for _, v := range records {
            batch = append(batch, v)
            if len(batch) == 150 {
                ch <- batch
                batch = nil
            }
        }
        close(ch)
    }()
}

我提交这些记录的 API 接受最多 150 个批次。为了加快速度,我想启动 4 个 go 例程来同时处理通道中的记录。一旦记录进入频道,它们被处理的顺序就无关紧要了。

目前我对上面单独运行的代码进行了以下更新:

func submitRecords(records []map[string]string) {
    batch := []map[string]string{}
    ch := make(chan []map[string]string)
    batchCt := 1
    go func() {
        for _, v := range records {
            batch = append(batch, v)
            if len(batch) == 150 {
                ch <- batch
                batch = nil
            }
        }
        close(ch)
    }()

    for b := range ch {
        str, _ := json.Marshal(b)
        fmt.Printf("Sending batch at line %d\n", (batchCt * 150))

        payload := strings.NewReader(string(str))
        client := &http.Client{}
        req, err := http.NewRequest(method, url, payload)
        if err != nil {
            fmt.Println(err)
        }

        login, _ := os.LookupEnv("Login")
        password, _ := os.LookupEnv("Password")
        req.Header.Add("user_name", login)
        req.Header.Add("password", password)
        req.Header.Add("Content-Type", "application/json")

        res, err := client.Do(req)
        if err != nil {
            fmt.Println(err)
        }
        batchCt++
    }
}

我将如何修改它以从通道中提取 4 个 go 例程并发送这些请求?或者,这是否可能/我是否误解了 goroutine 的功能?

【问题讨论】:

  • batch = nil 这行不行。
  • 为什么?根据我的研究,我发现这是清除切片以重复使用的方法,类似于 python 中的batch.clear() 操作。当我在将记录发送到 API 后检查记录时,看起来一切正常。
  • 您的异步循环缺少尾随 ch &lt;- batch,我猜。
  • 你的函数是颠倒的。当前正在异步的循环应该是同步的,应该有N个请求处理例程。该函数应该通过等待这 N 个 goroutine 完成来终止。
  • 基本上,我将记录添加到 batch 切片,直到切片有 150 条记录。一旦达到 150,我将其发送到通道并清除切片。我对 Go 中的并发性比较陌生,所以我可能把一些东西弄混了,它看起来像是在工作

标签: go concurrency


【解决方案1】:
func process(ch chan []map[string]string) {
    for b := range ch {
        str, _ := json.Marshal(b)

        // This wont work, or has to be included in the payload from the channel
        // fmt.Printf("Sending batch at line %d\n", (batchCt * 150))

        payload := strings.NewReader(string(str))
        client := &http.Client{}
        req, err := http.NewRequest(method, url, payload)
        if err != nil {
            fmt.Println(err)
        }

        login, _ := os.LookupEnv("Login")
        password, _ := os.LookupEnv("Password")
        req.Header.Add("user_name", login)
        req.Header.Add("password", password)
        req.Header.Add("Content-Type", "application/json")

        res, err := client.Do(req)
        if err != nil {
            fmt.Println(err)
        }
        // batchCt++
    }
    done <- 1
}

func submitRecords(records []map[string]string) {
    batch := []map[string]string{}
    ch := make(chan []map[string]string)
    
    go process(ch)
    go process(ch)
    go process(ch)
    go process(ch)

    // batchCt := 1
    for _, v := range records {
        batch = append(batch, v)
        if len(batch) == 150 {
            ch <- batch
            batch = []map[string]string{}
        }
    }
    // Send the last not to size batch
    ch <- batch
    close(ch)
}

使用 int 和 sleep https://play.golang.org/p/q5bUhXt9aUn 的 Playground 示例

package main

import (
    "fmt"
    "time"
)

func process(processor int, ch chan []int, done chan int) {
    for batch := range ch {
        // Do something.. sleep or http requests will let the other workers work as well
        time.Sleep(time.Duration(len(batch)) * time.Millisecond)
        fmt.Println(processor, batch)
    }
    done <- 1
}

const batchSize = 3

func main() {
    records := []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13}
    ch := make(chan []int)
    done := make(chan int)

    go process(1, ch, done)
    go process(2, ch, done)
    go process(3, ch, done)
    go process(4, ch, done)

    batch := make([]int, 0, batchSize)
    for _, v := range records {
        batch = append(batch, v)
        if len(batch) == batchSize {
            ch <- batch
            batch = make([]int, 0, batchSize)
        }
    }
    ch <- batch
    close(ch)

    <-done
    <-done
    <-done
    <-done
}

【讨论】:

  • 在上一个示例中使用 WaitGroup 而不是通道不是更好吗?
  • @Kent 我采用了这种方法,但并不是每条记录都被发送。在 100k 中发送了大约 95k。没有数据竞争或服务器端错误。有什么想法吗?
  • @DBA108642 猜测该程序在有时间发送所有消息之前就存在,即为什么我使用 done 通道,但正如 mh-cbon 提到的那样,可能应该使用 waitgroups 作为更清洁的解决方案.或者最后一批没有发送。因为它不是一个完整的 150(但是你应该看到几乎所有的记录,除了一个子 150)
猜你喜欢
  • 2018-10-23
  • 1970-01-01
  • 1970-01-01
  • 2020-08-02
  • 2015-06-03
  • 2019-10-08
  • 2022-01-02
  • 2019-04-08
  • 2018-09-29
相关资源
最近更新 更多