
我不是第一次被人问到“Tornado到底能不能扛百万级并发”这种问题了。说实话第一次看到这个数字时我也怀疑过但当我真正把一整套长连接网关拆开、压测、再部署到多机集群之后才理解了所谓“百万级实时服务”的真正含义——它不是一个孤零零的Python进程能创造的奇迹而是靠事件循环、连接管理、消息路由和系统调优一起堆出来的结果。Tornado作为Python生态里最早拥抱异步非阻塞IO的Web框架至今仍然是做WebSocket实时推送、IM、行情网关这类场景最顺手的工具之一。这篇内容不聊花架子就从怎么理解异步、怎么设计架构、怎么写代码、怎么压测调优一条线走完把我实际踩过的坑和验证过的方案原原本本放出来。1. 为什么是Tornado异步框架的底层逻辑1.1 先看同步框架的天花板在聊Tornado之前得先明白一个现实问题同样是Python为什么Flask和Django做长连接服务会被压垮传统WSGI框架的并发模型是“一个请求分一个线程”。同步框架处理请求时线程要一直在那儿等等数据库返回、等Redis返回、等上游HTTP响应。等的时候线程既不能干别的也不能被操作系统拿走只能空转。我打个比方这就像银行网点里一个柜员只服务一个客户客户填表填多久柜员就干等多久。来1000个客户就要招1000个柜员来10万个客户网点就得炸。服务器线程的创建、销毁、上下文切换都是有成本的内存再多也扛不住几十万条线程同时挂着。很多人抱怨“Django单机只能扛几百并发”真不是Django性能差而是同步模型为了让开发简单牺牲了资源的极致利用。1.2 IOLoop和epoll一个线程如何扛下十万连接Tornado走的是完全相反的路子单线程事件循环。它把所有的网络套接字都交给操作系统的epoll监听自己只在一个线程里循环处理“哪个连接有数据来了、哪个连接可以写数据了”。还是用服务员来类比一个服务员管好几桌客人谁举手就过去服务没举手的就继续等服务员永远不闲下来。每桌客人不用单独配一个服务员这就是非阻塞IOI/O多路复用的本质。在Tornado 5.0之前写异步代码主要靠回调函数或者gen.coroutine读起来相当痛苦。Tornado 5.0之后完整支持了原生async/await代码可读性好了一大截。核心组件就是IOLoop和IOLoop.current()整个进程的事件循环只有这一个事件循环负责调度所有协程和socket事件。关键点在于你在异步函数里写await遇到IO等待时事件循环会立刻切去处理别的连接而不是傻等当前这个。这就是“异步”和“并发”在单线程里能共存的原理。1.3 在Tornado、aiohttp、FastAPI与Go之间选型选型永远不是“谁最强选谁”而是“谁最匹配现有团队和场景”。框架模型适合场景注意点Tornado异步事件循环WebSocket网关、实时推送、长连接自带IOLoop模板老但稳定aiohttp纯异步HTTP服务、WebSocket更轻量生态相对分散FastAPIasync支持标准API、文档自动生成底层是Starlette长连接弱一点Go net/httpgoroutine高并发HTTP、长连接语言门槛团队要切换技术栈我在实际选择时有个经验如果团队已经全面Python化项目里又必须大量使用WebSocketTornado的成熟度和踩坑案例是最多的。aiohttp也很好但WebSocket的断线重连处理、连接管理、跨节点广播这些能力Tornado社区的实践资料更丰富。如果只是做纯HTTP API我可能会考虑FastAPI开发效率和类型提示更好。但如果要做百万连接级别的实时推送网关Tornado是Python阵营里最稳妥的答案。1.4 异步的本质代码上到底差在哪很多人卡在“异步”这个概念上说白了就是一句话同步是排队等结果异步是先去干别的结果好了再回来拿。举个例子。同步请求MySQL假设SQL要跑100毫秒这100毫秒线程啥也不干就在那等。异步请求MySQLawait交出去之后事件循环立刻处理其他连接的读写等数据库结果返回了再回来继续执行后面的代码。# 伪代码感受一下两种模型 # 同步方式 def handle_request(): result db.query(sql) # 线程阻塞等待100ms return result # Tornado异步方式 async def handle_request(): result await db.query(sql) # 挂起协程事件循环处理别人 return result区别虽然只有一行但承载能力天差地别。理解了这一点后面所有关于“事件循环不能阻塞”的纪律都会有据可依。2. 百万级实时服务架构拆解与容量评估2.1 先把“百万级并发”的定义说清楚做过线上压测的人都知道百万级并发在Tornado语境下指的几乎都是“百万级长连接数”而不是“百万级QPS”。这两个概念必须分开。QPS是每秒处理的请求数比如API网关给第三方提供接口高峰期每秒几万请求已经很吓人了。长连接数则是同一时刻保持在线状态的WebSocket连接数像行情推送、社交IM、群消息广播用户一连就是几小时甚至几天下不来在线连接数轻松超过百万。为什么说是百万连接而不是百万QPS因为Tornado这类事件循环框架处理单个消息的成本很低低到一次推送可能只涉及一次内存拷贝和一次socket写入。但连接多了之后内存、文件描述符、消息缓冲区才是真正的敌人。我做容量规划时一般先按业务场景计算“每秒每条连接平均收到多少条消息每条消息多大”再倒推需要的带宽和内存。100万连接每5秒推一条128字节的消息每秒就是20万条消息流量大约25MB/s这还没算WebSocket帧头、TCP开销和广播放大。数字一算出来就知道单机单进程不可能做到。2.2 单机容量预算内存、带宽与文件描述符很多人以为提升并发就是改几个内核参数实际上每个参数背后都有物理代价。一条空闲的TCP连接在用户态至少占几KB内存算上socket读写缓冲区、Tornado的handler对象、连接管理字典的开销给一条连接预算是5-10KB是比较合理的。100万连接就意味着至少5-10GB内存只用来“挂连接”这还没算消息队列、Redis缓存和业务状态。带宽也得算。刚才说的每5秒128字节消息100万连接就需要约25MB/s持续下行换算成带宽是200Mbps这还只是平均值如果出现万人同时发消息的峰值瞬时流量可能翻十倍。文件描述符同理。每个连接占用一个fdLinux默认每个进程只能开1024个fd必须调大。但fd数只是起点内核还需要维护TCP控制块、路由表项、软中断队列这些都会吃掉内存。所以我给团队的容量预估公式很简单单机8核16GB撑5-10万长连接是安全的压到20万需要精心调优再往上就交给多机。不要迷信“调几个参数就让单机百万连接”这是硬件的物理边界决定的。2.3 三层架构接入层、消息层与业务层想清楚单机边界之后百万级实时服务的架构就必须分层客户端 - Nginx/LVS负载均衡 - Tornado WebSocket集群接入层 - Redis Pub/Sub或消息队列消息层 - 业务服务集群业务层接入层的职责只有一个维护和客户端的WebSocket连接负责握手、心跳、断线清理。业务层的职责是处理消息、写数据库、做推荐、触发推送。中间的消息层负责跨节点通信。为什么接入层和业务层必须拆开因为WebSocket连接是有状态的比如用户A连接在节点1用户B连接在节点2A给B发消息节点1根本不知道B在哪台机器上。如果让节点1直接查自己的连接管理器只会得到一个“查无此人”。没有消息层跨节点推送只能靠某个中心节点集中路由那这个中心点又会变成单点和瓶颈。所以让所有节点都订阅同一个Redis频道或消息队列消息发进来每个节点都广播到本地连接谁有目标用户谁就推出去这才是标准解法。2.4 为什么接入层一定要无状态我见过很多初版架构会犯一个糊涂把用户上次的聊天记录、未读消息数、好友关系全塞进接入层进程里。这会导致节点之间状态不同步重启一台机器一部分用户的数据就丢了。接入层的正确姿势是“连接在本地状态在远端”。每个Tornado进程只维护一张本地连接表记录user_id到WebSocket handler的映射。用户的会话信息、离线消息、路由规则全部放到Redis或数据库里。这样做的核心价值是扩容和容灾。某个节点挂了Nginx直接摘掉它新连接分发到其他节点用户重新连接后从Redis恢复会话状态业务完全没有感知。如果接入层有状态节点迁移就是一场灾难你要先序列化所有连接状态再同步到新节点期间还不能断线。3. 从零搭建异步实时推送服务3.1 项目结构、依赖与最小启动先用最基础的方式把一个Tornado服务跑起来。项目结构保持精简push_service/ ├── app.py # 入口 ├── handlers/ │ ├── __init__.py │ ├── ws.py # WebSocket处理器 │ └── api.py # HTTP接口发消息/查询 ├── connection.py # 连接管理器 ├── redis_client.py # 异步Redis客户端 └── requirements.txt依赖就这四个数量非常克制tornado6.4 redis5.0.0 asyncpg0.29 sqlalchemy[asyncio]2.0最小启动代码先让Tornado听在一个端口上import tornado.web import tornado.ioloop class MainHandler(tornado.web.RequestHandler): async def get(self): self.write({message: tornado running}) def make_app(): return tornado.web.Application([ (r/, MainHandler), ]) if __name__ __main__: app make_app() app.listen(8888) tornado.ioloop.IOLoop.current().start()两个细节值得注意。第一app.listen(8888)默认绑定了所有网卡生产环境最好通过配置文件指定IP和端口方便用systemd/Nginx管理多实例。第二IOLoop.current().start()一旦启动就进入事件循环后面的代码不会再执行所以初始化要在start之前完成。3.2 连接管理器与WebSocket处理器Tornado真正适合做实时服务的核心组件是WebSocketHandler。要支撑百万级长连接所有连接必须集中管理不能散落在各个handler里。先写一个连接管理器按user_id保存当前连接class ConnectionManager: def __init__(self): self._connections {} async def add(self, user_id: str, handler): old self._connections.get(user_id) if old and old.ws_connection: old.close() # 顶号处理 self._connections[user_id] handler def remove(self, user_id: str): self._connections.pop(user_id, None) async def send_to_user(self, user_id: str, message: dict): handler self._connections.get(user_id) if handler: await handler.write_message(message) async def broadcast(self, message: dict): for handler in list(self._connections.values()): try: await handler.write_message(message) except Exception: continueWebSocket处理器本身逻辑不复杂import json import tornado.websocket class WsHandler(tornado.websocket.WebSocketHandler): def initialize(self, manager): self.manager manager self.user_id None async def open(self): self.user_id self.get_query_argument(user_id, ) await self.manager.add(self.user_id, self) await self.write_message({type: connected, msg: ok}) async def on_message(self, message): data json.loads(message) if data.get(type) ping: await self.write_message({type: pong}) elif data.get(to): await self.manager.send_to_user(data[to], data) def on_close(self): if self.user_id: self.manager.remove(self.user_id)踩坑经验on_close里remove用户时要确认这个连接还是不是当前连接否则用户A掉线后又立刻重连旧的on_close可能会把新连接误删。稳妥做法是add时记录一个session_idremove时对比session_id不一致就不删。这是小概率但真实会发生的问题。3.3 事件循环纪律那些让你“看似并发、实则串行”的坑Tornado的整个并发能力建立在“事件循环永不被阻塞”这个前提下。只要有一个请求让事件循环卡住所有连接全部陪葬。最典型的几个坑我在代码评审里几乎每次都能看到async函数里直接time.sleep(1)整个服务卡1秒。用requests.get调用外部HTTP接口同步等待3秒卡3秒。用psycopg2同步查数据库查询耗时500ms这500ms内所有WebSocket都收不到消息。大量解析超大JSON或执行正则处理CPU被抢走事件循环轮转变慢。正确姿势也很直接# 错误示范 async def bad_handler(self): time.sleep(1) self.write({msg: boom}) # 正确示范 async def good_handler(self): await asyncio.sleep(1) self.write({msg: ok})我记得第一次在生产环境排查“所有用户同时卡顿”的故障最后定位到的原因就是有人在异步请求里同步调了短信服务商的requests接口明明是50ms的IO等待却被放大成了整个服务的一次“集体暂停”。从此之后我定了一条铁律异步函数里禁止出现同步IO库代码评审一票否决。3.4 用异步Redis打通跨节点广播在3.2的例子里send_to_user只能发到本机连接。真正生产环境有几十台Tornado节点跨节点推送必须借助Redis。redis-py从4.2开始内置了redis.asyncio不需要再单独装aioredis库这个变化很多人还不知道。用法如下import redis.asyncio as aioredis redis_client aioredis.from_url( redis://127.0.0.1:6379, decode_responsesTrue )每个Tornado进程启动时开启一个订阅任务import asyncio import json class RedisSubscriber: def __init__(self, manager): self.manager manager async def start(self): pubsub redis_client.pubsub() await pubsub.subscribe(chat:global) asyncio.create_task(self._listen(pubsub)) async def _listen(self, pubsub): try: async for message in pubsub.listen(): if message[type] message: data json.loads(message[data]) await self.manager.broadcast(data) except asyncio.CancelledError: pass配套的HTTP发送接口进程收到业务请求后publish到Redis频道所有节点再统一广播class SendHandler(tornado.web.RequestHandler): async def post(self): body json.loads(self.request.body) await redis_client.publish(chat:global, json.dumps(body)) self.write({ok: True})Redis订阅端有一个非常隐蔽的坑如果连接意外断开pubsub不会自动重连listen()会静默结束消息从此再也进不来。必须在_listen里捕获异常并做指数退避重连或者定期检查pubsub.connection的状态。3.5 异步数据库读写从asyncpg到SQLAlchemy 2.0实时服务肯定要持久化消息记录、用户在线状态。数据库这一层同样必须异步。两个主流方案一个是用asyncpg直接写SQL另一个是用SQLAlchemy 2.0的异步扩展。后者对业务代码更友好工程化程度高我先给SQLAlchemy的配置from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker engine create_async_engine( postgresqlasyncpg://user:passwordlocalhost/push_db ) Session async_sessionmaker(engine) class UserHandler(tornado.web.RequestHandler): async def get(self): async with Session() as session: result await session.execute(text(SELECT 1)) value result.scalar_one() self.write({result: value})一个必须记住的规则异步Session是绑定协程的必须在使用它的那个连续await链里创建和关闭不能把session塞到全局变量或跨协程共享否则会碰到莫名其妙的连接池报错。SQLAlchemy异步和同步的争论在社区里一直存在我的实测感受是简单查询两者差异不大但高并发写场景下异步驱动能稳定占满数据库连接池而同步驱动在Tornado事件循环里会造成灾难性的阻塞。如果你的程序已经全面Tornado化数据库驱动没有理由不换异步。3.6 耗时任务怎么办run_in_executor与线程池事件循环不能阻塞但有些任务是纯CPU密集型比如图片压缩、加解密、复杂规则匹配。这种情况不适合放到事件循环里也不适合放到协程里因为协程不解决CPU占用问题。标准解法是用线程池兜底把CPU任务丢给线程池执行import asyncio from concurrent.futures import ThreadPoolExecutor executor ThreadPoolExecutor(max_workers8) class CpuHandler(tornado.web.RequestHandler): async def get(self): loop asyncio.get_running_loop() result await loop.run_in_executor( executor, heavy_cpu_func, request_data ) self.write({result: result})虽然线程池内的任务会阻塞工作线程但只要线程池大小合理事件循环本身不会被阻塞其他网络IO依旧顺畅。真正的终极方案是直接把这类任务下沉到Celery或独立worker进程Tornado只发任务编号异步等结果。线程池适合低频、轻量的CPU任务任务量大时一定要独立进程。4. 性能压测与工程调优4.1 压测工具与异步压测脚本没压测过就上线等于裸奔。压Tornado WebSocket服务工具有几个选择。wrk是压HTTP接口的老牌工具但对WebSocket支持不行。websocat可以模拟单个连接却做不了高并发。我自己最常用的还是自己写asyncio压测客户端可以精确控制连接数、消息频率和持续时间。import asyncio import websockets async def one_client(index): uri fws://127.0.0.1:8888/ws?user_iduser_{index} async with websockets.connect(uri) as ws: for _ in range(100): await ws.send(ping) await ws.recv() async def main(): tasks [one_client(i) for i in range(20000)] await asyncio.gather(*tasks) asyncio.run(main())这里有个很多人没意识到的细节压测脚本本身也必须是异步的如果你用同步多线程去模拟2万客户端客户端自己先把线程池和内存打满了服务端还没开始热身瓶颈就出现在压测端。压测时记录三个关键指标建立完整连接所需时间、每秒消息吞吐量、服务端内存增长曲线。连接建立速度如果断崖式下跌大多数时候是文件描述符或者内核队列满了。4.2 内核参数与文件描述符调优让服务端能开更多连接内核参数必须调整。# 文件描述符上限 echo fs.file-max 1000000 /etc/sysctl.conf # 全连接队列和半连接队列防止握手丢包 echo net.core.somaxconn 65535 /etc/sysctl.conf echo net.ipv4.tcp_max_syn_backlog 65535 /etc/sysctl.conf # 客户端主动连接所需的临时端口范围 echo net.ipv4.ip_local_port_range 1024 65535 /etc/sysctl.conf sysctl -p进程级别的fd限制也要改。如果是用systemd管理服务在service文件里加LimitNOFILE1048576。如果直接命令行启动需要先ulimit -n 1048576再启动进程。我提醒一句这些参数不是越大越好。somaxconn和backlog开太大会让TCP握手请求堆积在队列里反而增加内存和延迟。合理的做法是压测时观测ss -lnt的Recv-Q确认队列不积压就行。4.3 多进程部署用满多核CPUTornado的每个进程只有一条事件循环也就是说一个Tornado进程最多只吃满一个CPU核。8核机器只跑一个进程剩下7个核全空闲这是最大的资源浪费。标准做法是按CPU核数启动多个Tornado进程每个进程监听不同端口前面用Nginx做负载均衡。用supervisor管理最省心[program:tornado_ws] commandpython app.py --port800%(process_num)s numprocs4 process_name%(program_name)s_%(process_num)s autorestarttrue redirect_stderrtrue stdout_logfile/var/log/tornado_ws.log注意每个进程要有独立的--port参数否则都绑同一个端口必然报错。我之前踩过这个坑配置里写了固定端口结果4个进程只有一个起来剩下三个全部Address already in use。看到日志的一瞬间头都是大的。多进程带来的另一个变化是Redis订阅重复每个进程都会订阅同一个频道推送消息会被每个进程分别广播一次。这正是我们希望的广播语义但要注意广播消息的量级如果频道消息非常密集每个进程都全量广播网络开销会成倍放大。这时候要考虑按用户分组订阅特定频道而不是无脑全量broadcast。4.4 Nginx反向代理WebSocket的细节多进程服务前面必须有一个反向代理做负载均衡和TLS终止。Nginx配置里藏着两个常见的WebSocket坑。第一个坑是必须显式设置Upgrade头。普通HTTP代理不带这个头WebSocket握手会直接失败。第二个坑是proxy_read_timeout默认60秒。这个时间一到Nginx会主动断开空闲的WebSocket连接。如果你的客户端没有心跳机制就会发现连接每60秒准时断开一次监控图上线了一堆有规律的重连尖峰查起来非常诡异。实际可用的Nginx配置upstream tornado_ws { least_conn; server 127.0.0.1:8001; server 127.0.0.1:8002; server 127.0.0.1:8003; server 127.0.0.1:8004; } server { listen 80; location /ws { proxy_pass http://tornado_ws; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection upgrade; proxy_set_header Host $host; proxy_read_timeout 3600s; proxy_send_timeout 3600s; } }心跳机制也必须有推荐30秒一次。Tornado的WebSocketHandler支持周期性ping客户端收到ping要主动回pong服务端连续几次没收到pong就可以主动断掉避免僵尸连接占用资源。4.5 长连接服务的内存与GC调优长连接服务跑久了内存增长几乎是一定的只不过原因各不相同。最常见的原因是write_message时客户端不消费。Tornado的WebSocket写入是有缓冲区的如果客户端断网了但它又没触发on_close服务端持续往这个连接发消息缓冲区只会越积越大最终吃掉大量内存。处理方式是每次写入后检查缓冲区大小超过阈值直接强杀连接。async def safe_write(self, handler, message): await handler.write_message(message) buffer_size handler.ws_connection.get_write_buffer_size() if buffer_size 256 * 1024: handler.close()第二个常见原因是连接管理器里的handler泄漏。前面提到的on_close误删问题如果处理不当用户断开后handler还留在字典里长此以往内存只增不减。建议连接管理器里增加定时清理任务每隔几分钟就遍历一遍把ws_connection为空的连接清掉。说到GC调优我不建议在服务里定期调用gc.collect()这会触发全局停顿反而影响实时性。Python的引用计数能解决绝大多数回收问题真正要排查的是谁拿着对象的强引用不放用tracemalloc或者py-spy dump一下就能看到。5. 实战中踩过的坑问题排查速查表5.1 高频问题排查表现象可能原因排查方法解决方案连接数几千就卡死文件描述符不够、同步库阻塞ulimit -n看上限py-spy dump栈调fd上限禁止同步IO进async函数连接60秒准时断Nginx代理超时看重连时间间隔是否固定调proxy_read_timeout加心跳跨节点消息丢失Redis Pub/Sub断连不重连监控订阅连接状态指数退避重连升级Streams单进程CPU跑满事件循环里有死循环或CPU密集py-spy dump看当前栈run_in_executor或下沉任务队列内存持续上涨发送缓冲堆积、handler泄漏看write_buffer_size统计连接表长度阈值断开定时清理无效连接用户发消息对方收不到路由表只在本机或旧连接误删打印节点名和user_id归属通过Redis订阅广播session_id防误删5.2 三个真实故障复盘第一个故障发生在上线后第二天。某一台机器连接数涨到3万左右CPU突然冲到100%全服务卡顿。用py-spy dump一查栈里全是某同步加密库的执行帧。原因是一个旧业务方在异步处理器里引入了同步AES库加密不算慢但被高并发放大以后就把事件循环堵死了。后来把所有CPU密集操作全部下沉到线程池又加了熔断才解决。第二个故障是WebSocket连接“间歇性掉落”。客户端监控显示每隔59秒就重连一次时间准得像闹钟。顺着Nginx默认超时一查果然是proxy_read_timeout60s在作怪。心跳没配连接又空闲Nginx就“好心”帮我们断掉了。看过这个故障之后我再也不相信“架构能自动保持连接”这种话了。第三个故障是压测时候发现的。本地起了一个Tornado服务用2万个asyncio连接压测结果服务端没垮压测客户端先报Too many open files。当时根本没想过客户端也要调ulimit白白浪费了一个下午。这件事教会我一个道理压测之前先把两端的基础资源都调好不然问题根本分不清是服务端还是压测端。5.3 排查工具箱实时服务出问题时排查工具比猜重要得多。我常用的组合是py-spy、tracemalloc和ss命令。# 打印运行中进程的调用栈秒级定位“停在哪” py-spy dump --pid 12345 # 查看TCP连接状态分布 ss -s ss -lnt # 看每个连接接收队列是否积压 ss -lnt | awk {print $2} | sort | uniq -cpy-spy可以无侵入地查看一个正在运行的Python进程当前在干什么这比看日志高效太多。服务卡顿时先dump一下如果所有协程栈都停在同一个地方基本就是那里出了问题。tracemalloc适合内存问题import tracemalloc tracemalloc.start() # 运行一段时间后打印 snapshot tracemalloc.take_snapshot() top_stats snapshot.statistics(lineno) for stat in top_stats[:10]: print(stat)它能直接告诉你内存是哪个文件哪一行代码分配的比逐行读代码猜快多了。说到最后我还是想强调一点Tornado的“百万级并发”不是靠某一个黑科技参数实现的而是靠你对事件循环纪律的执念——永远不要让同步操作住进异步函数永远不要忽略长连接的内存边界永远先压测再上线。我见过太多人上来就改了一堆sysctl参数最后发现真正的瓶颈只是代码里多了一次requests.get。把地基打好再去追逐那些闪闪发光的大数字才有意义。