【问题标题】:Deadlock in event processing事件处理中的死锁
【发布时间】:2016-05-02 12:22:16
【问题描述】:

所以我有一个用于事件处理的通道,主服务器 goroutine 在这个通道上选择并在收到的每个事件上调用事件处理程序:

evtCh := make(chan Event)
// server loop:
for !quit {
    select {
    case e := <- evtCh:
        handleEvent(e)
        break
    case quit := <-quitCh:
        //finish
}

// for send a new event to processing
func addEvent(e Event) {
    evtCh <- e
}

handleEvent 将在事件类型上调用已注册的处理程序。我有func registerEventHandler(typ EventType, func(Event)) 来处理寄存器。该程序将支持用户编写扩展,这意味着他们可以注册自己的处理程序来处理事件。

现在问题出现在用户的事件处理程序中,他们可能通过调用addEvent向服务器发送新事件,这将导致服务器挂起,因为事件处理程序本身是在服务器主循环的上下文中调用的(在 for 循环中)。

我该如何优雅地处理这种情况?用切片建模的队列是个好主意吗?

【问题讨论】:

    标签: go goroutine


    【解决方案1】:

    这将导致服务器挂起,因为事件处理程序本身是在服务器主循环的上下文中调用的

    主循环不应该在调用 handleEvent 时阻塞,最常见的避免这种情况的方法是使用一个工作 goroutine 池。下面是一个未经测试的快速示例:

    type Worker struct {
        id int
        ch chan Event
        quit chan bool
    }
    
    func (w *Worker) start {
        for {
            select {
                case e := <- w.ch:
                    fmt.Printf("Worker %d called\n", w.id)
                    //handle event
                    break;
                case <- w.quit:
                    return
            }
        }
    }
    
    
    ch := make(chan Event, 100)
    quit := make(chan bool, 0)
    
    // Start workers
    for i:=0; i<10; i++{
        worker := &Worker{i,ch,quit}
        go worker.start()
    }
    
    // 
    func addEvent (e Event) {
        ch <- e
    }
    

    完成后,只需 close(quit) 杀死所有工人。

    编辑:来自下面的 cmets:

    在这种情况下,主循环是什么样的?

    视情况而定。如果您有固定数量的事件,则可以使用WaitGroup,如下所示:

    type Worker struct {
        id int
        ch chan Event
        quit chan bool
        wg *sync.WaitGroup
    }
    
    func (w *Worker) start {
        for {
            select {
                case e := <- w.ch:
                    //handle event
                    wg.Done()
    
                    break;
                case <- w.quit:
                    return
            }
        }
    }
    
    func main() {
        ch := make(chan Event, 100)
        quit := make(chan bool, 0)
    
        numberOfEvents := 100
    
        wg := &sync.WaitGroup{}
        wg.Add(numberOfEvents)
    
        // start workers
        for i:=0; i<10; i++{
            worker := &Worker{i,ch,quit,wg}
            go worker.start()
        }
    
    
        wg.Wait() // Blocks until all events are handled
    }
    

    如果事先不知道事件的数量,你可以在退出频道上阻塞:

    <- quit
    

    一旦另一个 goroutine 关闭了通道,你的程序也会终止。

    【讨论】:

    • 谢谢,在这种情况下主循环是什么样的?您的意思是以循环方式将 e 发送给工人?
    • 你不需要为循环赛做任何额外的事情。当您调用addEvent(myEvent) 时,第一个空闲的工作人员将获取该事件。我用一个主线程的例子扩展了答案。
    • 我想了解为什么将ch设为缓冲通道,与无缓冲通道相比有什么好处?
    • @fluter 缓冲通道将保护您的发送 goroutine 免受意外阻塞。如果偶尔,你所有的工作 goroutines 都忙并且你调用addEvent(),这个调用将会阻塞。即使所有工作人员都很忙,缓冲通道也可以让您向其推送事件。
    • 按照规范,频道满了也会阻塞,对吧?虽然我猜实际上机会很小。
    【解决方案2】:

    为了使事情更加异步,您可以

    • 为事件通道增加容量

      evtCh := make(chan Event, 10)

    • 异步调用handleEvent(e)

      go handleEvent(e)

    • 在处理程序中异步添加事件

      go addEvent(e)

    或者,如果您希望以确定的顺序处理事件,您可以直接在处理程序中调用 handleEvent(e) 而不是 addEvent(e)

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2011-07-15
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多