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

如何实现SFTP异步连接上传,且能通知客户端上传结果?

解决方案

1. 复用SFTP连接(连接池)

每次新建SFTP连接的握手、认证是耗时核心,复用已建立的连接能大幅压缩连接耗时。可以实现一个线程安全的连接池,维护一批存活的SftpClient实例:

public class SftpConnectionPool
{
    private readonly ConcurrentQueue<SftpClient> _pool = new ConcurrentQueue<SftpClient>();
    private readonly Func<SftpClient> _clientFactory;
    private readonly int _maxPoolSize;

    public SftpConnectionPool(Func<SftpClient> clientFactory, int maxPoolSize = 10)
    {
        _clientFactory = clientFactory;
        _maxPoolSize = maxPoolSize;
    }

    public async Task<SftpClient> GetClientAsync()
    {
        if (_pool.TryDequeue(out var client))
        {
            if (client.IsConnected)
                return client;
            client.Dispose();
        }

        if (_pool.Count >= _maxPoolSize)
            throw new InvalidOperationException("SFTP连接池已满");

        client = _clientFactory();
        await client.ConnectAsync();
        return client;
    }

    public void ReturnClient(SftpClient client)
    {
        if (client.IsConnected && _pool.Count < _maxPoolSize)
            _pool.Enqueue(client);
        else
            client.Dispose();
    }
}

使用示例:

// 全局初始化单例连接池
var sftpPool = new SftpConnectionPool(() => new SftpClient("host", "username", "password"));

// 异步上传方法
public async Task UploadAsync(string localPath, string remotePath)
{
    SftpClient client = null;
    try
    {
        client = await sftpPool.GetClientAsync();
        using var fileStream = File.OpenRead(localPath);
        await client.UploadFileAsync(fileStream, remotePath);
    }
    catch (Exception ex)
    {
        _logger.LogError(ex, "SFTP上传失败");
        throw; // 抛出异常让API返回错误给客户端
    }
    finally
    {
        if (client != null)
            sftpPool.ReturnClient(client);
    }
}

2. 使用原生异步SFTP方法

不要用Task.Run包装同步代码,改用支持原生异步的SFTP库(比如SSH.NET最新版本),这样API线程不会被阻塞,既能提升吞吐量,又能及时将错误反馈给客户端:

public async Task UploadAsync(string localPath, string remotePath)
{
    using var sftp = new SftpClient("host", "username", "password");
    try
    {
        await sftp.ConnectAsync(); // 原生异步连接,不阻塞线程
        using var fileStream = File.OpenRead(localPath);
        await sftp.UploadFileAsync(fileStream, remotePath);
        await sftp.DisconnectAsync();
    }
    catch (Exception ex)
    {
        _logger.LogError(ex, "Error Upload SFTPHandler");
        throw; // 抛出异常,API返回5xx错误给客户端
    }
}

3. 异步上传+客户端状态查询/回调

如果必须让客户端立即收到响应,同时要通知上传结果,可采用以下流程:

  • API接收请求后,生成唯一任务ID,将任务信息(本地文件路径、远程路径、状态)存入数据库/缓存
  • 立即返回客户端202 Accepted,附带任务ID
  • 后台用队列(比如Hangfire、RabbitMQ)异步执行SFTP上传,更新任务状态
  • 客户端通过任务ID轮询API查询结果,或提前提供回调URL,上传完成后API主动通知

示例伪代码:

// API接收请求接口
[HttpPost]
public IActionResult UploadMessage([FromBody] MessageDto message)
{
    var jsonPath = SaveMessageToLocalJson(message);
    var taskId = Guid.NewGuid().ToString();
    _taskRepository.Add(new UploadTask
    {
        Id = taskId,
        LocalPath = jsonPath,
        RemotePath = $"/remote/dir/{taskId}.json",
        Status = UploadStatus.Pending
    });

    // 加入后台执行队列
    _backgroundQueue.Enqueue(async () => await ProcessUploadTask(taskId));

    return Accepted(new { TaskId = taskId });
}

// 后台上传处理方法
private async Task ProcessUploadTask(string taskId)
{
    var task = _taskRepository.GetById(taskId);
    try
    {
        using var sftp = new SftpClient("host", "username", "password");
        await sftp.ConnectAsync();
        using var fileStream = File.OpenRead(task.LocalPath);
        await sftp.UploadFileAsync(fileStream, task.RemotePath);
        task.Status = UploadStatus.Success;
    }
    catch (Exception ex)
    {
        _logger.LogError(ex, "任务{TaskId}上传失败", taskId);
        task.Status = UploadStatus.Failed;
        task.ErrorMessage = ex.Message;
    }
    finally
    {
        _taskRepository.Update(task);
        // 若客户端提供回调URL,主动通知结果
        if (!string.IsNullOrEmpty(task.CallbackUrl))
            await _httpClient.PostAsync(task.CallbackUrl, JsonContent.Create(task));
    }
}

// 客户端查询状态接口
[HttpGet("upload-status/{taskId}")]
public IActionResult GetUploadStatus(string taskId)
{
    var task = _taskRepository.GetById(taskId);
    if (task == null) return NotFound();
    return Ok(new { task.Status, task.ErrorMessage });
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 09:55:03