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

如何在C#中实现第三方API长轮询?求100个SessionId请求示例

在C#中实现高效消息轮询的最佳方案

针对你的需求——每2秒触发一次轮询,且每次需并发处理100个不同Session ID的请求并将结果返回给对应用户,我推荐使用以下组合方案:

  1. PeriodicTimer(.NET 6+):替代传统的System.Timers.Timer,它支持异步等待,避免线程池阻塞,更符合现代C#异步编程范式。
  2. 异步并发请求:用Task.WhenAll配合复用的HttpClient(通过IHttpClientFactory),高效处理100个并行请求,确保能在1秒内完成(只要第三方API响应足够快)。
  3. 异常隔离与结果分发:每个请求单独捕获异常,避免单个请求失败影响整体轮询;通过回调、事件或依赖注入的用户服务,将结果推送给对应用户。

示例代码实现

1. 定义核心模型与服务接口

首先,我们需要一个模型来关联Session ID和用户的结果回调,以及一个处理轮询的服务:

// 用于关联Session ID和用户结果回调的模型
public class SessionPollingContext
{
    public string SessionId { get; set; } = string.Empty;
    // 可以替换为你的用户标识或回调委托
    public Action<string, string> OnResultReceived { get; set; } = (_, _) => { };
}

// 轮询服务接口
public interface ISessionPollingService
{
    void StartPolling();
    void StopPolling();
    void AddSession(SessionPollingContext sessionContext);
    void RemoveSession(string sessionId);
}

2. 实现轮询服务

使用PeriodicTimer定时触发,结合IHttpClientFactory复用HttpClient,并发处理请求:

public class SessionPollingService : ISessionPollingService, IDisposable
{
    private readonly IHttpClientFactory _httpClientFactory;
    private PeriodicTimer? _pollingTimer;
    private readonly List<SessionPollingContext> _sessions = new();
    private readonly object _sessionLock = new();
    private bool _isDisposed;

    public SessionPollingService(IHttpClientFactory httpClientFactory)
    {
        _httpClientFactory = httpClientFactory;
    }

    public void StartPolling()
    {
        if (_pollingTimer != null) return;

        // 每2秒触发一次轮询
        _pollingTimer = new PeriodicTimer(TimeSpan.FromSeconds(2));
        _ = RunPollingLoopAsync();
    }

    public void StopPolling()
    {
        _pollingTimer?.Dispose();
        _pollingTimer = null;
    }

    public void AddSession(SessionPollingContext sessionContext)
    {
        lock (_sessionLock)
        {
            if (!_sessions.Any(s => s.SessionId == sessionContext.SessionId))
            {
                _sessions.Add(sessionContext);
            }
        }
    }

    public void RemoveSession(string sessionId)
    {
        lock (_sessionLock)
        {
            _sessions.RemoveAll(s => s.SessionId == sessionId);
        }
    }

    private async Task RunPollingLoopAsync()
    {
        if (_pollingTimer == null) return;

        try
        {
            while (await _pollingTimer.WaitForNextTickAsync())
            {
                await ProcessAllSessionsAsync();
            }
        }
        catch (OperationCanceledException)
        {
            // 定时器被停止时会抛出此异常,无需处理
        }
    }

    private async Task ProcessAllSessionsAsync()
    {
        List<SessionPollingContext> currentSessions;
        lock (_sessionLock)
        {
            currentSessions = new List<SessionPollingContext>(_sessions);
        }

        if (!currentSessions.Any()) return;

        // 并发处理所有Session的请求
        var tasks = currentSessions.Select(session => ProcessSingleSessionAsync(session));
        await Task.WhenAll(tasks);
    }

    private async Task ProcessSingleSessionAsync(SessionPollingContext session)
    {
        var httpClient = _httpClientFactory.CreateClient("ThirdPartyApiClient");
        try
        {
            // 构造请求(根据第三方API的实际要求调整)
            var response = await httpClient.GetAsync($"/api/data?sessionId={session.SessionId}");
            response.EnsureSuccessStatusCode();

            var result = await response.Content.ReadAsStringAsync();
            // 将结果返回给对应用户(通过回调/事件/用户服务)
            session.OnResultReceived(session.SessionId, result);
        }
        catch (Exception ex)
        {
            // 单独处理每个Session的请求异常,避免影响其他请求
            Console.WriteLine($"处理Session {session.SessionId} 时出错: {ex.Message}");
            // 可以在这里通知用户请求失败
            session.OnResultReceived(session.SessionId, $"请求失败: {ex.Message}");
        }
    }

    public void Dispose()
    {
        Dispose(true);
        GC.SuppressFinalize(this);
    }

    protected virtual void Dispose(bool disposing)
    {
        if (_isDisposed) return;

        if (disposing)
        {
            _pollingTimer?.Dispose();
        }

        _isDisposed = true;
    }
}

3. 配置与使用

在你的依赖注入容器中注册服务和HttpClient:

// 在Program.cs中
builder.Services.AddHttpClient("ThirdPartyApiClient", client =>
{
    client.BaseAddress = new Uri("https://your-third-party-api-url.com/");
    client.Timeout = TimeSpan.FromSeconds(1); // 设置超时,确保1秒内完成请求
});

builder.Services.AddSingleton<ISessionPollingService, SessionPollingService>();

然后在业务代码中使用:

// 示例:添加一个Session并启动轮询
var pollingService = serviceProvider.GetRequiredService<ISessionPollingService>();
pollingService.AddSession(new SessionPollingContext
{
    SessionId = "user-session-001",
    OnResultReceived = (sessionId, result) =>
    {
        // 这里处理返回给用户的逻辑,比如推送SignalR消息、更新数据库等
        Console.WriteLine($"Session {sessionId} 收到结果: {result}");
    }
});

pollingService.StartPolling();

关键注意事项

  1. HttpClient复用:绝对不要在每次请求时创建新的HttpClient,否则会导致套接字耗尽。使用IHttpClientFactory是最佳实践。
  2. 并发控制:如果第三方API有并发请求限制,可以添加SemaphoreSlim来控制并发数,比如:
    private readonly SemaphoreSlim _semaphore = new SemaphoreSlim(50); // 限制同时50个请求
    // 在ProcessSingleSessionAsync中
    await _semaphore.WaitAsync();
    try
    {
        // 执行请求逻辑
    }
    finally
    {
        _semaphore.Release();
    }
    
  3. 结果分发:示例中用了回调委托,实际项目中可以结合SignalR、WebSocket或消息队列来实时推送给用户。
  4. 定时器精度:PeriodicTimer的精度依赖于系统时钟,若需要更高精度,可以考虑使用System.Threading.Timer,但PeriodicTimer的异步模型更易用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 08:52:25