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

Go协程为何挂起?SQS消息轮询场景排查求助

问题排查:SQS轮询方法多次调用导致程序挂起

我尝试轮询SQS队列,将消息体反序列化为JSON事件负载后,仅筛选符合特定条件的消息。测试调用GetMessages(内部调用pollQueue)时,控制台输出如下:

Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Length of messages 1Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Length of messages 1Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Message number 1 sent
Length of messages 1Message number 1 sent

当多次调用该方法时,程序直接挂起,相关代码如下:

func GetMessages[T any](ctx context.Context, awsconfig aws.Config, queueName string, maxNumberOfMessages int32,
    filter func(t T) bool) ([]T, error) {
    var wg sync.WaitGroup

    ch := make(chan *types.Message, maxNumberOfMessages)

    wg.Add(1)
    go func() {
        defer wg.Done()
        pollQueue(ctx, ch, awsconfig, queueName, maxNumberOfMessages)
        wg.Wait()
    }()
    

    var payload T
    var messages []T

    for message := range ch {
        err := json.Unmarshal([]byte(*message.Body), &payload)
        if err != nil {
            logging.WithError(err).Errorf("Failed to parse JSON for message %s", *message.MessageId)
        }
        if filter(payload) {
            messages = append(messages, payload)
            if len(messages) == int(maxNumberOfMessages) {
                break
            }
        }
    }

    fmt.Printf("Length of messages %d", len(messages))
    return messages, nil
}

func pollQueue(ctx context.Context, ch chan<- *types.Message, awsconfig aws.Config, queueName string, maxNumberOfMessages int32) {
    sqsClient := sqs.NewFromConfig(awsconfig)
    queueURLOut, err := sqsClient.GetQueueUrl(ctx, &sqs.GetQueueUrlInput{
        QueueName: aws.String(queueName),
    })
    if err != nil {
        _ = fmt.Errorf("could not get queue, error is : %w", err)
    }

    for {
        msgs, err := sqsClient.ReceiveMessage(ctx, &sqs.ReceiveMessageInput{
            QueueUrl:            queueURLOut.QueueUrl,
            MaxNumberOfMessages: maxNumberOfMessages,
            WaitTimeSeconds:     1,
        })
        if err != nil {
            _ = fmt.Errorf("failed to fetch messages, error is: %w", err)
        }

        for i, message := range msgs.Messages {
            ch <- &message
            fmt.Printf("Message number %d sent\n", i+1)
        }
    }
}

问题根源分析

  1. 协程死锁与资源泄漏:GetMessages中启动的协程调用了wg.Wait(),但WaitGroup仅添加了1个计数器,导致协程无限等待自身完成,永远无法退出。同时pollQueue是无限循环无终止条件,每次调用都会启动一个无法停止的协程,多次调用后协程堆积导致程序挂起。
  2. 通道阻塞:当GetMessages收集到足够消息跳出循环后,pollQueue仍持续向通道发消息,通道缓冲填满后ch <- &message会阻塞,进一步加剧协程堆积。
  3. payload复用问题:payload在循环外声明,反序列化会覆盖原有值,若T是引用类型,切片中所有元素会指向同一个实例,导致数据错误。
  4. 错误处理无效:pollQueue仅创建错误对象但未处理,队列获取、消息接收的错误被静默忽略,无法排查潜在问题。

修复方案

func GetMessages[T any](ctx context.Context, awsconfig aws.Config, queueName string, maxNumberOfMessages int32,
    filter func(t T) bool) ([]T, error) {
    ch := make(chan *types.Message, maxNumberOfMessages)
    stopCh := make(chan struct{}) // 用于终止poll协程的信号通道

    // 启动poll协程,传入stopCh控制终止
    go func() {
        defer close(ch) // 协程退出时关闭通道,避免主循环阻塞
        pollQueue(ctx, ch, stopCh, awsconfig, queueName, maxNumberOfMessages)
    }()

    var messages []T

    for message := range ch {
        var payload T // 每次循环创建新实例,避免复用导致的数据问题
        err := json.Unmarshal([]byte(*message.Body), &payload)
        if err != nil {
            logging.WithError(err).Errorf("Failed to parse JSON for message %s", *message.MessageId)
            continue
        }
        if filter(payload) {
            messages = append(messages, payload)
            if len(messages) == int(maxNumberOfMessages) {
                close(stopCh) // 发送终止信号,让poll协程退出
                break
            }
        }
    }

    fmt.Printf("Length of messages %d\n", len(messages))
    return messages, nil
}

func pollQueue(ctx context.Context, ch chan<- *types.Message, stopCh <-chan struct{}, awsconfig aws.Config, queueName string, maxNumberOfMessages int32) {
    sqsClient := sqs.NewFromConfig(awsconfig)
    queueURLOut, err := sqsClient.GetQueueUrl(ctx, &sqs.GetQueueUrlInput{
        QueueName: aws.String(queueName),
    })
    if err != nil {
        logging.WithError(err).Error("Could not get queue URL")
        return
    }

    for {
        select {
        case <-stopCh: // 收到终止信号,退出循环
            return
        case <-ctx.Done(): // 上下文取消,退出循环
            logging.WithError(ctx.Err()).Error("Context cancelled, stopping poll")
            return
        default:
            msgs, err := sqsClient.ReceiveMessage(ctx, &sqs.ReceiveMessageInput{
                QueueUrl:            queueURLOut.QueueUrl,
                MaxNumberOfMessages: maxNumberOfMessages,
                WaitTimeSeconds:     1,
            })
            if err != nil {
                logging.WithError(err).Error("Failed to fetch messages")
                continue
            }

            for i, message := range msgs.Messages {
                select {
                case ch <- &message:
                    fmt.Printf("Message number %d sent\n", i+1)
                case <-stopCh: // 发送消息时也监听终止信号,避免阻塞
                    return
                }
            }
        }
    }
}

修复说明

  • 新增stopCh信号通道,收集到足够消息后关闭通道,通知pollQueue终止循环。
  • pollQueue用select监听stopCh和ctx.Done(),实现优雅退出,避免协程泄漏。
  • 每次循环创建新的payload实例,解决引用类型数据覆盖问题。
  • 修复错误处理逻辑,用日志记录错误而非静默忽略。
  • 移除多余的WaitGroup代码,避免死锁。
  • 协程退出时关闭通道ch,确保主循环的for range正常结束。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 00:18:11