You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何为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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.19 08:43:16