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

Go中用Context超时处理AMQP消息时出现send on closed channel恐慌

解决Go中Context超时处理AMQP消息时的panic: send on closed channel错误

问题场景

在Go应用中使用Context超时处理AMQP消息消费时,遇到了panic: send on closed channel错误,简化复现代码如下:

业务逻辑代码

func (i interactor) ScanExtension() (*AmqpMailSummaryModel, error) {
    resultChan := make(chan AmqpMailSummaryModel)
    errorChan := make(chan error)

    ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
    defer cancel()

    go func() {
        defer close(resultChan)
        defer close(errorChan)

        err := i.services.AmqpConsumer.StartConsuming(
            func(body []byte) {
                // body 转换为 summary model 的逻辑
                resultChan <- summary
            })
        if err != nil {
            errorChan <- fmt.Errorf("error starting consumer: %v", err)
        }
    }()

    // messageBody 处理逻辑...
    err = i.services.AmqpProducer.Publish(messageBody)
    if err != nil {
        return nil, fmt.Errorf("failed to publish a message: %v", err)
    }

    select {
    case result := <-resultChan:
        return &result, nil
    case err := <-errorChan:
        return nil, err
    case <-ctx.Done():
        return nil, nil
    }
}

AMQP消费者代码

func (c *RabbitMQConsumer) StartConsuming(handleMessage func(body []byte)) error {
    msgs, err := c.Channel.Consume(
        c.Queue,
        "",
        true,  // auto-ack
        false, // exclusive
        false, // no-local
        false, // no-wait
        nil,
    )
    if err != nil {
        return err
    }

    go func() {
        for msg := range msgs {
            handleMessage(msg.Body)
        }
    }()

    return nil
}

错误堆栈

panic: send on closed channel

goroutine 45 [running]:
tScanExtension.func1.1({0xc000782000, 0x56425a, 0x56425a})
        interactor.go:306 +0x32b
(*RabbitMQConsumer).StartConsuming.func1()
       Consumer.go:52 +0x1a8
created by (*RabbitMQConsumer).StartConsuming in goroutine 16
       Consumer.go:48 +0x15a
exit status 2

错误原因分析

当Context超时触发ctx.Done()时,主函数直接return,同时defer cancel()执行。此时启动的goroutine会执行defer close(resultChan)和defer close(errorChan)关闭结果通道,但AMQP消费者的goroutine是独立运行的,不会感知Context的取消,后续收到消息后仍然会调用handleMessage往已关闭的resultChan发送数据,从而触发send on closed channel的panic。

解决方案

1. 让AMQP消费者响应Context取消信号

修改StartConsuming方法,传入Context,让内部的消费goroutine监听Context的取消信号,及时停止消费循环:

func (c *RabbitMQConsumer) StartConsuming(ctx context.Context, handleMessage func(body []byte)) error {
    msgs, err := c.Channel.Consume(
        c.Queue,
        "",
        true,  // auto-ack
        false, // exclusive
        false, // no-local
        false, // no-wait
        nil,
    )
    if err != nil {
        return err
    }

    go func() {
        for {
            select {
            case msg, ok := <-msgs:
                if !ok {
                    return
                }
                handleMessage(msg.Body)
            case <-ctx.Done():
                // Context取消,退出消费循环
                return
            }
        }
    }()

    return nil
}

2. 在消息处理函数中使用非阻塞发送

修改ScanExtension中的消息处理闭包,用select实现非阻塞发送,避免通道关闭后的panic:

go func() {
    defer close(resultChan)
    defer close(errorChan)

    err := i.services.AmqpConsumer.StartConsuming(ctx, // 传入Context
        func(body []byte) {
            // body 转换为 summary model 的逻辑
            select {
            case resultChan <- summary:
                // 发送成功
            case <-ctx.Done():
                // Context已取消,放弃发送
            }
        })
    if err != nil {
        select {
        case errorChan <- fmt.Errorf("error starting consumer: %v", err):
        case <-ctx.Done():
        }
    }
}()

3. 同步通道关闭与goroutine生命周期

通过让消费者监听Context,确保在Context取消后停止消息处理,避免后续的通道发送操作,配合defer的通道关闭逻辑,保证时序合理性。

Go中结合Context超时处理通道的注意事项

  • 通道关闭后禁止发送数据:通道一旦关闭,再往里面发送数据会直接触发panic,仅能执行接收操作(接收会返回零值和false)。
  • 绑定goroutine生命周期到Context:所有长时间运行的goroutine都应该监听ctx.Done()信号,在Context取消时及时退出,避免僵尸goroutine和无效操作。
  • 优先使用非阻塞发送或带超时的select:当不确定通道是否可用时,用select结合通道发送和ctx.Done()/超时分支,避免阻塞或panic。
  • 避免重复关闭通道:同一个通道只能被关闭一次,多次关闭会触发panic,要确保关闭操作的唯一性(比如只在一个goroutine中执行关闭)。
  • 谨慎使用无缓冲通道:无缓冲通道的发送会阻塞直到有接收方,在Context超时场景下,若接收方已经退出,发送操作会永久阻塞,必须结合select处理。

内容的提问来源于stack exchange,提问作者Talha K.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 04:10:00