简介:爱投票系统FastAPI后端项目(二)是上一期《爱投票系统 - FastApi后端项目(一)》的功能扩展包,需配合前置资源才能完整运行,面向已初步掌握FastAPI、希望深入异步项目开发的Python工程师。该包新增Celery定时任务与fund_shares两个核心模块,覆盖投票系统的异步任务调度、基金份额数据处理等场景,并附有CentOS下独立部署Celery的文字步骤教程。压缩包共59个文件,包含33个Python源文件(路由、模型、Celery配置与任务、基金份额相关逻辑)和26个pyc编译文件,整体仅64KB,体量轻便。目前已有1083人学习下载,适合作为实战参考,帮助读者理解FastAPI项目模块划分、Celery集成、Linux服务器部署细节,也可为后续扩展投票功能提供可复用的代码结构。
1. 爱投票系统(二)开工:fastApi 后端项目这一篇填哪个坑
投票系统听起来是个最小业务:用户发起一个投票,别人点一下选项,页面上票数涨一点。真去实现会发现,大部分时间不是在写 POST 接口,而是在回答三个问题:一次投票落了几张表、怎么拦掉同一个人的第二次点击、票数涨起来之后怎么不让数据库成为瓶颈。这个“爱投票系统(二)”把你当已经跑通 FastAPI 骨架的人,第二篇专门填充投票业务本体,范围锁在三块:数据建模、创建投票的完整链路、高并发下的投票与实时结果。方案不绑云厂商,一个 MySQL 一个 Redis 就能复现,初学的人能照着跑起来,写过几个接口的人也能在这里面找到边界参数和踩坑点。网上 FastAPI 教程大多是单个接口演示,投票这种“读写都重”的场景很少串起来讲,这篇算是把它补上。
2. 爱投票系统的数据模型:poll、option、vote 三张表怎么建
建表之前先想清楚:投票页面最终要回答两类问题,一是“这个活动有哪些选项”,二是“每个选项得了多少票”。如果只有一张表,第二个问题必然要走 JSON 解析;如果拆成两张表,又丢了“是谁投的”这个审计信息。poll、option、vote 三张表的分法是投票系统最朴素的落法,也是后面所有查询和约束的地基。
2.1 三张表的职责边界与选型理由
一次投票活动在系统里是一个 poll:标题、类型、起止时间、状态。poll 下面挂着投票人能看到的所有候选答案,也就是 option。用户真正投出的一票是一条 vote,它记录投给了哪些选项,以及可选的账号、IP、设备指纹。三者的关系是 poll 一对多 option,poll 一对多 vote。
| 表 | 一行代表什么 | 核心字段 | 数据特点 |
|---|---|---|---|
| poll | 一次投票活动 | title、poll_type、start_at/end_at、status | 低频变化 |
| poll_option | 一个候选选项 | poll_id、label、display_order、is_active | 仅在创建和关闭时变化 |
| vote | 一条投票记录 | vote_no、poll_id、option_ids、user_id | 只插入,不修改 |
为什么不把 option 塞进 poll 的一个 JSON 字段?因为运营后期会调整选项顺序、下架某个违规选项,拆成表之后这些操作是对独立行更新;vote 里的 option_ids 存的是投票当时的选项快照,即使选项后来改名,历史票数也不受影响。为什么 vote 不用一张 option-vote 关联表?一次投票通常是批量写入几个选项 ID,关联表在读取“这一次投了谁”时要多两次 join,而按选项聚合统计始终走 Redis 和缓存,关联表的结构收益不高。
2.2 用 SQLAlchemy 2.0 定义投票模型
# app/models/poll.py from datetime import datetime from typing import Optional from sqlalchemy import String, Text, Integer, DateTime, ForeignKey, JSON, Boolean, UniqueConstraint from sqlalchemy.orm import Mapped, mapped_column, relationship from app.db.base import Base class Poll(Base): __tablename__ = "poll" id: Mapped[int] = mapped_column(primary_key=True, autoincrement=True) title: Mapped[str] = mapped_column(String(128), nullable=False) description: Mapped[Optional[str]] = mapped_column(Text, default=None) poll_type: Mapped[str] = mapped_column(String(16), default="single") # single / multi max_options: Mapped[int] = mapped_column(Integer, default=1) status: Mapped[str] = mapped_column(String(16), default="draft") # draft / running / closed start_at: Mapped[datetime] = mapped_column(DateTime, nullable=False) end_at: Mapped[datetime] = mapped_column(DateTime, nullable=False) created_by: Mapped[Optional[str]] = mapped_column(String(64), default=None) created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow) options: Mapped[list["PollOption"]] = relationship( back_populates="poll", cascade="all, delete-orphan" ) class PollOption(Base): __tablename__ = "poll_option" id: Mapped[int] = mapped_column(primary_key=True, autoincrement=True) poll_id: Mapped[int] = mapped_column(ForeignKey("poll.id"), index=True) label: Mapped[str] = mapped_column(String(512), nullable=False) display_order: Mapped[int] = mapped_column(Integer, default=0) is_active: Mapped[bool] = mapped_column(Boolean, default=True) poll: Mapped["Poll"] = relationship(back_populates="options") class Vote(Base): __tablename__ = "vote" __table_args__ = ( UniqueConstraint("poll_id", "user_id", name="uk_vote_poll_user"), ) id: Mapped[int] = mapped_column(primary_key=True, autoincrement=True) vote_no: Mapped[str] = mapped_column(String(64), unique=True, index=True) poll_id: Mapped[int] = mapped_column(ForeignKey("poll.id"), index=True) option_ids: Mapped[list] = mapped_column(JSON, nullable=False) user_id: Mapped[Optional[str]] = mapped_column(String(64), index=True, default=None) client_ip: Mapped[Optional[str]] = mapped_column(String(45), default=None) ua_hash: Mapped[Optional[str]] = mapped_column(String(64), default=None) created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)模型里几个点值得注意。option_ids 用 JSON 而非字符串拼 ID:读取时直接是列表,写入也避免自己处理分隔符,MySQL 5.7+ 的 JSON 列还能走函数索引。联合唯一约束 uk_vote_poll_user 是“同一个人同一个投票只能投一次”的数据库兜底;匿名投票时 user_id 是 NULL,MySQL 唯一约束允许多个 NULL 并存,正好满足游客参与。created_by 没有建外键,投票系统允许账号体系外的人发起,存一个不依赖用户表的业务标识更省事。
字段长度方面,title 用 String(128) 足够覆盖大多数标题,超出部分要么截断要么在 schema 层报 422;label 放宽到 String(512) 是因为选项里可能出现长文案;client_ip 用 String(45) 是 IPv6 的长度上限。如果业务要求开奖时间精确到毫秒,把 DateTime 换成 DateTime(6),但迁移文件要手动改,autogenerate 不保证生成对。
2.3 Alembic 建表与 pydantic-settings 读取配置
常见做法是项目里用 Alembic 管理表结构变更,而不是在启动时 create_all。原因很简单:投票系统上线后一定会加字段,Alembic 的迁移链可以保留每次变更历史,create_all 做不到。初始化三步:
alembic init alembic alembic revision --autogenerate -m "create poll/option/vote" alembic upgrade head在 alembic/env.py 里把 target_metadata 指向 Base.metadata,Alembic 才能识别模型变化。autogenerate 对 JSON 字段的默认值检测不可靠,生成完打开迁移文件确认 option_ids 是否被标成 nullable=False。如果项目里多个环境连不同库,配置用 pydantic-settings 读 .env:
# app/core/config.py from pydantic_settings import BaseSettings class Settings(BaseSettings): database_url: str = "mysql+pymysql://root:pass@127.0.0.1:3306/vote?charset=utf8mb4" redis_url: str = "redis://127.0.0.1:6379/0" model_config = {"env_file": ".env", "env_file_encoding": "utf-8"} settings = Settings()用 pydantic-settings 而不是直接 os.getenv 的好处是自动读取 .env、自动转类型、支持嵌套配置。依赖安装用 uv sync 还是 pip install -r requirements.txt 都行,关键是 .env 不要提交到仓库。FastAPI 初始化读取配置文件的入口可以这么做:在 main.py 里创建 app 之后,从 core.config 引入 settings,再在 lifespan 里初始化 Redis 连接池,配置读取不会散落在各路由文件里。
3. 创建投票接口:把校验、事务和路由拆开
创建投票是第一个要写的业务接口。最怕把它写成一个肥大函数:请求体校验、数据库操作、返回结果全挤在路由里。FastAPI 项目里我把职责拆成三层,Pydantic 管请求结构,service 管业务规则和事务,路由只做依赖注入和状态码映射。分包清晰之后,后面加一个“复制投票”功能时,service 可以直接复用。这种分层在前后端分离项目实战里比一层到底好维护得多,前端只需要 POST 一个 JSON,其它都交给后端。
3.1 Pydantic Schema 的校验边界有哪些
# app/schemas/poll.py from datetime import datetime from typing import Literal from pydantic import BaseModel, Field, model_validator class PollItemIn(BaseModel): label: str = Field(..., min_length=1, max_length=512) display_order: int = Field(0, ge=0) class CreatePollIn(BaseModel): title: str = Field(..., min_length=1, max_length=128) description: str | None = Field(None, max_length=2000) poll_type: Literal["single", "multi"] = "single" max_options: int = Field(1, ge=1, le=10) start_at: datetime end_at: datetime options: list[PollItemIn] = Field(..., min_length=2, max_length=20) @model_validator(mode="after") def check_dates(self): if self.end_at <= self.start_at: raise ValueError("end_at must be later than start_at") return self这里我关心的不是校验写得对不对,而是校验放在哪一层。日期顺序、选项数量这类业务规则放在 Pydantic 的 model_validator 里,请求到路由前就被拦住,接口返回的是标准的 422,而不是 service 里抛出的 500。max_options 限定 1 到 10,防止多选投票被塞进几十个合法选项,前端下拉框渲染反而被拖慢。poll_type 用 Literal 而不是普通字符串,等于在类型层面就排除了“双选”这种脏数据。
3.2 service 层的 commit 时机与默认值兜底
# app/services/poll_service.py from fastapi import HTTPException from sqlalchemy.orm import Session from app.models.poll import Poll, PollOption from app.schemas.poll import CreatePollIn def create_poll(db: Session, data: CreatePollIn, creator: str | None = None) -> Poll: if data.poll_type == "single": data.max_options = 1 poll = Poll( title=data.title, description=data.description, poll_type=data.poll_type, max_options=data.max_options, status="draft", start_at=data.start_at, end_at=data.end_at, created_by=creator, ) db.add(poll) db.flush() # 拿到自增主键,但事务未提交 for idx, item in enumerate(data.options): db.add(PollOption( poll_id=poll.id, label=item.label, display_order=item.display_order or idx, )) db.commit() db.refresh(poll) return pollservice 里最关键的一步是 db.flush()。flush 会把数据发送到数据库拿到自增主键,但事务并未提交;之后的选项插入和主表共用同一个事务,最后 commit 一次。如果先 commit 主表再插入 option,任一选项失败都会留下一个没有选项的投票活动,这种脏数据在页面上很难排查。单选强制 max_options=1 是后端兜底,前端可能因为旧版本缓存传了 3 过来,但后端必须保证落库的数据合理。投票创建后的 status 固定为 draft,等到后台确认选项无误后再手动改为 running,这给运营留了准备时间。
3.3 挂载路由:CORS 配置与接口调试
# app/api/v1/poll.py from fastapi import APIRouter, Depends, status from sqlalchemy.orm import Session router = APIRouter(prefix="/api/v1/polls", tags=["poll"]) @router.post("", response_model=CreatePollOut, status_code=status.HTTP_201_CREATED) def create_poll_endpoint( data: CreatePollIn, db: Session = Depends(get_db), current_user: str | None = Depends(get_current_user), ): poll = create_poll(db, data, creator=current_user) return poll创建接口用 201 而不是 200;输出 Schema 用 ConfigDict(from_attributes=True) 从 ORM 对象转换,避免把 created_by 这种内部字段透传给前端。前后端分离项目里最常见的调试障碍是 CORS:Vue3 开发服务器跑在 5173,接口在 8000,跨域请求要先过浏览器预检。如果 allow_origins 设置成 ["*"] 又开了 allow_credentials=True,浏览器的 fetch 会直接拦截,接口本身没报错但前端就是拿不到响应。本地开发阶段把 5173 写进 allow_origins 就行:
app.add_middleware( CORSMiddleware, allow_origins=["http://localhost:5173"], allow_credentials=True, allow_methods=["*"], allow_headers=["*"], )如果你是 FastAPI 返回 HTML + Layui 的旧写法,切到纯 API 之后,CORS 和事务边界会突然变得显眼,这一节相当于把最常踩的两个坑提前拆掉。
4. 投一票要过几道关:防重、计数和缓存回写
投票接口和一般 CRUD 最大的不同,是写路径上挤了三个动作:记一条明细、给选项计数、标记已投票。几十个人同时投没问题,几千人同时点一个选项时,连 UPDATE 同一行的锁竞争就能把接口拖到超时。这一章先把写入放大讲清楚,再用 Redis 把防重和计数从数据库里挪出来,最后说定时回写的取舍。后端开发除了增删改查,并发边界才是这个接口真正要处理的业务。
4.1 一次投票动了多少数据:写入放大与瓶颈
一次普通投票如果全走 MySQL,至少要操作三处:vote 新增一行;poll 表的 total_votes 加一;poll_option 表的 votes 加一。三者在一个事务里,数据库的锁范围覆盖了 poll 和 poll_option 两行,并发越高锁等待越明显。更麻烦的是,票数展示接口还要读同一批行,读和写互相挤占连接池。常见做法是把“计数”这种可重算的数据挪到 Redis,投票明细保留在 MySQL 作为事实来源。票数本质是趋势数据,允许一个短暂的滞后窗口,这个语义和 Redis INCR 天然匹配。
| 数据 | 存储 | 写入方式 | 读取方式 |
|---|---|---|---|
| 投票明细 | MySQL | INSERT 追加 | 后台审计/对账 |
| 实时票数 | Redis | INCR | 前端展示/列表页 |
| 已投标记 | Redis | SETNX | 防重复投票 |
4.2 Redis 防重与幂等的最小实现
Redis key 的设计直接决定这个接口能不能扛住压测。防重 key 的过期时间设为到投票结束,活动期间用户只能投一次;幂等 key 独立于防重 key,TTL 设为 300 秒,客户端网络重试时相同 vote_no 直接返回 429。
| key 示例 | 类型 | 过期 | 作用 |
|---|---|---|---|
| vote:user:1:u-10001 | string | 投票结束 | 用户级防重 |
| idem:20240601-uuid | string | 300s | 请求幂等 |
| vote:count:1:100 | string(int) | 投票结束+7天 | 选项实时计数 |
# app/services/vote_service.py import redis from sqlalchemy.orm import Session from app.models.poll import Vote redis_client = redis.Redis.from_url(settings.redis_url, decode_responses=True) def cast_vote(db: Session, poll_id: int, option_ids: list[int], user_id: str, vote_no: str) -> None: user_key = f"vote:user:{poll_id}:{user_id}" ok = redis_client.set(user_key, "1", nx=True, ex=ttl_until_end(poll_id)) if not ok: raise HTTPException(status_code=409, detail="你已经投过票") idem_key = f"idem:{vote_no}" ok = redis_client.set(idem_key, "1", nx=True, ex=300) if not ok: raise HTTPException(status_code=429, detail="重复请求,请稍后重试") pipe = redis_client.pipeline() for oid in option_ids: pipe.incr(f"vote:count:{poll_id}:{oid}") pipe.execute() try: db.add(Vote( vote_no=vote_no, poll_id=poll_id, option_ids=option_ids, user_id=user_id, )) db.commit() except Exception: db.rollback() redis_client.delete(user_key, idem_key) raise在服务启动时用 ConnectionPool 创建 redis_client,避免每个请求新建连接。ttl_until_end 根据 poll 的 end_at 计算剩余秒数,下线前至少要保留 60 秒。投票流程的顺序有讲究:先 SETNX 防重,并发时只有第一个请求能占位成功,其它请求直接拿 409;再对每个选项做 INCR,把写入压力卸到 Redis;最后写 MySQL 明细做持久化。一旦 MySQL 写入失败,要补偿删除前面两个 key,否则用户会“被投票成功”但实际没有记录。这里不需要跨组件事务,Redis 操作可以在 MySQL 失败后反向执行,这个补偿窗口的概率远比“先写 MySQL 后写 Redis 丢计数”低。幂等判断放在防重之后,因为重复请求和已经投过票是两件事:前者是网络重试,后者是业务越权。
路由函数这里不要写成 async def,用普通 def 即可。FastAPI 会把同步路由放进线程池执行,Redis 的阻塞调用不会卡住事件循环;一旦写成 async def,redis_client.set 这种同步调用会把整个进程的事件循环拖停。
4.3 定时回写与缓存一致性
Redis 里的计数在 MySQL 里必须有个归宿,否则 Redis 重启就全丢。常见做法是起一个后台任务,每分钟用 SCAN 扫描 vote:count 前缀,把增量 update 回 poll_option,然后重置为 0。代码大致如下:
# app/tasks/flush_counts.py from sqlalchemy import update from app.models.poll import PollOption def flush_vote_counts(): cursor = "0" while True: cursor, keys = redis_client.scan(cursor=cursor, match="vote:count:*", count=100) for key in keys: _, poll_id, option_id = key.split(":") delta = int(redis_client.get(key) or 0) if delta: db.execute( update(PollOption) .where(PollOption.id == option_id) .values(votes=PollOption.votes + delta) ) redis_client.set(key, 0, ex=ttl) if cursor == "0": break db.commit()为什么用 SCAN 而不是 KEYS:KEYS 在键很多时会把 Redis 单线程阻塞几百毫秒,流量尖峰时这是不能接受的;SCAN 用游标分批取,count=100 表示每批取 100 个键。回写成功后把计数归零而不是删 key,是为了下一轮还能继续累加,同时也让 SCAN 的集合大小保持稳定。归零与数据库更新之间有约 60 秒的窗口,进程在窗口内崩溃会丢失部分秒级增量,投票系统通常接受这个口径;如果产品要求绝对不丢,就要把 vote 明细作为唯一事实来源每小时重放一次明细,代价高很多,活动投票一般不做。键数量很大时,循环里逐 key GET 会有 N 次网络往返,优化方式是先把 keys 收进列表,再走一次 pipeline 批量取回,核心逻辑不变。
提示:Redis 里的计数只用于前端“看起来实时”,报表和最终结果以 MySQL 为准。上线后加一个对账脚本,遍历 MySQL 票数总合和 Redis 剩余计数,偏差超过阈值报警。
5. 票数实时推送与自测:WebSocket 和防重验证
投票页的体验要求是投完马上看到票数涨。轮询每两秒拉一次接口,投票高峰期会产生大量空转请求。FastAPI 原生支持 WebSocket,配合 Redis Pub/Sub 可以实现增量广播:每个浏览器连上专属的 ws 通道,服务端只在有人投票后把增量推给所有订阅者。
# app/api/v1/ws.py import json import redis.asyncio as aioredis from fastapi import APIRouter, WebSocket, WebSocketDisconnect router = APIRouter() redis_client = aioredis.from_url(settings.redis_url) @router.websocket("/ws/polls/{poll_id}/votes") async def poll_votes_ws(websocket: WebSocket, poll_id: int): await websocket.accept() pubsub = redis_client.pubsub() await pubsub.subscribe(f"poll:{poll_id}") try: while True: msg = await pubsub.get_message(ignore_subscribe_messages=True, timeout=30) if msg: await websocket.send_text(msg["data"].decode()) else: await websocket.send_json({"type": "ping"}) except WebSocketDisconnect: await pubsub.unsubscribe(f"poll:{poll_id}") await pubsub.aclose()每个 WebSocket 连接占用一个 pubsub 订阅,连接断开时一定要 aclose,否则 Redis 侧连接会慢慢堆积。广播端在投票接口 INCR 完成后,往对应通道 publish 一个增量 JSON,实例间通过 Redis 转发,多进程部署时所有节点都能收到。注意 Pub/Sub 不保证消息不丢,极端情况下前端要用轮询兜底,ws 断了自动重连并拉一次全量快照。
没有前端页面也能验证防重是否生效。连续发两个相同用户请求,第二次应返回 409:
curl -X POST http://127.0.0.1:8000/api/v1/polls/1/vote \ -H "Content-Type: application/json" \ -H "X-User-ID: u-10001" \ -d '{"option_ids":[1],"vote_no":"v-1001"}' sleep 0.1 curl -X POST http://127.0.0.1:8000/api/v1/polls/1/vote \ -H "Content-Type: application/json" \ -H "X-User-ID: u-10001" \ -d '{"option_ids":[1],"vote_no":"v-1002"}'注意第二个请求的 vote_no 用了新值,这样命中的是用户防重而不是幂等,能分辨 409 来自哪一层。压测用 Locust 时,把 X-User-ID 用固定集合循环复用,而不是每次随机生成,否则整场压测都在测防重 key 的新增,测不出真实选票并发。wait_time 从 0.01 起到 0.05,起步带宽更接近真实用户快速连点。观察两个指标:Redis 的 used_memory 是否平稳,MySQL 的写入 QPS 是否被控在预期水位。压测后如果 Redis 连接超限,把 ConnectionPool 的 max_connections 从 20 调到 50,再把 SCAN 的 count 从 100 调到 500,重新跑一轮看轮询兜底接口的响应时间。
本文还有配套的精品资源,点击获取