SFTP大文件上传时如何并行执行CPU密集型与I/O密集型操作
优化实现方案
针对提出的两个并行优化点,可按以下方式落地,全程不会重复读取NFS上的源文件,不破坏原有多连接并发上传的逻辑。
两个open操作的并行
本地文件打开是磁盘/NFS阻塞I/O,远端SFTP文件打开是网络阻塞I/O,二者无依赖关系,可直接并发执行。
注意原生阻塞的文件操作不能直接跑在asyncio事件循环线程里,否则会卡住所有并发任务,需要丢到默认线程池执行;如果SFTP客户端提供的是原生协程版的open方法,可直接参与并发调度:
import asyncio import os import hashlib # 并发打开本地、远端文件 inp, out = await asyncio.gather( asyncio.to_thread(open, fName, "rb"), SFTP.open(os.path.split(fName)[-1], "w") ) bsize = os.stat(inp.fileno()).st_blksize
文件关闭逻辑同理,也可以用同样的gather方式并发执行,减少等待时间。
digest.update()与out.write()的并行
hashlib.sha256().update()的C底层实现计算时会主动释放GIL,属于可在线程池中安全并发的CPU密集操作;SFTP的write()是网络I/O密集操作,二者操作的是已经读入内存的同一份文件块,不存在数据竞争,完全可以并行。
实现时需要保证文件块的写入顺序不混乱,因此每读入一个块后,同时启动hash计算和块写入两个任务,等两个任务都完成后再读取下一个块即可,此时CPU计算hash的时间和网络传输块的时间完全重叠,不会串行等待。
另外注意NFS的read操作也是阻塞I/O,同样需要丢到线程池执行,避免卡事件循环。
完整优化后代码
async def upload(fName, SFTP): # 并发打开两端文件 inp, out = await asyncio.gather( asyncio.to_thread(open, fName, "rb"), SFTP.open(os.path.split(fName)[-1], "w") ) digest = hashlib.sha256() bsize = os.stat(inp.fileno()).st_blksize try: while True: # 读文件块(NFS I/O丢线程池) buf = await asyncio.to_thread(inp.read, bsize) if not buf: break # 并行执行hash计算和远端写入 await asyncio.gather( asyncio.to_thread(digest.update, buf), out.write(buf) ) finally: # 并发关闭文件 await asyncio.gather( asyncio.to_thread(inp.close), out.close() ) print('SHA256 (%s) = %s' % (fName, digest.hexdigest()))
补充说明
- 如果你使用的SFTP库的
write/open方法是阻塞式实现而非原生协程,需要同样用asyncio.to_thread()包装后再参与并发,否则依然会串行阻塞。 - 块大小不必完全拘泥于文件系统的
st_blksize,针对NFS+广域网上传的场景,把块大小调整到1MB~4MB通常能获得更高的吞吐量,减少系统调用和网络包调度的开销。 - 该实现全程仅对源文件的每个块读取一次,完全满足不重复读取低速NFS存储的要求,同时CPU计算和网络I/O的重叠执行能大幅提升单连接的上传效率,也不会影响多upload任务的并发性能。
内容的提问来源于stack exchange,提问作者Mikhail T.
相关产品推荐
相关产品推荐

