You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Go中带错误处理的通道读取批处理实现方案问询

Go通道批量收集的优雅实现优化

原始代码

func concurrentProducer() (<-chan int, <-chan error) {
    prodCh, errCh := make(chan int, 5), make(chan error, 1)

    go func() {
        defer close(prodCh)
        for _ = range 13 {
            prodCh <- rand.Int()
        }
    }()

    return prodCh, errCh
}

func batchConsumer(batch []int) {
    fmt.Println(batch)
}

func batchCollector(ctx context.Context, batchSize int) error {
    prodCh, errCh := concurrentProducer()
    batch := make([]int, 0, batchSize)

    // ... Read a batch but also handle errors!
    batchConsumer(batch)

    return nil
}

func main() {
    err := batchCollector(context.Background(), 5)
    if err != nil {
        log.Fatalln(err)
    }
}

需求说明

实现batchCollector函数需满足:

  • 批量收集prodCh中的数据:当批量大小达到batchSize,或通道无更多数据时,将当前批量传给batchConsumer
  • 响应Context的取消信号,及时返回对应错误
  • 处理errCh中传入的错误

现有实现

本人当前的实现版本如下:

func batchCollector(ctx context.Context, batchSize int) error {
    prodCh, errCh := concurrentProducer()
    batch := make([]int, 0, batchSize)

    elem, consuming := 0, true
    for consuming {
        select {
        case <-ctx.Done():
            return ctx.Err()
        case err := <-errCh:
            return err
        case elem, consuming = <-prodCh:
            if !consuming {
                break
            }
            batch = append(batch, elem)
            if len(batch) < batchSize {
                continue
            }
        }

        batchConsumer(batch)
        batch = batch[:0]
    }

    return nil
}

疑问

当前实现的控制流较为复杂,请问是否存在更符合Go惯用模式的优雅实现方式?

内容的提问来源于stack exchange,提问作者Niklas H.

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.14 04:48:22