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

使用Channel实现Go观察者模式时遇所有goroutine休眠致命错误求助

用Go Channel实现观察者模式的死锁问题修复

问题根源

你的代码触发死锁的核心原因有两个:

  1. 值传递导致Channel失效:Register方法接收的是Observer值类型,修改的只是副本的ListenChannel,原对象的channel始终是nil;Subscribe用值接收者,同样操作的是副本,goroutine里读取nil channel会永久阻塞。
  2. 同步发送的阻塞风险:无缓冲channel发送时必须有接收goroutine同步等待,哪怕修复了值传递问题,一旦某个观察者的goroutine没准备好,Publish就会卡住。

修复后的完整代码

observer/subject.go

package observer

type ISubject interface {
    Register(observer *Observer) error
    CancelRegister(observer *Observer) error
    Publish(msg string) error
}

type Subject struct {
    observers []*Observer // 直接存观察者指针,不用单独存channel
}

func (s *Subject) Register(observer *Observer) error {
    observer.ListenChannel = make(chan string)
    s.observers = append(s.observers, observer)
    return nil
}

func (s *Subject) CancelRegister(observer *Observer) error {
    for idx, obs := range s.observers {
        if obs == observer {
            s.observers = append(s.observers[:idx], s.observers[idx+1:]...)
            close(observer.ListenChannel) // 注销时关闭channel,让监听goroutine退出
            break
        }
    }
    return nil
}

func (s *Subject) Publish(msg string) error {
    // 每个消息发送都开goroutine,避免单个阻塞影响全局
    for _, obs := range s.observers {
        go func(ch chan string) {
            ch <- msg
        }(obs.ListenChannel)
    }
    return nil
}

observer/observer.go

package observer

import "fmt"

type IObserver interface {
    Subscribe()
}

type Observer struct {
    ListenChannel chan string `json:"listenChannel"`
}

// 改用指针接收者,确保操作原对象的channel
func (observer *Observer) Subscribe() {
    // 用for range持续监听,直到channel关闭
    for msg := range observer.ListenChannel {
        fmt.Println("收到消息:", msg)
    }
    fmt.Println("该观察者已停止监听")
}

测试代码

import (
    "testing"
    "time"
)

func Test_Sub(t *testing.T) {
    s1 := &Subject{}
    o1 := &Observer{}
    
    if err := s1.Register(o1); err != nil {
        t.Fatal(err)
    }
    
    // 启动监听goroutine
    go o1.Subscribe()
    
    fmt.Println("发布消息中...")
    if err := s1.Publish("hahaha"); err != nil {
        t.Fatal(err)
    }
    
    // 测试环境等待消息处理完成,生产环境可改用WaitGroup等同步方式
    time.Sleep(100 * time.Millisecond)
    
    // 注销观察者
    if err := s1.CancelRegister(o1); err != nil {
        t.Fatal(err)
    }
}

关键调整说明

  • 全指针操作:所有涉及Observer的方法和参数都改用指针,保证修改和访问的是同一个对象的channel,彻底解决值传递导致的副本问题。
  • 异步消息发送:Publish时给每个发送操作单独开goroutine,避免某个观察者处理慢阻塞整个发布流程。
  • channel生命周期管理:注销观察者时关闭对应的channel,让监听goroutine正常退出,防止内存泄漏。
  • 持续监听:Subscribe改用for range循环监听,支持接收多条消息,直到channel被关闭。

内容的提问来源于stack exchange,提问作者y Zz

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 18:36:25