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

Golang基于Channel与Select实现Pub/Sub并发的异常问题求助

问题:Golang并发Pub/Sub模型中已移除的订阅者仍输出空消息

我用Golang实现了一个基于并发的Pub/Sub模型,代码有时能正常执行,但偶尔会出现异常输出:

异常输出

message received on channel 2: Hello World
message received on channel 3: Hello World
message received on channel 1: Hello World
message received on channel 1: 
subscriber 1's context cancelled
message received on channel 3: Only channels 2 and 3 should print this
message received on channel 2: Only channels 2 and 3 should print this

关键问题是:订阅者1在调用RemoveSubscriber移除后,仍打印了message received on channel 1: 。订阅者1应该只接收第一条消息,之后其context被取消、goroutine退出。目前推测是订阅者的dctx被取消前,从subscriber.out通道收到了消息(通道关闭后接收会得到零值)。

预期执行结果:
通道1仅打印一次消息,之后context被取消,goroutine退出,不再接收任何消息。


main.go

package main

import (
    "fmt"
    "net/http"
)

func main() {
  var err error
  
    publisher := NewPublisher()

  publisher.AddSubscriber()
  publisher.AddSubscriber()
  publisher.AddSubscriber()

  publisher.Start()
  
  err = publisher.Publish("Hello World")
  if err != nil {
    fmt.Printf("could not publish: %v\n", err)
  }

  err = publisher.RemoveSubscriber(1)
  if err != nil {
    fmt.Printf("could not remove subscriber: %v\n", err)
  }
  
  err = publisher.Publish("Only channels 2 and 3 should print this")
  if err != nil {
    fmt.Printf("could not publish: %v\n", err)
  }

  // 保持服务运行
  http.ListenAndServe(":8080", nil)
}

publisher.go

package main

import (
    "context"
    "errors"
    "fmt"
    "sync"
)

// Publisher 将in通道收到的消息发送给所有订阅者
type Publisher struct {
    sequence    uint
    in          chan string
    subscribers map[uint]*Subscriber
    sync.RWMutex
    ctx         context.Context
    cancel      *context.CancelFunc
}

// NewPublisher 返回一个空订阅者的Publisher实例
// 必须调用Start()方法后才能开始发布消息
func NewPublisher() *Publisher {
    ctx, cancel := context.WithCancel(context.Background())
    return &Publisher{subscribers: map[uint]*Subscriber{}, ctx: ctx, cancel: &cancel}
}

// AddSubscriber 创建新订阅者并开始监听Publisher的消息
func (p *Publisher) AddSubscriber() {
    dctx, cancel := context.WithCancel(p.ctx)
    p.Lock()
    nextId := p.sequence + 1
    subscriber := NewSubscriber(nextId, dctx, &cancel)
    p.subscribers[nextId] = subscriber
    p.sequence = p.sequence + 1
    p.Unlock()

    go func() {
        for {
            select {
            case <-p.ctx.Done():
                fmt.Printf("parent context cancelled\n")
                (*subscriber.cancel)()
                return
            case <-dctx.Done():
                fmt.Printf("subscriber %d's context cancelled\n", subscriber.id)
                return
            case msg := <-subscriber.out:
                fmt.Printf("message received on channel %d: %s\n", subscriber.id, msg)
            }
        }
    }()  
}

// Publish 将消息发送给所有订阅者
func (p *Publisher) Publish(msg string) error {
    // 未启动则返回错误
    if p.in == nil {
        return errors.New("publisher not started yet")
    }
    // 无订阅者则返回错误
    p.RLock()
    if len(p.subscribers) == 0 {
        return errors.New("no subscribers to receive the message")
    }
    p.RUnlock()

    // 发送消息到in通道,加锁防止并发修改
    p.Lock()
    p.in <- msg
    p.Unlock()

    return nil
}

// Start 初始化in通道,使Publisher可以接收并发布消息
func (p *Publisher) Start() {
    in := make(chan string)
    p.in = in

    go func() {
        for {
            select {
            case <-p.ctx.Done():
                fmt.Printf("done called on publisher\n")
                return
            case msg := <-p.in:
                p.RLock()
                for _, subscriber := range p.subscribers {
                    subscriber.out <- msg
                }
                p.RUnlock()
            }
        }
    }()
}

// Stop 终止Publisher的消息监听
func (p *Publisher) Stop() {
    (*p.cancel)()
}

// RemoveSubscriber 根据ID移除订阅者,找不到则返回错误
func (p *Publisher) RemoveSubscriber(id uint) error {
    p.Lock()
    defer p.Unlock()
    subscriber, ok := p.subscribers[id]

    if !ok {
        return errors.New("could not find subscriber")
    }

    (*subscriber.cancel)()
    delete(p.subscribers, id)
    close(subscriber.out)

    return nil
}

问题分析与修复方案

问题根源

出现空消息的核心原因是:在RemoveSubscriber中,你先取消订阅者的context,紧接着关闭了subscriber.out通道。但订阅者goroutine的select分支中,context取消信号和通道读取是平等竞争的,有可能在context取消逻辑执行前,goroutine先读取到了关闭通道返回的零值(空字符串),从而打印出异常内容。

修复步骤

  1. 删除手动关闭通道的操作:不需要主动关闭subscriber.out,订阅者goroutine会在context取消后自动退出,通道会被Go的垃圾回收机制处理,避免了读取零值的问题。
  2. 增强订阅者goroutine的可靠性:在读取消息后,额外检查context状态,确保即使消息在context取消前到达,也不会被处理。

修改后的关键代码

修改RemoveSubscriber方法

移除close(subscriber.out)语句:

// RemoveSubscriber 根据ID移除订阅者,找不到则返回错误
func (p *Publisher) RemoveSubscriber(id uint) error {
    p.Lock()
    defer p.Unlock()
    subscriber, ok := p.subscribers[id]

    if !ok {
        return errors.New("could not find subscriber")
    }

    (*subscriber.cancel)()
    delete(p.subscribers, id)
    // 移除关闭通道的操作,避免读取零值
    // close(subscriber.out)

    return nil
}

优化订阅者goroutine逻辑

在处理消息前再次检查context状态,确保不会处理取消后的消息:

go func() {
    for {
        select {
        case <-p.ctx.Done():
            fmt.Printf("parent context cancelled\n")
            (*subscriber.cancel)()
            return
        case <-dctx.Done():
            fmt.Printf("subscriber %d's context cancelled\n", subscriber.id)
            return
        case msg := <-subscriber.out:
            // 再次检查context是否已取消,避免处理延迟到达的消息
            select {
            case <-dctx.Done():
                return
            default:
                fmt.Printf("message received on channel %d: %s\n", subscriber.id, msg)
            }
        }
    }
}()

修复原理

  • 不再关闭订阅者通道,彻底避免了读取通道关闭后零值的情况。
  • 订阅者goroutine会在context取消信号触发后立即退出,不会再处理任何后续消息。
  • 额外的context检查确保即使消息在取消信号前已经到达通道,也会直接退出,不执行打印逻辑。

内容的提问来源于stack exchange,提问作者Zaid Sheikh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 11:01:35