基于Twisted与txrequests实现异步文件下载的问题咨询
看起来你在Twisted应用里用txrequests做流式下载时遇到了反应器阻塞、没法更新GUI动画的问题对吧?我来逐个解答你的疑问:
问题1:让阻塞生成器在Twisted中让出控制权
你的核心问题在于r.iter_content()是同步阻塞的——当它等待网络数据时,整个生成器会卡在那里,不会让出控制权给反应器,哪怕你加了yield None也没用,因为yield只有在iter_content()返回一个chunk之后才会执行。
twisted.flow确实是早年用来处理这类迭代式异步逻辑的模块,但它早就被废弃了,现在Twisted推荐用Deferred结合线程或者原生异步IO原语来处理。针对你的场景,有两个可行方案:
方案A:用deferToThread包装阻塞的迭代
把iter_content()的迭代放到线程里,这样每次获取chunk的操作不会阻塞反应器:
from twisted.internet.threads import deferToThread def _stream_download(r, f): def read_one_chunk(): try: return next(r.iter_content(chunk_size=128)) except StopIteration: return None def write_and_continue(chunk): if chunk is None: return f.write(chunk) # 递归调用,读下一个chunk前先回到反应器 return deferToThread(read_one_chunk).addCallback(write_and_continue) return deferToThread(read_one_chunk).addCallback(write_and_continue)
这样每次读取一个chunk都会回到反应器,GUI动画就能正常更新了。
方案B:放弃iter_content,改用异步响应API
txrequests底层还是依赖requests的同步IO,所以更彻底的方式是直接用Twisted原生的异步响应处理(比如后面提到的treq),但如果一定要用txrequests,方案A是最直接的修复方式。
问题2:Twisted下的全功能异步文件下载方案
当然有!treq是Twisted官方推荐的、模仿requests API的异步HTTP客户端,它完全基于Twisted的异步IO,支持流式响应、代理、重试,而且和Twisted生态完美兼容。
用treq实现流式下载的示例:
import treq from twisted.internet import defer def http_download(url, dst, callback, errback=None): def handle_response(response): if response.code != 200: return defer.fail(Exception(f"HTTP Error: {response.code}")) filehandle = open(dst, 'wb') def write_chunk(chunk): if chunk: filehandle.write(chunk) # 这里可以插入GUI进度更新逻辑 return response.read() else: filehandle.close() return callback(url, dst) return response.read().addCallback(write_chunk) d = treq.get(url, stream=True) d.addCallback(handle_response) if errback: d.addErrback(errback) return d
treq支持stream=True,response.read()会异步返回下一个chunk,完全不会阻塞反应器。另外,你可以结合Twisted的twisted.web.client配置代理,或者用txretry实现重试逻辑,这些都是成熟的Twisted生态组件。
问题3:纯Twisted实现HTTP异步文件下载的思路
如果完全不用requests/treq,用Twisted原生的Agent组件实现,核心思路是:
- 创建
Agent实例(可配置代理、SSL等) - 发送HTTP GET请求,获取
Response对象 - 实现一个
Protocol子类,用来处理流式响应数据:- 当
dataReceived被调用时,把数据写入文件 - 当
connectionLost被调用时,完成下载并触发回调
- 当
- 用
response.deliverBody(protocol_instance)把响应数据交给自定义Protocol处理
示例代码:
from twisted.web.client import Agent from twisted.web.http_headers import Headers from twisted.internet import reactor, defer from twisted.protocols.basic import FileSender class DownloadProtocol(FileSender): def __init__(self, filehandle, callback): self.filehandle = filehandle self.callback = callback def connectionLost(self, reason): self.filehandle.close() self.callback(reason) def http_download(url, dst, callback, errback=None): agent = Agent(reactor) def handle_response(response): if response.code != 200: return defer.fail(Exception(f"HTTP Error: {response.code}")) filehandle = open(dst, 'wb') protocol = DownloadProtocol(filehandle, callback) response.deliverBody(protocol) return protocol.whenDone() d = agent.request( b'GET', url.encode('utf-8'), Headers({'User-Agent': ['Twisted Downloader']}), None # GET请求不需要请求体 ) d.addCallback(handle_response) if errback: d.addErrback(errback) return d
这个方案完全异步,所有IO操作都由Twisted反应器调度,不会阻塞GUI或其他任务。你可以在DownloadProtocol的dataReceived方法里添加GUI进度更新逻辑,每次收到数据块就触发一次更新。
对你现有代码的修改建议
你的现有代码中,cooperate(_stream_download(...))无法解决阻塞问题,因为iter_content()是同步的。你可以把_stream_download替换成上面方案A中的线程包装版本,或者直接切换到treq方案,这样就能让反应器正常处理GUI动画了。
内容的提问来源于stack exchange,提问作者Chintalagiri Shashank

