如何用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
相关产品推荐
相关产品推荐

