ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Python脚本改造FastAPI服务:分层架构与异步重构实战

Python脚本改造FastAPI服务:分层架构与异步重构实战 我把一个跑了大半年的 Python 脚本改成了 FastAPI 服务整个重构过程踩了不少坑也总结了一些可以复用的套路。如果你手里也有一堆“能跑但不好维护”的脚本正琢磨着要不要改成服务、怎么改才不翻车这篇记录应该能帮上忙。先说背景。这个脚本原本是个独立的定时任务负责从几个数据源拉取信息做清洗和聚合最后把结果写进数据库同时在本地生成一份报告。单看功能它完成得挺好问题出在三个地方一是数据源和下游系统的对接越来越多每次加一个对接方都要在脚本里翻来覆去找修改点二是业务方不满足于每天定时跑一次想要按需触发甚至要能传不同参数跑不同逻辑三是报告结果希望能通过接口对外提供方便其他系统直接拉取。这些需求堆在一起脚本模式的扩展成本就变得很高于是决定用 FastAPI 把它重构成一个带分层架构的小服务。1. 从脚本到服务先想清楚“为什么改”和“怎么算成功”1.1 脚本模式的天花板在哪不是所有脚本都需要改服务但如果你正面临下面几种情况就得认真考虑重构了逻辑复用困难脚本里的函数和全局变量纠缠在一起想单独调用某一段逻辑必须理解整个文件的执行顺序甚至得把文件整个 import 进来。多入口需求增多原来只有一个main()现在要支持手动触发、定时触发、外部系统调用每个入口都有自己的参数和返回格式要求脚本里靠argparse兜底已经越来越勉强。无法并发和隔离脚本是长进程里的一段顺序执行处理耗时任务时会阻塞而且不同业务方共用同一份全局状态参数互相污染。可观测性差打印日志、try-except 后继续跑、跑挂了靠 supervisor 重启线上出了什么问题很难定位。如果一条都不占继续用脚本完全没问题没必要为了“上微服务”而上微服务。但一旦占了三条以上改造的收益就很明显了。1.2 为什么选择 FastAPI 而不是 Flask、Django选型这件事我个人的判断标准是能解决眼前的问题同时不给未来留下大坑。对比一下主要几个选项框架异步支持类型提示自动交互文档适合场景综合推荐度Flask仅同步异步需额外配置弱需要集成 flasgger小项目、传统 WSGI 服务一般Django DRF同步为主异步支持较新中等有但集成较重大而全的管理系统偏重FastAPI原生 async/await强基于 Pydantic自带 Swagger UI 和 ReDocAPI 服务、微服务、数据接口高当时手上有两类任务一类是 IO 密集型拉数据、查数据库一类是轻计算聚合、清洗。这两类用async/await都能吃到红利——不是说得多少并发而是同样的代码结构异步版本在查询外部接口时的等待时间可以还给调度器而不是白占一个线程。FastAPI 的异步支持是原生设计写起来很自然。另一个关键点是类型提示。FastAPI 把“入参校验”和“数据模型定义”结合得非常好你用 Pydantic 定义好模型FastAPI 自动帮你做参数校验、类型转换和文档生成。这对从脚本迁移过来的项目尤其宝贵因为脚本时代往往没有接口契约到处是裸字典和魔法数字。再加上 FastAPI 自带 OpenAPI 交互文档后端写好了文档也好了联调时候给前端丢个/docs链接就行省了好大功夫。1.3 重构成功怎么度量重构不是“代码跑通了”就算完事。我会在动手前先列一份检查清单盯着这几件事原有脚本的核心逻辑全部保留输出结果和旧版本一致至少关键字段一致。新的服务接口可以替代原有的手动触发命令且支持传参。新增一个数据源时改动范围只落在少数几个文件内而不是全世界都要动。服务能同时接受多个请求不会因为一个慢查询堵死其他请求。异常能被记录、被感知而不是静默丢弃。也就是说重构的验收标准不是“接口能通”而是“后续迭代的成本有没有降下来”。分层架构的价值恰恰体现在这里。2. 分层架构落地目录、依赖和边界划分2.1 三层架构的设计思路这次重构采用的是最朴素的三层架构没有引入六边形架构、DDD 那些偏重的模式因为项目规模还远没到那个程度。三层分别叫路由层、服务层、数据访问层Data Access Layer, DAL每一层有自己的职责路由层api只负责和 HTTP 打交道。接收请求、返回响应不做业务判断不做数据处理。服务层service承载业务逻辑。编排多个数据源的拉取动作、做清洗和聚合、决定何时入库、何时生成报告。这一层是核心也是可测试性最高的地方。数据访问层repository封装所有和存储相关的操作。包括关系型数据库、文件系统、外部缓存等。分层之后上层不需要知道数据到底存在哪、用的是什么驱动只调用 repository 的方法。这样的分层还有一个额外好处如果你将来要换数据库或者把某个存储换成另一套方案改动可以收敛在 repository 层业务逻辑几乎不用动。2.2 推荐的目录结构很多从脚本改造成服务的项目第一个坑就是目录结构没有规划所有文件堆在一个包里。我这次使用的结构如下可以参考my_service/ ├── app/ │ ├── __init__.py │ ├── main.py # FastAPI 实例入口 │ ├── core/ │ │ ├── config.py # 配置项读取 env / settings │ │ ├── logging.py # 日志配置 │ │ └── exceptions.py # 自定义异常与全局处理器 │ ├── api/ │ │ ├── __init__.py │ │ ├── deps.py # 依赖注入会话、鉴权等 │ │ └── endpoints/ │ │ ├── __init__.py │ │ └── trigger.py # 触发任务接口 │ ├── service/ │ │ ├── __init__.py │ │ ├── task_orchestrator.py # 任务编排服务 │ │ ├── data_cleaner.py # 数据清洗服务 │ │ └── report_generator.py # 报告生成服务 │ ├── repository/ │ │ ├── __init__.py │ │ ├── database.py # 数据库连接与会话管理 │ │ └── result_store.py # 结果数据读写 │ └── models/ │ ├── __init__.py │ ├── api_models.py # 接口入参/出参模型 │ └── data_models.py # 数据实体模型 ├── tests/ │ ├── test_trigger.py │ └── test_cleaner.py ├── scripts/ │ └── run_legacy_check.py # 旧逻辑兼容性校验脚本 ├── requirements.txt └── README.md这里有几个细节特别值得说core/里放的是横切关注点配置、日志、异常定义。这些是所有层都可能用到的但放这里不会破坏依赖方向——依赖始终是从上层指向下层。api/deps.py这一层用来写 FastAPI 的依赖函数比如创建数据库会话、鉴权校验。用 FastAPI 的Depends机制可以把数据库会话的生命周期管理得很好不需要在业务代码里手动处理。models/分为两组api_models是面向接口的请求/响应结构data_models是数据实体的数据类。分开写看起来很啰嗦但在实际迭代中价值巨大。接口结构字段可以精简数据实体字段可以很丰富二者并不总是一一对应的。2.3 边界守则跨层调用怎么约束有了目录还要有纪律否则分层会形同虚设。我给自己定了几条硬规则路由层里不写业务逻辑。路由函数体内只做参数获取、调用服务层、返回响应这三件事。超过五行就当业务逻辑处理拆到 service 去。服务层不直接接触 HTTP 对象。Request、Response对象不应进入 service 层。入参是普通的 Python 数据类或 Pydantic 模型出参是普通对象。repository 层不向上抛细节异常。数据库连接断开、表锁冲突这类底层异常按场景包装为领域异常或统一格式的异常由全局异常处理器统一转成对应的 HTTP 状态码和响应体。依赖方向单向。api - service - repository反向的 import 一出现就说明边界破了需要重新审视设计。这些规则听起来像老生常谈但真正能从脚本思维转过来的人不多。脚本里的函数天然是“从上到下顺序执行”全局变量多模块之间动不动就互相引用分层架构要求的是“单向依赖、逐层隔离”这一步迈不过去后面所有重构都会白做。3. 核心重构实操一步步把脚本拆开再织起来3.1 第一步先把“主流程”换成“用例”老脚本长这样简化演示def main(): raw_list [] for source in SOURCES: data fetch_from_source(source) raw_list.extend(data) cleaned [] for item in raw_list: if item.get(score, 0) 60: cleaned.append(clean_item(item)) save_to_db(cleaned) generate_report(cleaned) print(fdone, total {len(cleaned)})重构第一步不是写类而是把这个主流程中每一步拆成独立的“用例”。所谓用例就是一个动词场景比如fetch_latest_data、clean_raw_data、store_results、generate_report。这一步的目的是确定职责归属fetch_from_source属于数据访问层吗不完全。因为“从哪些源拉取、拉什么格式”是业务决策但“怎么建立连接、怎么发请求”是技术细节。所以我拆成两层service 里定义data_source_manager.py负责知道“要拉哪几个源、参数是什么”repository 里的external_api_client.py负责具体的 HTTP 请求和重试。clean_item放在 service 层没问题但要注意它应该是纯函数输入一个数据项输出一个清洗后的数据项不做 IO、不读全局状态。这样方便单测。save_to_db显然属于 repository 层。generate_report严格来说既有业务决策生成什么内容、什么格式又有技术操作写文件、渲染模板我把它拆成service/report_generator.py决定报告内容和repository/report_storage.py负责写文件或上传对象存储。拆的过程别急着写实现先在纸上或思维导图里把每个函数从“紧挨着调用”变成“各归各位”。这一步做完你已经完成了 60% 的重构工作量。剩下的就是把函数组装起来。3.2 第二步定义清晰的调用链服务层是核心我以一个顶层用例run_task举例看它是怎么把多个环节编排起来的# service/task_orchestrator.py from typing import List from models.data_models import RawItem, CleanedItem from service.data_cleaner import DataCleaner from service.report_generator import ReportGenerator from repository.external_api_client import ExternalApiClient from repository.result_store import ResultStore class TaskOrchestrator: def __init__( self, cleaner: DataCleaner, report_gen: ReportGenerator, api_client: ExternalApiClient, store: ResultStore, ): self._cleaner cleaner self._report_gen report_gen self._api_client api_client self._store store async def run(self, source_names: List[str]) - dict: # 1. 并行拉取多个源这里才是真正的 async 收益点 raw_items: List[RawItem] await self._api_client.fetch_many(source_names) # 2. 清洗 cleaned: List[CleanedItem] self._cleaner.clean_batch(raw_items) # 3. 存储 stored_count await self._store.save_cleaned(cleaned) # 4. 生成报告 report_path await self._report_gen.generate(cleaned) # 返回值给上层路由转成接口响应 return { stored_count: stored_count, report_path: report_path, }为什么要在构造函数里把依赖传进来而不是在 service 里直接import一句话回答为了测试的时候能替换真实依赖为 Mock 对象。这是我重构过程中收益最大的一点。旧脚本时代要测试clean_item就必须确保fetch_from_source真的能连上外网现在用unittest.mock或者pytest的 monkeypatch 就能轻松替换class FakeApiClient: async def fetch_many(self, sources): return [RawItem(sourcea, content...)]这个替换能力是分层架构带给你最直接的“投资回报”。3.3 第三步用路由层暴露能力FastAPI 的接口定义非常简洁。以一个触发任务的接口为例# api/endpoints/trigger.py from fastapi import APIRouter, Depends, HTTPException from models.api_models import RunTaskRequest, RunTaskResponse from service.task_orchestrator import TaskOrchestrator from api.deps import get_orchestrator router APIRouter(prefix/api/task, tags[task]) router.post(/run, response_modelRunTaskResponse) async def run_task( payload: RunTaskRequest, orchestrator: TaskOrchestrator Depends(get_orchestrator), ): try: result await orchestrator.run(source_namespayload.sources) return RunTaskResponse(**result) except Exception as e: # 全局异常也会兜底这里可以加更细的业务日志 raise HTTPException(status_code500, detailstr(e))入口本身不关心TaskOrchestrator内部拉数据用了几个线程、DB 用的是 PostgreSQL 还是 SQLite只管收请求、调服务、返结果。依赖注入的写法要注意get_orchestrator放在api/deps.py里返回的是一个组装好的TaskOrchestrator实例。实例里面依赖的 Repository 对象可以通过 FastAPI 的lifespan在应用启动时统一初始化也可以在每个请求级别创建会话看业务需要。我建议先不要过度设计。打个比方你手里有一堆零件先按它们的功能分类放进不同的抽屉分层然后找一个“流水线班长”TaskOrchestrator把它们按顺序组织起来最后给班长配一个“前台接待”路由层跟外部沟通。这个模型简单清晰够用就好不用一上来就上依赖注入框架。3.4 第四步迁移配置与启动逻辑脚本时代配置散落在各种环境变量和config.py的全局变量里。我重构时用 Pydantic 的BaseSettings统一管理配置# core/config.py from pydantic_settings import BaseSettings class Settings(BaseSettings): db_url: str source_api_keys: dict[str, str] {} report_output_dir: str ./reports log_level: str INFO class Config: env_prefix MY_SERVICE_ env_file .env settings Settings()优势在于环境变量、.env 文件、默认值三者结合部署时不用改代码。类型校验在加载时就完成避免了脚本时代“配置项写成字符串运行时才炸”的问题。维护成本很低加新配置只需要在 Settings 里加一个带默认值的字段。启动逻辑放到main.py用 FastAPI 的lifespan管理应用生命周期from contextlib import asynccontextmanager from fastapi import FastAPI asynccontextmanager async def lifespan(app: FastAPI): # 启动时初始化 Repository 实例 await init_async_db() yield # 关闭时释放资源 await dispose_async_db() app FastAPI(titleTask Service, version1.0.0, lifespanlifespan) app.include_router(trigger.router)这里有几个小细节数据库连接如果用的是 SQLAlchemy 2.0 的 async 引擎需要在 lifespan 里调用await engine.dispose()否则进程退出时可能报“Task was destroyed but it is pending”。有些异步客户端如httpx.AsyncClient也需要在退出时 close可以把它们的实例挂在某个对象下统一管理别在函数里反复创建。4. 分层重构的常见问题和排查心得4.1 异步改造的连锁反应脚本改成服务最大的隐性成本是把 IO 操作都改成异步。requests.get()换成httpx.AsyncClient.get()psycopg2的同步连接换成asyncpg或 SQLAlchemy 的 async 引擎文件读取用aiofiles。我在这个阶段踩过一个大坑外部 API 的客户端只支持同步比如某些 SDK。这时不要硬改可以在anyio.to_thread.run_sync里包装一下。但这里有个注意点并发量大的时候to_thread 会占用线程池资源如果线程不够性能反而下降。所以要么把这种调用限定为低频场景要么换一个真的支持异步的客户端库。另外async/await有“传染性”一个函数要是 async调用它的函数也得是 async一路传染到路由层。这既是压力也是分层的机遇——如果你把“数据拉取”合理地封装在 repository 层那么只有这一层需要关心同步/异步的转换service 层只盯着接口用就好。4.2 依赖注入最常见的坑循环导入分层架构里api/deps.py依赖serviceservice依赖repository。如果谁手一滑在 repository 里反向 import 了 service 的某个类循环导入就来了。症状是启动时报ImportError: cannot import name ... from partially initialized module。排查方式很简单python -c from app.main import app报错信息会给出首次 warning 的位置直接用 IDE 按引用关系搜一下是哪个模块在启动时反向引用。规避办法所有跨层依赖只通过构造函数注入不要在模块顶部互相 import 对方的具体实现。把接口抽象放单独目录比如repository/base.py里放抽象基类service 层依赖“抽象”而不是“具体实现”。4.3 异常处理要统一否则接口直接裸奔脚本时代无所谓异常处理报错了无非是 console 里打一段堆栈。但作为服务异常如果不统一用户拿到的响应就是 200 OK 一段 Python 异常文本或者 500 里什么都没有排查全靠猜。我的做法是三步走第一定义业务异常基类# core/exceptions.py class AppError(Exception): status_code 500 code internal_error message Internal Server Error class DataSourceError(AppError): status_code 502 code data_source_error message Failed to fetch data from source第二在 FastAPI 里注册全局异常处理器from fastapi import FastAPI, Request from fastapi.responses import JSONResponse from core.exceptions import AppError app.add_exception_handler(AppError, app_error_handler) async def app_error_handler(request: Request, exc: AppError): return JSONResponse( status_codeexc.status_code, content{code: exc.code, message: exc.message}, )第三service 内部把底层异常转换成业务异常try: await self._api_client.fetch_many(sources) except ResourceReadTimeoutError as e: raise DataSourceError(上游接口超时) from e这样用户拿到的是结构化的错误信息而完整堆栈记录在我们的日志里两不耽误。4.4 旧脚本还能不能留怎么共存重构过程中不要让旧脚本突然消失。比较稳的方案是“双轨运行”服务先上线接口先提供旧脚本留着做对账。对账怎么做我写了一个scripts/run_legacy_check.py它的逻辑很简单用旧脚本跑一遍当天任务再调新服务的接口跑一遍比对两份输出文件的关键字段。比对通过后再把定时调度切到新服务的接口上。这个对账脚本只存在于过渡期但它给了项目组一个“安全网”。哪怕新服务上线了某一天发现了数据不一致的问题我们可以立刻定位是哪一层的问题而不会被业务方质疑“是不是改了之后变坏了”。4.5 日志和链路追踪不要等到上线后才补脚本时代的 print 输出到了服务阶段得升级成标准结构化日志。我在core/logging.py里用标准库logging.config.dictConfig配置了 JSON 格式的输出每行日志包含时间、级别、logger 名称、事件和关键字段。一个实用技巧是在路由层给每个请求生成一个request_id用contextvars管理在日志里打出来。这样即便用户反馈“某个任务失败了”你也只需要搜索那一个 request_id就能看到这个请求完整经过了哪些层。import contextvars _request_id_var: contextvars.ContextVar contextvars.ContextVar(request_id, default) def get_request_id(): return _request_id_var.get()结合 FastAPI 的中间件在请求进来时生成 ID响应时输出到 header前端联调时也能用同一个 ID 在后端日志里检索。5. 测试策略、性能优化和部署细节5.1 用 pytest 搭一个分层测试框架脚本时代几乎无测试可言。重构后我补了三类用例单元测试针对 service 层的纯函数清洗、聚合直接给定输入断言输出。接口测试用 FastAPI 的TestClient配合 Fake Repository 返回假数据验证路由层的响应结构、状态码、异常格式。集成测试连接真实的测试数据库或测试桩服务跑完整流程。其中接口测试是最能体现分层价值的。因为依赖通过构造函数注入测试时只需要传入 fake 实现不需要真的连数据库和外部 API。如果服务里写死了直连数据库的代码这类测试基本没法写。# tests/test_trigger.py def test_trigger_run_success(): fake_orch FakeOrchestrator(result{stored_count: 10, report_path: /tmp/x}) app.dependency_overrides[get_orchestrator] lambda: fake_orch with TestClient(app) as client: resp client.post(/api/task/run, json{sources: [a]}) assert resp.status_code 200 assert resp.json() {stored_count: 10, report_path: /tmp/x} app.dependency_overrides.clear()5.2 并发性能的实测数据我实测过一个小场景同时触发 10 个任务每个任务拉取 3 个源清洗 5000 条数据写库。同步版本耗时约 210 秒异步版本约 32 秒。主要提升在于外部请求的并发等待时间被释放了瓶颈从“IO 串行等待”转移到“数据库写入并发”。当然这样的提升有条件数据库连接池要够用SQLAlchemy 的pool_size和max_overflow要根据并发请求数设置否则会大量报“connection pool exhausted”。外部 API 有频率限制的话并发反而会触发 429需要在上游客户端里做限流。如果你的任务本身是 CPU 密集型的比如大量正则跑在几十万条文本上放 async 里也没用这时需要用多进程或者把任务切出去单独跑。5.3 部署方式Docker systemd 的一点点经验FastAPI 服务部署首选 Docker。一个精简的 DockerfileFROM python:3.11-slim WORKDIR /app COPY requirements.txt . RUN pip install --no-cache-dir -r requirements.txt COPY . . EXPOSE 8000 CMD [uvicorn, app.main:app, --host, 0.0.0.0, --port, 8000]注意 uvicorn 的 worker 数量单机部署随口配 4 个 worker 不一定是好事。如果服务里大量用了异步 IO单 worker 往往就能打满 CPU 的作用不大如果业务不是 IO 密集而是计算密集多 worker 才有价值。我一般建议先从 2 个 worker 起步压测后看曲线再调不要一上来就复制粘贴网上的四 worker 配置。另外--reload只用于开发环境生产环境千万不要开。这个方向不用太纠结但如果你真在公网环境开着 reload 跑日志里会不断出现 “WatchFiles detected changes, reloading”那多半是被外部触发了文件变更非常危险。5.4 健康检查和优雅退出服务上线后一定要加一个健康检查接口router.get(/healthz) async def healthz(): return {status: ok}健康检查要真实反映依赖状态如果数据库连不上这个接口必须返回非 200。否则调度器以为服务健康结果任务全部失败用户直到数据缺失才感知到问题。优雅退出也值得配置。Docker 里 uvicorn 会监听 SIGTERM停止接收新连接等处理中的请求完成后退出。为了让这个“处理中的请求”不无限拖延我建议为数据库客户端和外部 API 客户端都设置超时时间一般 30 秒足够。6. 重构全过程的节奏与方法论小结这次重构大概分四步走梳理现状确定边界把脚本里的主流程画成功能清单标出哪些是业务逻辑、哪些是 IO 细节、哪些是配置。搭建骨架先不还债先把目录建好、模型定义好、配置迁移好再用最原始的方式让路由能调到 serviceservice 能调到 repository哪怕 repository 里先直接返回假数据。逐层迁移边迁移边测试先把纯函数逻辑清洗、聚合搬过去再搬外部接口和数据存储每次迁移完跑一遍对应用例对照输出。上线切换保留后备旧脚本保留一份服务灰度运行用对账脚本验证输出一致后再切换调度。我个人这几年做脚本改造项目下来最大的体会是重构的核心不是“代码写得更漂亮”而是“让下一次需求变更的成本显著下降”。如果你改了架构下一次加需求却还是到处打补丁那这个重构就是表面工程没意义。还想提醒一句不要试图在重构中同时解决所有问题——顺便修 bug、顺便换数据库、顺便统一历史遗留的设计缺陷——一次只做一件事。把“重构”和“修 bug”混在一起出问题了你根本不知道是重构引入的还是原有 bug 暴露的一次只做一件事你永远知道当前风险在哪。另外提一个小技巧在拆函数的过程中我会用一个文本文件记录“函数前后对应关系”例如“旧脚本的clean_item→ 新服务的DataCleaner.clean_one”。这个表看着无聊但和旧脚本对账的时候它就是最快的索引。等到全部迁移完成、对账通过后再销毁顺手还能当迁移文档用。把脚本改造成 FastAPI 服务的路径没有唯一标准答案。只要是“依赖清晰、边界完整、测试可跑、上线可查”就已经赢了。如果你正在做类似的改造希望这篇记录能帮你避开几个坑等你的服务跑起来之后再回看原来的脚本多半会感慨一句——早该这么干了。
返回列表