【问题标题】:Should one drain a buffered channel when closing it关闭缓冲通道时是否应该排出缓冲通道
【发布时间】:2016-06-20 00:59:22
【问题描述】:

在 Go 中给定一个(部分)填充的缓冲通道

ch := make(chan *MassiveStruct, n)
for i := 0; i < n; i++ {
    ch <- NewMassiveStruct()
}

是否建议在关闭通道时(由作者)也将其排空,以防读者何时读取它是未知的(例如,这些通道的数量有限并且他们目前很忙)?那是

close(ch)
for range ch {}

如果通道上有其他并发读者,这样的循环是否保证结束?

上下文:具有固定数量的工作人员的队列服务,它应该在服务停止时停止处理排队的任何内容(但不一定在之后立即被 GC)。因此,我将关闭以向工作人员表明该服务正在终止。我可以立即耗尽剩余的“队列”,让 GC 释放分配的资源,我可以读取并忽略工作人员中的值,我可以离开通道,因为正在运行阅读器并将通道设置为 nil,以便GC 清理一切。我不确定哪种方法最干净。

【问题讨论】:

  • 这完全取决于你程序中的逻辑。关闭通道不是清理操作,黄油通道不需要为空即可进行 GC,因此这些都不是必需的。

标签: go channel


【解决方案1】:

我觉得除了提示既不需要排水也不需要关闭之外,提供的答案实际上并没有澄清太多。因此,对于所描述的上下文,以下解决方案对我来说看起来很干净,它终止了工作人员并删除了对他们或有问题的通道的所有引用,因此,让 GC 清理通道及其内容:

type worker struct {
    submitted chan Task
    stop      chan bool
    p         *Processor
}

// executed in a goroutine
func (w *worker) run() {
    for {
        select {
        case task := <-w.submitted:
            if err := task.Execute(w.p); err != nil {
                logger.Error(err.Error())
            }
        case <-w.stop:
            logger.Warn("Worker stopped")
            return
        }
    }
}

func (p *Processor) Stop() {
    if atomic.CompareAndSwapInt32(&p.status, running, stopped) {
        for _, w := range p.workers {
            w.stop <- true
        }
        // GC all workers as soon as goroutines stop
        p.workers = nil
        // GC all published data when workers terminate
        p.submitted = nil
        // no need to do the following above:
        // close(p.submitted)
        // for range p.submitted {}
    }
}

【讨论】:

  • 我写答案时,您的上下文段落不在这里。
  • 现在,我觉得同时使用原子操作和同步例程是错误的:如果您要在此之后立即锁定,则无需执行 CompareAndSwap... 在您的情况下,您可能只是希望您的Stop 方法将布尔标志设置为true,并让您的工作人员检查该布尔值以跳出循环。根据您的要求,这甚至不需要任何同步。
  • @Elwinar 抱歉,您发布答案时我没有上下文。我只是在阅读后才意识到它会有所帮助。关于您关于锁定的评论:锁定与这个问题无关,但这里原子在提交新任务(没有锁定)时被检查,锁定仅用于启动和停止。但是感谢您的评论,我意识到如何通过更改启动方法(此处未显示)将其省略。
【解决方案2】:

有更好的方法来实现您想要实现的目标。您当前的方法可能会导致丢弃一些记录,并随机处理其他记录(因为耗尽循环正在与所有消费者竞争)。这并没有真正实现目标。

您想要的是取消。这是来自Go Concurrency Patterns: Pipelines and cancellation的示例

func sq(done <-chan struct{}, in <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for n := range in {
            select {
            case out <- n * n:
            case <-done:
                return
            }
        }
    }()
    return out
}

您将done 通道传递给所有goroutine,并在您希望它们全部停止处理时关闭它。如果你经常这样做,你可能会发现 golang.org/x/net/context 包很有用,它形式化了这种模式,并添加了一些额外的功能(如超时)。

【讨论】:

  • 我现在使用这种模式并且会坚持下去。毕竟感觉我唯一需要做的就是在发布者中将通道设置为 nil ,以便在工作人员终止后,GC 清理通道及其内容。
  • 为什么要使用和传递done频道?我想只需要带有in 频道的 for 循环。
【解决方案3】:

这取决于你的程序,但一般来说,我倾向于说不(你不需要在关闭之前清除频道):如果你的频道在你关闭时有项目,任何读者仍在阅读通道将接收项目,直到通道为空。

这是一个例子:

package main

import (
    "sync"
    "time"
)

func main() {

    var ch = make(chan int, 5)
    var wg sync.WaitGroup
    wg.Add(1)

    for range make([]struct{}, 2) {
        go func() {
            for i := range ch {
                wg.Wait()
                println(i)
            }
        }()
    }

    for i := 0; i < 5; i++ {
        ch <- i
    }
    close(ch)

    wg.Done()
    time.Sleep(1 * time.Second)
}

在这里,程序将输出所有项目,尽管通道在任何阅读器甚至可以从通道中读取之前就已严格关闭。

【讨论】:

    猜你喜欢
    • 2014-01-07
    • 2014-01-30
    • 2018-12-01
    • 2021-10-13
    • 2014-11-05
    • 2017-09-06
    • 2016-08-30
    • 1970-01-01
    • 2019-01-07
    相关资源
    最近更新 更多