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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 11:49:59