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

合并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...")

关键修改点说明

  1. 全局消息队列:message_queue是asyncio的异步队列,安全地在两个服务之间传递数据,避免直接调用带来的协程冲突。
  2. TCP协议类调整:在data_received中,改为保留长连接(若需要短连接可取消注释关闭代码),并用asyncio.create_task将消息放入队列,保证异步操作不阻塞。
  3. WebSocket Handler扩展:添加listen_message_queue协程任务,在客户端连接后自动启动,持续监听队列,一旦有新消息就推送给所有浏览器客户端。
  4. 统一启动流程:用asyncio.run(main())启动所有服务,确保TCP和WebSocket服务都运行在同一个事件循环中。
  5. 错误处理:在write_to_clients中添加了异常捕获,避免单个客户端连接异常导致整个推送流程崩溃。

测试方法

  1. 运行上述代码,启动两个服务。
  2. 用TCP客户端(比如telnet 127.0.0.1 8889或自定义TCP客户端)连接并发送消息。
  3. 打开多个浏览器窗口访问http://127.0.0.1:8888,可以看到TCP客户端发送的消息会实时推送到所有浏览器窗口。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:50:27