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

如何用Twisted/Asyncio实现接收数据并批量推送的Buffer应用?

用Twisted和Asyncio实现数据接收+批量推送的方案

嘿,我来帮你搞定这个批量推送的需求!先从你已经写了一部分的Twisted代码入手完善,再给你Asyncio的实现方案,这样你可以根据自己的技术栈选合适的。

一、完善你的Twisted实现

你现有的代码已经搭好了接收数据的架子,咱们把批量推送的逻辑补全,还要注意Twisted的异步特性——绝对不能在回调里写阻塞代码哦。

完整实现示例

from twisted.internet import protocol, reactor, endpoints
from twisted.protocols import basic
from twisted.web.client import Agent, readBody
from twisted.web.http_headers import Headers

class FirehoseProtocol(basic.LineReceiver):
    def __init__(self, batch_size=10):
        self.data = []
        self.batch_size = batch_size  # 批量推送的阈值
        self.target_service_url = "http://your-target-service/push"  # 目标服务地址

    def lineReceived(self, line):
        self.data.append(line)
        # 当数据达到批量阈值时触发推送
        if len(self.data) >= self.batch_size:
            self.push_to_firehose()

    def connectionLost(self, reason):
        # 连接断开时,把剩余的未推送数据发出去
        if self.data:
            self.push_to_firehose()

    def push_to_firehose(self):
        # 异步推送数据到目标服务,这里用Twisted的Agent做HTTP请求示例
        agent = Agent(reactor)
        # 把批量数据拼接成适合目标服务的格式,比如换行分隔的字符串
        payload = "\n".join(self.data).encode("utf-8")
        # 发送POST请求
        d = agent.request(
            b'POST',
            self.target_service_url.encode("utf-8"),
            Headers({'Content-Type': ['text/plain']}),
            protocol.FileBodyProducer(payload)
        )

        def on_push_success(response):
            print(f"批量推送成功,状态码: {response.code}")
            # 推送成功后清空本地缓存
            self.data = []
            return readBody(response)  # 读取响应体(可选)

        def on_push_failure(failure):
            print(f"批量推送失败: {failure.getErrorMessage()}")
            # 这里可以根据需求做重试逻辑,比如把数据放回队列

        d.addCallback(on_push_success)
        d.addErrback(on_push_failure)

class EchoFactory(protocol.ServerFactory):
    def buildProtocol(self, addr):
        # 这里可以传入批量大小等配置
        return FirehoseProtocol(batch_size=10)

# 启动服务,监听本地8000端口
endpoints.serverFromString(reactor, "tcp:8000").listen(EchoFactory())
reactor.run()

关键说明

  • 批量触发逻辑:我加了batch_size参数,每收到batch_size条数据就自动推送;另外在连接断开时会把剩余数据推送,避免丢失。
  • 异步推送:用Twisted的Agent做HTTP客户端,全程异步不阻塞reactor,这是Twisted的核心要求——任何IO操作都不能阻塞事件循环。
  • 错误处理:添加了回调和错误回调,你可以根据实际需求扩展重试、告警等逻辑。

二、Asyncio实现方案

如果你更熟悉Asyncio的生态,也可以用它来实现,逻辑和Twisted类似,但写法更贴近Python原生的异步语法。

完整实现示例

import asyncio
from aiohttp import ClientSession

class BatchPushServer:
    def __init__(self, host='0.0.0.0', port=8000, batch_size=10, target_url="http://your-target-service/push"):
        self.host = host
        self.port = port
        self.batch_size = batch_size
        self.target_url = target_url
        self.data_buffer = []

    async def handle_client(self, reader, writer):
        while True:
            line = await reader.readline()
            if not line:
                # 客户端断开连接,推送剩余数据
                if self.data_buffer:
                    await self.push_batch()
                break
            # 把字节转成字符串,去掉末尾换行
            line_str = line.decode('utf-8').strip()
            self.data_buffer.append(line_str)
            # 达到批量阈值就推送
            if len(self.data_buffer) >= self.batch_size:
                await self.push_batch()

    async def push_batch(self):
        async with ClientSession() as session:
            payload = "\n".join(self.data_buffer)
            try:
                async with session.post(self.target_url, data=payload, headers={'Content-Type': 'text/plain'}) as resp:
                    if resp.status == 200:
                        print(f"批量推送成功,共推送{len(self.data_buffer)}条数据")
                        self.data_buffer = []
                    else:
                        print(f"推送失败,状态码: {resp.status}")
                        # 这里可以处理重试逻辑
            except Exception as e:
                print(f"推送请求出错: {str(e)}")

    async def start(self):
        server = await asyncio.start_server(self.handle_client, self.host, self.port)
        async with server:
            await server.serve_forever()

if __name__ == "__main__":
    server = BatchPushServer(batch_size=10)
    asyncio.run(server.start())

关键说明

  • 原生异步语法:用async/await写法,更符合Python开发者的习惯。
  • aiohttp客户端:用aiohttp做异步HTTP请求,和Asyncio生态完美适配。
  • 连接断开处理:同样在客户端断开时推送剩余数据,保证数据不丢失。

不管你选Twisted还是Asyncio,核心思路都是异步接收数据+批量缓存+异步推送,避免阻塞事件循环,这样才能保证服务的高并发能力。

内容的提问来源于stack exchange,提问作者Alex Tonkonozhenko

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 08:27:36