构建永不崩溃的可扩展 Python 系统
接口超时、队列堆积、内存一点点往上爬,这三样东西一旦同时出现,我一般先不怀疑“流量变大了”,我先怀疑系统设计里有几个地方把失败当成了小概率事件。
Python 系统要想真扛事,靠的不是“代码不报错”,而是某一段挂了、慢了、重复了,系统还能不能继续活着,继续扩,继续把账对上。
“永不崩溃”当然不是字面意思。机器会挂,网络会抖,依赖方会抽风。真正该追的,是别让一个点的小毛病,顺着调用链滚成整条链路的事故。
我见过不少 Python 服务,业务逻辑不复杂,崩得却很有规律:同步调用太深,重试没边界,任务全堆内存,日志一多 I/O 跟着炸。看着像偶发,复盘基本都能落到一句话:没有隔离,没有兜底,没有限流。
先看一个很常见的写法,支付成功后顺手做三件事:发券、发消息、记审计日志。
import requests
defon_order_paid(order_id: str, user_id: str, amount: int) -> None:
requests.post("http://coupon-service/grant", json={"user_id": user_id, "amount": amount}, timeout=3)
requests.post("http://message-service/notify", json={"user_id": user_id, "order_id": order_id}, timeout=3)
requests.post("http://audit-service/record", json={"order_id": order_id, "amount": amount}, timeout=3)
这段代码平时跑得挺像样,一到高峰就开始恶心人。优惠券服务抖一下,主流程跟着超时;消息服务慢一点,接口 RT 直接翻倍;再碰上调用方重试,雪上加霜。
这种地方我一般不继续加 try/except 自我安慰,先拆主链路。主流程只做“最小闭环”,剩下的异步化,能落库就先落库,能进队列就先别同步等。
import json
import time
from dataclasses import dataclass
@dataclass
classOutboxEvent:
event_id: str
topic: str
payload: dict
created_at: float
defcreate_outbox_event(conn, event: OutboxEvent) -> None:
conn.execute(
"""
insert into app_outbox(event_id, topic, payload, status, created_at)
values (%s, %s, %s, 'NEW', %s)
""",
(event.event_id, event.topic, json.dumps(event.payload), event.created_at)
)
defon_order_paid(conn, order_id: str, user_id: str, amount: int) -> None:
conn.execute(
"update orders set status='PAID', paid_at=now() where order_id=%s and status='INIT'",
(order_id,)
)
create_outbox_event(conn, OutboxEvent(
event_id=f"pay:{order_id}",
topic="order_paid",
payload={"order_id": order_id, "user_id": user_id, "amount": amount},
created_at=time.time()
))
conn.commit()
这里不是炫架构,核心就一个:业务状态和待发送事件放进同一个事务。订单先落稳,后面谁慢谁补,别把主请求绑死在下游身上。
然后再起一个后台投递器,失败重试,但次数要有数,别重试到把自己打死。
import json
import logging
import time
import requests
log = logging.getLogger("dispatcher")
defdispatch_loop(conn):
whileTrue:
rows = conn.query(
"""
select id, event_id, topic, payload, retry_count
from app_outbox
where status in ('NEW', 'RETRY')
order by id
limit 100
for update skip locked
"""
)
for row in rows:
try:
requests.post(
"http://mq-gateway/publish",
json={"topic": row["topic"], "payload": json.loads(row["payload"])},
timeout=1.5
)
conn.execute("update app_outbox set status='DONE' where id=%s", (row["id"],))
except Exception as ex:
conn.execute(
"""
update app_outbox
set status=%s, retry_count=retry_count+1, last_error=%s
where id=%s
""",
("DEAD"if row["retry_count"] >= 5else"RETRY", str(ex)[:200], row["id"])
)
log.warning("dispatch failed event_id=%s err=%s", row["event_id"], ex)
conn.commit()
time.sleep(0.2)
注意这里有两个味道很重的点。
一个是 skip locked,说明别让多个 worker 抢同一批任务。另一个是 DEAD,说明重试不是信仰,重试是成本。超过阈值就进死信,留给人查,不要无限转圈。
再往前一步,系统为什么会在高峰期突然“崩”?很多时候不是 CPU 满了,是你来多少收多少,线程、协程、队列、连接池一起被灌爆。
所以入口必须限流,内部必须背压。尤其 Python,I/O 场景里大家很爱上 asyncio,但协程多不等于系统就稳。没有并发上限,照样把 Redis、DB、第三方接口一锅端。
import asyncio
import httpx
semaphore = asyncio.Semaphore(100)
asyncdeffetch_profile(client: httpx.AsyncClient, user_id: str) -> dict:
asyncwith semaphore:
resp = await client.get(f"http://user-service/users/{user_id}", timeout=1.0)
resp.raise_for_status()
return resp.json()
很多人看到 async 就兴奋,结果一口气起几万个任务。代码没崩,依赖先崩了。并发控制这种东西,平时看着碍事,出事时就是保险丝。
另外一个常被低估的点是幂等。
支付回调重放一次,任务补偿跑两遍,消费者 rebalance 之后重复消费,这些都太正常了。你要是还拿“这个请求理论上只会来一次”当设计前提,系统迟早给你上一课。
defhandle_order_paid(conn, event_id: str, order_id: str) -> None:
inserted = conn.execute(
"insert into consume_log(event_id, created_at) values (%s, now()) on conflict do nothing",
(event_id,)
)
if inserted.rowcount == 0:
return
conn.execute(
"update orders set settled = true where order_id=%s and settled = false",
(order_id,)
)
conn.commit()
我一般不太信“消息系统保证不重复”这种说法。业务自己做幂等,晚上睡得踏实一点。
最后说扩展。很多 Python 系统扩不起来,不是代码性能差,是状态全塞进进程里。缓存放本地 dict,任务队列放内存 list,用户会话绑单机,扩容当然难看。新机器一加,状态没了;老机器一挂,数据也跟着没了。
能外置的状态尽量外置:配置进配置中心,会话进 Redis,任务进 MQ,文件进对象存储。应用实例尽量做成“随时能死、死了能起”的样子。真要扩容的时候,不用先祈祷。
我现在看一个 Python 服务稳不稳,基本就扫这几件事:主链路是不是够短,失败有没有边界,消息能不能补,消费会不会重,状态是不是外置,并发有没有上限。
这些东西平时不显山不露水,一到线上抖两下,差距就全出来了。
系统不是靠一次写对就永不崩溃,系统是靠你提前承认它一定会出错,然后把出错这件事,设计成不会把整盘棋带走。