17 KiB
17 KiB
FastAPI + Docker 迁移设计文档
日期: 2026-05-09 版本: 1.0 状态: 设计阶段
1. 背景与目标
将现有的广西政府采购网公告监控系统从 Flask CLI 架构迁移到 FastAPI + Docker 部署方案,同时提升代码质量、安全性和可维护性。
核心目标
- FastAPI 替代 Flask 回调服务器 + argparse CLI,统一为一个 Web 应用
- Docker 单容器部署,连接现有 PostgreSQL
- SQLAlchemy ORM 替代原始 psycopg2 操作
- APScheduler 内置定时任务
- 环境变量管理配置,敏感信息不提交 git
- 异步改造(httpx + asyncpg)
2. 新项目结构
gx-gp-notify/
├── app/
│ ├── __init__.py
│ ├── main.py # FastAPI 应用入口 + lifespan
│ ├── config.py # Pydantic Settings (环境变量驱动)
│ ├── api/
│ │ ├── __init__.py
│ │ ├── router.py # 统一路由注册
│ │ ├── deps.py # 依赖注入 (get_db, get_scheduler)
│ │ ├── announcements.py # 公告 CRUD 端点
│ │ ├── crawl.py # 爬取触发端点
│ │ └── wechat.py # 企业微信回调端点
│ ├── models/
│ │ ├── __init__.py
│ │ ├── announcement.py # SQLAlchemy 模型
│ │ └── schemas.py # Pydantic 请求/响应模型
│ ├── crawler/
│ │ ├── __init__.py
│ │ ├── base.py # Spider 抽象基类 + Pipeline 配置
│ │ ├── gxgp_spider.py # 广西政府采购网爬虫 (API 模式)
│ │ ├── dahuagov_spider.py # 大化县政府网爬虫 (HTML 模式)
│ │ └── parsers.py # 数据解析器
│ ├── services/
│ │ ├── __init__.py
│ │ ├── crawl_service.py # 爬取编排器 (注册/调用所有 Spider)
│ │ ├── pipeline.py # 统一后处理管道 (存储→筛选→推送)
│ │ ├── filter_service.py # 筛选逻辑
│ │ └── notification_service.py # 通知服务
│ ├── wechat/
│ │ ├── __init__.py
│ │ ├── handler.py # 消息处理器
│ │ ├── menu.py # 菜单管理
│ │ ├── crypto.py # WXBizMsgCrypt 加解密
│ │ └── client.py # 企业微信 API 客户端
│ └── scheduler/
│ ├── __init__.py
│ └── jobs.py # APScheduler 定时任务
├── tests/
│ ├── __init__.py
│ ├── conftest.py # pytest fixtures
│ ├── test_crawler/ # 爬虫解析器测试
│ ├── test_api/ # API 端点测试
│ └── test_services/ # 业务逻辑测试
├── alembic/ # 数据库迁移
│ └── versions/
├── docker/
│ ├── Dockerfile
│ └── docker-compose.yml
├── pyproject.toml # 项目元数据 + 依赖 + 工具配置
├── .env.example # 环境变量模板
├── .gitignore
└── logs/ # 日志目录 (gitignore)
3. 爬虫模块化设计
3.1 设计动机
系统爬取两个不同来源的公告:广西政府采购网(API 模式,13 种公告类型)和大化县政府网(HTML 解析,单一来源)。两个爬虫的后处理流程不同:
| 对比 | 广西政府采购网 | 大化县政府网 |
|---|---|---|
| 爬取方式 | POST JSON API | GET HTML + BeautifulSoup |
| 反爬机制 | 敏感词检查 | 无 |
| 分页 | 多页轮询 | 单页 |
| 筛选策略 | 关键词 + 日期 + 去重 | 不过滤 |
| 推送策略 | 只推送关键词匹配的 | 全部推送 |
| 已发送标记 | is_new 标记 | is_new → mark_sent |
为避免两个爬虫各自维护一套存储/筛选/推送逻辑,采用 Spider 接口 + Pipeline 策略模式 分离关注点。
3.2 架构
┌─────────────────────────────────────────────────┐
│ CrawlOrchestrator │
│ (注册所有 Spider,统一调度爬取) │
└─────────────────────┬───────────────────────────┘
│
┌─────────────┼─────────────┐
│ │ │
┌───────▼──────┐ ┌────▼──────┐ ┌────▼──────────┐
│ GXGPSpider │ │ Dahuagov │ │ (FutureSpider) │
│ (API 模式) │ │ Spider │ │ │
│ 13 sources │ │ (HTML模式) │ │ │
└───────┬──────┘ └────┬──────┘ └────┬───────────┘
│ │ │
└─────────────┼─────────────┘
│ List[Announcement]
│
┌─────────────────────▼───────────────────────────┐
│ PostCrawlPipeline │
│ │
│ 每个 Spider 声明 PipelineConfig: │
│ - filter_enabled: 是否启用关键词筛选 │
│ - keywords: 筛选关键词列表 │
│ - dedup_enabled: 是否去重 │
│ - notify_mode: "filtered" | "all" │
│ - mark_sent: 推送后是否标记已发送 │
│ │
│ 统一处理流程: │
│ 存储 → 去重 → 筛选(可选) → 推送(可选) → 标记(可选) │
└──────────────────────────────────────────────────┘
3.3 Spider 基类
from abc import ABC, abstractmethod
from dataclasses import dataclass
from typing import List
@dataclass
class PipelineConfig:
"""Spider 后处理策略配置"""
filter_enabled: bool = True
keywords: List[str] = None
dedup_enabled: bool = True
notify_mode: str = "filtered" # "filtered" | "all"
mark_sent: bool = False
class BaseSpider(ABC):
"""爬虫基类 — 只负责爬取+解析,不管后续如何处理"""
name: str # spider 名称
source_code: str # 来源代码
source_name: str # 来源名称
@abstractmethod
async def crawl(self) -> CrawlResult:
"""执行爬取,返回包含 Announcement 列表的 CrawlResult"""
...
def get_pipeline_config(self) -> PipelineConfig:
"""返回后处理策略(子类可覆盖)"""
return PipelineConfig()
3.4 两个 Spider 的 PipelineConfig
# GXGPSpider — 关键词筛选模式
PipelineConfig(
filter_enabled=True,
keywords=["大化"], # 从配置读取
dedup_enabled=True,
notify_mode="filtered", # 只推送匹配关键词的
mark_sent=False, # 用 is_new 标记,不单独 mark_sent
)
# DahuagovSpider — 全量推送模式
PipelineConfig(
filter_enabled=False,
keywords=[],
dedup_enabled=True,
notify_mode="all", # 全部推送
mark_sent=True, # 推送后标记已发送
)
3.5 统一后处理管道
class PostCrawlPipeline:
"""统一后处理管道 — 所有 Spider 共用"""
async def process(self, announcements: List[Announcement],
config: PipelineConfig) -> PipelineResult:
# 1. 统一存入 announcements 表 (按 source_code 区分)
# 2. 去重 (基于 content_hash)
# 3. 根据 config.filter_enabled 决定是否筛选
# 4. 根据 config.notify_mode 决定推送策略
# - "all": 推送全部新公告
# - "filtered": 只推送匹配关键词的
# 5. 根据 config.mark_sent 标记已发送
# 6. 生成 Markdown (可选)
...
3.6 优势
- 新增爬虫源只需实现
BaseSpider.crawl()+ 声明PipelineConfig,推送/存储逻辑零改动 - 推送逻辑单一入口,修改一处对所有源生效
- 每个 Spider 文件只关注"怎么爬"和"怎么解析",不再包含存储/筛选/推送代码
dahuagov_announcements独立表可合并到统一的announcements表,用source_code='dahuagov'区分
4. API 端点设计
3.1 公告查询
| 方法 | 路径 | 说明 |
|---|---|---|
GET |
/api/v1/announcements |
查询公告列表(分页、来源筛选、日期筛选、关键词搜索) |
GET |
/api/v1/announcements/{id} |
获取单条公告详情 |
GET |
/api/v1/announcements/today |
获取今日公告 |
GET |
/api/v1/announcements/stats |
获取公告统计信息 |
3.2 爬取控制
| 方法 | 路径 | 说明 |
|---|---|---|
POST |
/api/v1/crawl/trigger |
手动触发一次爬取(支持关键词/来源参数) |
GET |
/api/v1/crawl/status |
查看最近爬取状态 |
GET |
/api/v1/crawl/sources |
获取已配置的公告来源列表 |
3.3 定时任务
| 方法 | 路径 | 说明 |
|---|---|---|
GET |
/api/v1/scheduler/jobs |
查看活跃的定时任务 |
POST |
/api/v1/scheduler/pause/{job_id} |
暂停某个定时任务 |
POST |
/api/v1/scheduler/resume/{job_id} |
恢复某个定时任务 |
3.4 企业微信
| 方法 | 路径 | 说明 |
|---|---|---|
GET POST |
/api/v1/wechat/callback |
企业微信回调入口(URL验证 + 消息接收) |
3.5 系统
| 方法 | 路径 | 说明 |
|---|---|---|
GET |
/health |
健康检查(包含数据库连通性) |
3.6 通用特性
- FastAPI 自动生成 Swagger UI (
/docs) 和 ReDoc (/redoc) - 统一分页格式:
{ total, page, page_size, items } - 统一错误响应:
{ detail: "error message" } - 依赖注入:数据库 session、配置、调度器均通过 FastAPI
Depends()注入 - CORS 默认关闭,通过
CORS_ORIGINS环境变量可选开启
5. 数据模型
5.1 公告表 (announcements)
合并现有的 announcements、auto_announcements、manual_announcements、dahuagov_announcements 四张表为一张,通过 source_code 字段区分来源,通过 PipelineConfig 控制每个来源的推送策略。
| 字段 | 类型 | 说明 |
|---|---|---|
| id | Integer | 主键 |
| title | String(500) | 公告标题 |
| publish_date | DateTime | 发布时间 |
| purchase_name | String(200) | 发布单位 |
| content_url | Text | 内容链接 |
| source_code | String(50) | 来源代码(ZcyAnnouncement1 / dahuagov 等) |
| source_name | String(100) | 来源名称 |
| announcement_type | String(50) | 公告类型枚举值 |
| content_hash | String(64) | SHA256 内容哈希,UNIQUE INDEX |
| crawl_mode | String(20) | "auto" / "manual" |
| is_new | Boolean | 是否新公告 |
| is_sent | Boolean | 是否已推送 |
| keyword_matched | Boolean | 是否匹配关键词 |
| created_at | DateTime | 创建时间 |
| updated_at | DateTime | 更新时间 |
5.2 设计决策
- 四表合一:不再为不同来源或抓取模式创建独立表。
source_code区分来源(ZcyAnnouncement1~dahuagov),crawl_mode区分自动/手动 - is_sent 字段:替代原来大化县的
is_new → mark_sent模式,同时为广西采购网的公告提供统一的已推送标记 - keyword_matched 字段:在筛选阶段标记,持久化到数据库方便后续查询
5.3 变更说明
content_hash从 MD5 (32位) 改为 SHA256 (64位)- 去掉
announcement_sources表,来源信息从环境变量配置动态读取 - 去掉
crawl_results表,爬取结果通过日志记录 - 使用 Alembic 管理数据库迁移
- 旧数据迁移:通过 Alembic migration 脚本将现有四张表的数据合并到新表
5. 配置管理
5.1 pydantic-settings
所有配置通过 pydantic-settings 的 BaseSettings 加载,来源优先级:环境变量 > .env 文件 > 默认值。
5.2 环境变量
# 应用
DEBUG=false
LOG_LEVEL=INFO
# 数据库(连接现有 PostgreSQL)
DATABASE_URL=postgresql+asyncpg://gx-gp-notify:password@10.10.10.14:5432/gx-gp-notify
# 爬虫
CRAWLER_BASE_URL=https://zfcg.gxzf.gov.cn
CRAWLER_KEYWORDS=["大化"]
CRAWLER_MAX_PAGES=10
CRAWLER_TIMEOUT=30
# 企业微信
WECHAT_ENABLED=true
WECHAT_CORP_ID=ww69e8e44636f47780
WECHAT_AGENT_ID=1000007
WECHAT_SECRET=<secret>
WECHAT_TOKEN=<token>
WECHAT_ENCODING_AES_KEY=<aes_key>
# 定时任务
SCHEDULER_ENABLED=true
SCHEDULER_CRON=0 8,14,18 * * *
# Markdown 输出
MARKDOWN_OUTPUT_FILE=onu.md
5.3 安全措施
.env文件加入.gitignoreconfig.yaml中的敏感信息不再提交到仓库- 使用 SHA256 替代 MD5 做内容哈希
- 建议清理 git 历史中的敏感信息
6. Docker 部署
6.1 架构
docker-compose.yml
┌─────────────────────────────┐
│ app (FastAPI) │
│ - uvicorn 服务器 │
│ - APScheduler (内置) │
│ - 端口 8000 │
│ - 挂载: ./logs │
└──────────┬──────────────────┘
│ TCP 连接
┌──────────▼──────────────────┐
│ 外部 PostgreSQL (现有) │
│ 10.10.10.14:5432 │
└─────────────────────────────┘
6.2 Dockerfile
FROM python:3.12-slim
WORKDIR /app
COPY pyproject.toml .
RUN pip install --no-cache-dir .
COPY . .
CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000"]
6.3 docker-compose.yml
services:
app:
build: .
ports:
- "8000:8000"
env_file:
- .env
volumes:
- ./logs:/app/logs
restart: unless-stopped
6.4 启动步骤
cp .env.example .env # 编辑 .env 填入真实配置
docker compose up -d # 启动服务
docker compose logs -f app # 查看日志
7. 技术栈
| 组件 | 当前 | 迁移后 |
|---|---|---|
| Web 框架 | Flask (仅回调用) | FastAPI |
| CLI | argparse | 不再需要,API 替代 |
| 数据库驱动 | psycopg2 (同步) | asyncpg (异步) + SQLAlchemy 2.0 |
| HTTP 客户端 | requests (同步) | httpx (异步) |
| 定时任务 | schedule 库 | APScheduler |
| 配置 | PyYAML + config.yaml | pydantic-settings + .env |
| 加密 | pycryptodome + cryptography | 保留 |
| 部署 | 裸机 + uWSGI | Docker + uvicorn |
| 代码检查 | 无 | ruff |
| 测试 | 无 | pytest + httpx |
| 迁移 | 手动 CREATE TABLE | Alembic |
8. 额外改进
8.1 异步改造
- 爬虫网络请求:
httpx.AsyncClient替代requests - 数据库:
asyncpg+ SQLAlchemy async engine - 通知发送:
httpx.AsyncClient - API 端点使用
async def
8.2 自动化测试
- 爬虫解析器单元测试(Mock 网页响应,验证解析逻辑)
- 筛选逻辑单元测试
- API 端点集成测试(用
httpx.AsyncClient+ 测试数据库) - pytest + pytest-asyncio + pytest-cov
8.3 日志改进
- 结构化日志(JSON 格式)
- Docker 环境下输出到 stdout
- 保留文件日志供宿主机查看(挂载
./logs目录)
8.4 错误处理
- 统一异常处理中间件
- 健康检查
/health含数据库连通性验证 - 关键操作结构化日志记录
8.5 代码质量
ruff替代 flake8 + blackpyproject.toml统一管理依赖和工具配置- 全面 Type Hints
8.6 .gitignore
.env
logs/
*.log
__pycache__/
*.pyc
.venv/
.ruff_cache/
.pytest_cache/
*.egg-info/
9. 迁移策略
分步执行
- 创建新项目结构骨架(FastAPI app + config)
- 数据库模型 + Alembic 迁移
- 迁移爬虫模块(异步改造)
- 迁移服务层(筛选、通知、Markdown)
- 迁移 WeChat 回调(Flask → FastAPI route)
- 添加 APScheduler 定时任务
- 编写测试
- Docker 化
- 清理旧文件
向后兼容
- 保留
dahuagov_announcements表结构不变 - 企业微信回调 URL 路径保持不变 (
/api/v1/wechat/callback) - 数据库迁移使用 Alembic,不丢失现有数据
10. 未包含的内容
- 不做多用户认证/授权(个人使用)
- 不做 Celery 分布式任务(个人使用,APScheduler 足够)
- 不做 Redis 缓存
- 不做前端 UI
11. 风险与缓解
| 风险 | 缓解 |
|---|---|
| 数据库迁移丢失数据 | Alembic 自动生成迁移脚本,先在测试环境验证 |
| 异步爬虫被目标站限流 | 保留现有延迟/重试机制,httpx 支持同样的超时配置 |
| 企业微信回调兼容性 | 回调路径不变,加解密逻辑完全复用现有代码 |
| Docker 网络访问 10.10.10.14 | 确认 Docker 宿主机可访问该 IP,必要时用 host.docker.internal |