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) } } }
问题根源分析
- 协程死锁与资源泄漏:
GetMessages中启动的协程调用了wg.Wait(),但WaitGroup仅添加了1个计数器,导致协程无限等待自身完成,永远无法退出。同时pollQueue是无限循环无终止条件,每次调用都会启动一个无法停止的协程,多次调用后协程堆积导致程序挂起。 - 通道阻塞:当
GetMessages收集到足够消息跳出循环后,pollQueue仍持续向通道发消息,通道缓冲填满后ch <- &message会阻塞,进一步加剧协程堆积。 - payload复用问题:
payload在循环外声明,反序列化会覆盖原有值,若T是引用类型,切片中所有元素会指向同一个实例,导致数据错误。 - 错误处理无效:
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
相关产品推荐
相关产品推荐

