如何为Twisted中基于Deferred的HTTP客户端实现限流?
问题:Twisted Deferred HTTP客户端实现请求限流与排队
我用Twisted框架写了一个基于Deferred的HTTP客户端,用来向某网站API发送请求,简化后的代码如下:
from json import loads from core import output from twisted.python.log import msg from twisted.internet import reactor from twisted.web.client import Agent, HTTPConnectionPool, _HTTP11ClientFactory, readBody from twisted.web.http_headers import Headers from twisted.internet.ssl import ClientContextFactory class WebClientContextFactory(ClientContextFactory): def getContext(self, hostname, port): return ClientContextFactory.getContext(self) class QuietHTTP11ClientFactory(_HTTP11ClientFactory): # 关闭日志冗余输出 noisy = False class Output(output.Output): def start(self): myQuietPool = HTTPConnectionPool(reactor) myQuietPool._factory = QuietHTTP11ClientFactory self.agent = Agent( reactor, contextFactory=WebClientContextFactory(), pool=myQuietPool ) def stop(self): pass def write(self, event): messg = 'Whatever' self.send_message(messg) def send_message(self, message): headers = Headers({ b'User-Agent': [b'MyApp'] }) url = 'https://api.somesite.com/{}'.format(message) d = self.agent.request(b'GET', url.encode('utf-8'), headers, None) def cbBody(body): return processResult(body) def cbPartial(failure): failure.printTraceback() return processResult(failure.value) def cbResponse(response): if response.code in [200, 201]: return else: msg('Site response: {} {}'.format(response.code, response.phrase)) d = readBody(response) d.addCallback(cbBody) d.addErrback(cbPartial) return d def cbError(failure): failure.printTraceback() def processResult(result): j = loads(result) msg('Site response: {}'.format(j)) d.addCallback(cbResponse) d.addErrback(cbError) return d
这个客户端运行正常,但目标网站有限流机制,请求发太快会被丢弃。我需要给客户端实现限流,保证请求发送速度合理,同时对请求缓冲排队避免丢失。不需要精确限流(比如每秒不超过X次),只要每个请求后设置1秒左右的延迟就行。
但Deferred里不能用sleep(),查资料知道大概思路类似:
self.transport.pauseProducing() delay = 1 # seconds self.reactor.callLater(delay, self.transport.resumeProducing)
但这段参考代码没法直接运行(SlowDownloader需要传入reactor参数),另外找到的都是服务端限流方案,我需要的是客户端限流,不知道怎么整合出可用代码,求帮助。
解决方案
要实现客户端的请求排队与延迟,核心思路是维护一个请求队列,每次只处理队列中的一个请求,处理完成后延迟1秒再处理下一个。具体修改如下:
1. 添加队列与状态变量
在Output类的start方法中初始化请求队列和是否正在处理请求的标记:
def start(self): myQuietPool = HTTPConnectionPool(reactor) myQuietPool._factory = QuietHTTP11ClientFactory self.agent = Agent( reactor, contextFactory=WebClientContextFactory(), pool=myQuietPool ) # 初始化请求队列 self.request_queue = [] # 标记是否正在处理请求 self.is_processing = False
2. 重写write方法,将请求加入队列
原来的write直接调用send_message,现在改为把请求加入队列,然后触发队列处理:
def write(self, event): messg = 'Whatever' # 将请求加入队列 self.request_queue.append(messg) # 触发队列处理 self.process_queue()
3. 实现队列处理逻辑
添加process_queue方法,负责逐个处理队列中的请求,处理完成后延迟1秒再处理下一个:
def process_queue(self): # 如果正在处理或队列为空,直接返回 if self.is_processing or not self.request_queue: return self.is_processing = True # 取出队列第一个请求 message = self.request_queue.pop(0) # 发送请求,获取Deferred对象 d = self.send_message(message) # 请求完成后的回调:标记处理结束,延迟1秒后继续处理队列 def on_request_done(_): self.is_processing = False # 延迟1秒后处理下一个请求 reactor.callLater(1, self.process_queue) # 不管成功失败,都执行后续逻辑 d.addBoth(on_request_done)
4. 优化send_message的异常处理(可选)
给processResult添加异常捕获,避免解析响应失败时阻塞队列:
def processResult(result): try: j = loads(result) msg('Site response: {}'.format(j)) except Exception as e: msg('Failed to parse response: {}'.format(e))
完整修改后的代码
from json import loads from core import output from twisted.python.log import msg from twisted.internet import reactor from twisted.web.client import Agent, HTTPConnectionPool, _HTTP11ClientFactory, readBody from twisted.web.http_headers import Headers from twisted.internet.ssl import ClientContextFactory class WebClientContextFactory(ClientContextFactory): def getContext(self, hostname, port): return ClientContextFactory.getContext(self) class QuietHTTP11ClientFactory(_HTTP11ClientFactory): # 关闭日志冗余输出 noisy = False class Output(output.Output): def start(self): myQuietPool = HTTPConnectionPool(reactor) myQuietPool._factory = QuietHTTP11ClientFactory self.agent = Agent( reactor, contextFactory=WebClientContextFactory(), pool=myQuietPool ) # 初始化请求队列和处理状态 self.request_queue = [] self.is_processing = False def stop(self): pass def write(self, event): messg = 'Whatever' self.request_queue.append(messg) self.process_queue() def process_queue(self): if self.is_processing or not self.request_queue: return self.is_processing = True message = self.request_queue.pop(0) d = self.send_message(message) def on_request_done(_): self.is_processing = False reactor.callLater(1, self.process_queue) d.addBoth(on_request_done) def send_message(self, message): headers = Headers({ b'User-Agent': [b'MyApp'] }) url = 'https://api.somesite.com/{}'.format(message) d = self.agent.request(b'GET', url.encode('utf-8'), headers, None) def cbBody(body): return processResult(body) def cbPartial(failure): failure.printTraceback() return processResult(failure.value) def cbResponse(response): if response.code in [200, 201]: return else: msg('Site response: {} {}'.format(response.code, response.phrase)) d_body = readBody(response) d_body.addCallback(cbBody) d_body.addErrback(cbPartial) return d_body def cbError(failure): failure.printTraceback() def processResult(result): try: j = loads(result) msg('Site response: {}'.format(j)) except Exception as e: msg('Failed to parse response: {}'.format(e)) d.addCallback(cbResponse) d.addErrback(cbError) return d
关键说明
- 用
request_queue缓存所有待发送的请求,避免请求丢失 is_processing标记确保同一时间只处理一个请求- 每个请求处理完成(无论成功失败)后,通过
reactor.callLater(1, ...)延迟1秒再启动下一个请求的处理,实现限流效果 addBoth方法保证无论请求成功还是失败,都会触发后续的队列处理逻辑,不会导致队列阻塞
内容的提问来源于stack exchange,提问作者bontchev
相关产品推荐
相关产品推荐

