NET 7中Amazon SQS千级以上队列消费遇套接字类异常求助
问题描述
我是Amazon SQS新手,基于.NET 7开发,需要实现同时读取千余个动态创建队列的功能。我编写了SQSQueueConsumer类,通过ReceiveMessageAsync方法读取队列并触发消息接收事件,启动多任务进行队列消费。
当消费队列数量达到1500及以上时,Amazon对象会抛出System.Net.Http.HttpRequestException、System.Net.Sockets.SocketException等异常,且队列消费速度极慢;但队列数量在1000左右时运行完全正常。
我怀疑该问题与HTTP Handler套接字耗尽有关,尝试自定义SqsHttpClientFactory结合IHttpClientFactory来处理,但仍出现相同异常。
相关代码如下:
public class SQSQueueConsumer : IQueueConsumer { public const int MIN_TIMEOUT = 0; public const int MAX_TIMEOUT = 20; private const string QUEUE_NOT_EXISTS_CODE = "AWS.SimpleQueueService.NonExistentQueue"; private readonly IAmazonSQS _sqs; private int _messageTimeout = MAX_TIMEOUT; public event EventHandler<MessageReceivedEventArgs> MessageReceived; public event EventHandler<QueueErrorEventArgs> ErrorOccurred; public SQSQueueConsumer(IAmazonSQS sqs) { ArgumentNullException.ThrowIfNull(sqs); _sqs = sqs; } public void SetMessageTimeout(int timeout) { if(timeout < MIN_TIMEOUT || timeout > MAX_TIMEOUT) { throw new ArgumentException($"Timeout should be between {MIN_TIMEOUT} and {MAX_TIMEOUT}."); } _messageTimeout = timeout; } public async Task ConsumeAsync(string queueUrl, CancellationToken ct) { var receivedRequest = new ReceiveMessageRequest() { QueueUrl = queueUrl, MaxNumberOfMessages = 10, WaitTimeSeconds = _messageTimeout }; while(!ct.IsCancellationRequested) { try { var messageResponse = await _sqs.ReceiveMessageAsync(receivedRequest, ct); if (messageResponse.HttpStatusCode != HttpStatusCode.OK) { OnError(new QueueErrorEventArgs(queueUrl, "HTTP Error.")); continue; } foreach (var message in messageResponse.Messages) { OnMessageReceived(new MessageReceivedEventArgs(queueUrl, message.Body)); await _sqs.DeleteMessageAsync(queueUrl, message.ReceiptHandle, ct); } } catch(AmazonSQSException ex) { if (ex.ErrorCode == QUEUE_NOT_EXISTS_CODE) { throw new Exceptions.QueueDoesNotExistException(queueUrl); } } } } protected virtual void OnMessageReceived(MessageReceivedEventArgs e) => MessageReceived?.Invoke(this, e); protected virtual void OnError(QueueErrorEventArgs e) => ErrorOccurred?.Invoke(this, e); }
自定义HTTP工厂代码:
public class SqsHttpClientFactory : Amazon.Runtime.HttpClientFactory { private readonly IHttpClientFactory _httpClientFactory; public SqsHttpClientFactory(IHttpClientFactory httpClientFactory) { _httpClientFactory = httpClientFactory; } public override HttpClient CreateHttpClient(IClientConfig clientConfig) { return _httpClientFactory.CreateClient(); } }
依赖注入配置代码:
services.AddSingleton<IAmazonSQS>(s => { var settings = s.GetRequiredService<IOptions<SQSSettings>>().Value; AWSOptions awsOptions = new AWSOptions { Credentials = new Amazon.Runtime.BasicAWSCredentials(settings.AccessKeyId, settings.SecretAccessKey), Region = RegionEndpoint.GetBySystemName(settings.Region), }; var httpFactory = s.GetRequiredService<IHttpClientFactory>(); var config = new AmazonSQSConfig() { RegionEndpoint = awsOptions.Region, HttpClientFactory = new SqsHttpClientFactory(httpFactory), }; var client = new AmazonSQSClient(awsOptions.Credentials, config); return client; });
解决方案
核心原因
1500+队列的长轮询(WaitTimeSeconds=20)会占用大量HTTP连接,默认的连接池配置无法承载,导致套接字耗尽;自定义工厂仅复用了客户端实例,但未针对长连接场景调整连接池参数。
1. 调整HTTP连接池配置
针对SQS场景自定义HttpClient连接池,扩容连接数并优化闲置连接回收:
services.AddHttpClient("SQSClient") .ConfigurePrimaryHttpMessageHandler(() => new SocketsHttpHandler { // 适配1500+队列的长轮询需求 MaxConnectionsPerServer = 2000, // 缩短闲置连接存活时间,释放套接字 PooledConnectionIdleTimeout = TimeSpan.FromMinutes(5), // 启用HTTP/2提升连接复用效率 EnableMultipleHttp2Connections = true });
修改自定义工厂使用命名客户端:
public override HttpClient CreateHttpClient(IClientConfig clientConfig) { return _httpClientFactory.CreateClient("SQSClient"); }
2. 优化AWS SDK配置
禁用SDK自带连接池,完全依赖IHttpClientFactory的配置,并调整超时适配长轮询:
services.AddSingleton<IAmazonSQS>(s => { var settings = s.GetRequiredService<IOptions<SQSSettings>>().Value; var httpFactory = s.GetRequiredService<IHttpClientFactory>(); var config = new AmazonSQSConfig() { RegionEndpoint = RegionEndpoint.GetBySystemName(settings.Region), HttpClientFactory = new SqsHttpClientFactory(httpFactory), UseSdkHttpClientFactory = false, Timeout = TimeSpan.FromSeconds(30), UseHttp2 = true }; return new AmazonSQSClient( new BasicAWSCredentials(settings.AccessKeyId, settings.SecretAccessKey), config); });
3. 优化消费逻辑
- 批量删除消息:减少HTTP请求次数,降低连接压力
// 替换ConsumeAsync中的消息处理逻辑 if (messageResponse.Messages.Any()) { var deleteRequests = messageResponse.Messages.Select(m => new DeleteMessageBatchRequestEntry { Id = Guid.NewGuid().ToString(), ReceiptHandle = m.ReceiptHandle }).ToList(); await _sqs.DeleteMessageBatchAsync(new DeleteMessageBatchRequest { QueueUrl = queueUrl, Entries = deleteRequests }, ct); foreach (var message in messageResponse.Messages) { OnMessageReceived(new MessageReceivedEventArgs(queueUrl, message.Body)); } }
- 添加异常重试:针对套接字/HTTP异常添加指数退避重试,避免单个队列异常扩散
// 在ConsumeAsync的try块中添加重试逻辑 var retryCount = 0; const int maxRetries = 3; do { try { var messageResponse = await _sqs.ReceiveMessageAsync(receivedRequest, ct); // 处理消息逻辑... retryCount = maxRetries; // 成功后退出重试 } catch (HttpRequestException ex) when (retryCount < maxRetries) { retryCount++; await Task.Delay(TimeSpan.FromSeconds(Math.Pow(2, retryCount)), ct); } catch (SocketException ex) when (retryCount < maxRetries) { retryCount++; await Task.Delay(TimeSpan.FromSeconds(Math.Pow(2, retryCount)), ct); } } while (retryCount < maxRetries && !ct.IsCancellationRequested);
内容的提问来源于stack exchange,提问作者Matheus Constantino
相关产品推荐
相关产品推荐

