如何在Python sshfs(fsspec)中实现连接池并自动重连?
问题
我通过sshfs从远程SSH存储拉取视频文件,接口代码如下:
@app.get("/video/{filename}") async def video_endpoint( filename, range: str = Header(None), db=Depends(get_db) ): # pylint: disable=redefined-builtin """ Endpoint for video streaming Accepts the UUID and an arbitrary extension """ # Requesting video with uuid (uuid, extension) = filename.split(".") # pylint: disable=unused-variable # Try to get the file info from database file = File.get_by_public_uuid(db, uuid) # Return 404 if file not found if not file: raise HTTPException(404, "File not found") # Connect with a password ssh_fs = SSHFileSystem( settings.ssh_host, username=settings.ssh_username, password=settings.ssh_password, ) start, end = range.replace("bytes=", "").split("-") start = int(start) end = int(end) if end else start + settings.chunk_size with ssh_fs.open(file.path) as video: video.seek(start) data = video.read(end - start) filesize = file.size headers = { "Content-Range": f"bytes {str(start)}-{str(end)}/{filesize}", "Accept-Ranges": "bytes", } return Response(data, status_code=206, headers=headers, media_type="video/mp4")
接口重启后能正常工作数小时,但后续调用会报错asyncssh.sftp.SFTPNoConnection: Connection not open。排查发现,尽管每次API调用都初始化SSHFileSystem,但fsspec会缓存实例并创建asyncio事件循环,远端断开连接后无法自动重连。调用ssh_fs.clear_instance_cache()能避免错误,但会导致每个请求(包括分片请求)都新建连接,效率太低。
请问如何通过连接池实现连接保持与自动重连,解决SFTPNoConnection问题?
解决方案
1. 配置fsspec连接池参数
SSHFileSystem原生支持连接池配置,通过max_connections和cache_timeout参数控制连接复用逻辑:
ssh_fs = SSHFileSystem( settings.ssh_host, username=settings.ssh_username, password=settings.ssh_password, max_connections=10, # 限制同时存在的最大连接数 cache_timeout=300 # 缓存超时时间(秒),到期自动清理失效连接 )
max_connections:避免无限制创建连接导致资源耗尽cache_timeout:让fsspec定期清理过期或已断开的连接缓存
2. 主动检查连接有效性并自动重连
在使用连接前添加有效性校验,若连接失效则清除缓存并重新创建:
from asyncssh.sftp import SFTPNoConnection async def get_valid_ssh_fs(): # 尝试获取缓存连接 ssh_fs = SSHFileSystem( settings.ssh_host, username=settings.ssh_username, password=settings.ssh_password, max_connections=10, cache_timeout=300 ) # 执行简单操作验证连接状态 try: ssh_fs.exists("/") except (SFTPNoConnection, ConnectionError): # 清除失效缓存并重建连接 ssh_fs.clear_instance_cache() ssh_fs = SSHFileSystem( settings.ssh_host, username=settings.ssh_username, password=settings.ssh_password, max_connections=10, cache_timeout=300 ) return ssh_fs
在接口中替换原有的初始化逻辑:
# 原SSHFileSystem初始化代码替换为 ssh_fs = await get_valid_ssh_fs()
3. 用依赖注入统一管理连接池
将连接池封装为FastAPI依赖项,实现单例连接池共享,同时统一处理连接失效:
async def get_ssh_fs_pool(): # 单例模式创建连接池配置的SSHFileSystem if not hasattr(get_ssh_fs_pool, "fs"): get_ssh_fs_pool.fs = SSHFileSystem( settings.ssh_host, username=settings.ssh_username, password=settings.ssh_password, max_connections=10, cache_timeout=300 ) fs = get_ssh_fs_pool.fs # 验证连接有效性 try: fs.exists("/") except (SFTPNoConnection, ConnectionError): fs.clear_instance_cache() get_ssh_fs_pool.fs = SSHFileSystem( settings.ssh_host, username=settings.ssh_username, password=settings.ssh_password, max_connections=10, cache_timeout=300 ) fs = get_ssh_fs_pool.fs yield fs # 修改接口依赖 @app.get("/video/{filename}") async def video_endpoint( filename, range: str = Header(None), db=Depends(get_db), ssh_fs=Depends(get_ssh_fs_pool) ): # 后续业务逻辑保持不变 ...
所有请求共享同一连接池,无需重复初始化,同时自动处理连接失效重连。
4. 优化分片请求的连接复用
视频流的分片请求属于同一会话,连接池会自动复用空闲连接,无需额外处理。若需更精细控制,可结合请求会话ID关联连接,但默认的连接池逻辑已能满足分片请求的复用需求。
内容的提问来源于stack exchange,提问作者Daniel
相关产品推荐
相关产品推荐

