合并Asyncio TCP Server与Tornado WebSocket,实现数据异步转发
解决方案:整合Asyncio TCP服务器与Tornado WebSocket服务
要实现你的需求,核心是让两个服务共享同一个事件循环,并建立一个异步消息通道来传递TCP客户端的数据到WebSocket客户端。下面是完整的实现方案,我会一步步解释关键修改点:
关键思路
- 共用事件循环:Tornado 5+ 默认使用asyncio的事件循环,所以我们可以直接把TCP服务和WebSocket服务绑定到同一个loop上,避免多循环的冲突。
- 异步消息队列:用
asyncio.Queue作为TCP服务和WebSocket服务之间的桥梁,TCP收到数据后放入队列,WebSocket服务持续监听队列并推送数据给所有浏览器客户端。 - 调整端口:原来两个服务都用8888,需要修改其中一个的端口(比如TCP服务用8889),避免端口冲突。
完整代码实现
import asyncio import tornado.httpserver import tornado.websocket import tornado.ioloop import tornado.web # 全局异步队列,用于TCP和WebSocket之间传递消息 message_queue = asyncio.Queue() class EchoServerClientProtocol(asyncio.Protocol): def connection_made(self, transport): peername = transport.get_extra_info('peername') print('TCP Connection from {}'.format(peername)) self.transport = transport def data_received(self, data): message = data.decode().strip() print('TCP Data received: {!r}'.format(message)) # 将收到的消息放入队列,交给WebSocket服务推送 asyncio.create_task(message_queue.put(message)) # 可选:给TCP客户端回复确认 self.transport.write(f"Received: {message}".encode()) # 这里保留长连接,若需要短连接可取消注释下面一行 # self.transport.close() class WSHandler(tornado.websocket.WebSocketHandler): clients = [] def check_origin(self, origin): return True def open(self): print('New WebSocket Connection Established') self.write_message("Connected to WebSocket Server...") WSHandler.clients.append(self) # 启动监听队列的后台任务 asyncio.create_task(self.listen_message_queue()) async def listen_message_queue(self): """持续监听消息队列,有新消息就推送给所有客户端""" while True: message = await message_queue.get() print(f"Pushing message to WebSocket clients: {message}") self.write_to_clients(f"From TCP Client: {message}") message_queue.task_done() def on_message(self, message): print('WebSocket Message Received: {}'.format(message)) def on_close(self): print('WebSocket Connection Closed...') if self in WSHandler.clients: WSHandler.clients.remove(self) @classmethod def write_to_clients(cls, message): """给所有连接的WebSocket客户端发消息""" for client in cls.clients: try: client.write_message(message) except Exception as e: print(f"Failed to send message to client: {e}") application = tornado.web.Application([ (r'/', WSHandler), ]) async def main(): # 获取共用的事件循环 loop = asyncio.get_running_loop() # 启动TCP服务器 tcp_server = await loop.create_server(EchoServerClientProtocol, '127.0.0.1', 8889) print('TCP Server serving on {}'.format(tcp_server.sockets[0].getsockname())) # 启动Tornado WebSocket服务器 http_server = tornado.httpserver.HTTPServer(application) http_server.listen(8888) print('WebSocket Server serving on http://127.0.0.1:8888') # 保持服务运行 await asyncio.Event().wait() if __name__ == "__main__": try: asyncio.run(main()) except KeyboardInterrupt: print("Shutting down servers...")
关键修改点说明
- 全局消息队列:
message_queue是asyncio的异步队列,安全地在两个服务之间传递数据,避免直接调用带来的协程冲突。 - TCP协议类调整:在
data_received中,改为保留长连接(若需要短连接可取消注释关闭代码),并用asyncio.create_task将消息放入队列,保证异步操作不阻塞。 - WebSocket Handler扩展:添加
listen_message_queue协程任务,在客户端连接后自动启动,持续监听队列,一旦有新消息就推送给所有浏览器客户端。 - 统一启动流程:用
asyncio.run(main())启动所有服务,确保TCP和WebSocket服务都运行在同一个事件循环中。 - 错误处理:在
write_to_clients中添加了异常捕获,避免单个客户端连接异常导致整个推送流程崩溃。
测试方法
- 运行上述代码,启动两个服务。
- 用TCP客户端(比如
telnet 127.0.0.1 8889或自定义TCP客户端)连接并发送消息。 - 打开多个浏览器窗口访问
http://127.0.0.1:8888,可以看到TCP客户端发送的消息会实时推送到所有浏览器窗口。
内容的提问来源于stack exchange,提问作者user9018881
相关产品推荐
相关产品推荐

