如何实现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
相关产品推荐
相关产品推荐

