新闻早报生成系统 - 运作原理详解
本文档深入解释新闻早报生成系统的内部工作机制,帮助你理解系统是如何从原始新闻数据到最终推送到飞书群聊的完整流程。
一、系统架构概览
┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐
│ 外部数据源 │ │ 外部 AI 服务 │ │ 外部推送服务 │
│ ┌───────────┐ │ │ ┌───────────┐ │ │ ┌───────────┐ │
│ │ CBS News │ │ │ │ 通义千问 │ │ │ │ 飞书 API │ │
│ │ RSS源 │ │ │ │(DashScope)│ │ │ │ │ │
│ ├───────────┤ │ │ └───────────┘ │ │ └───────────┘ │
│ │ 中新网 │ │ └─────────────────┘ └─────────────────┘
│ │ RSS源 │ │
│ └───────────┘ │
└────────┬────────┘
│
▼
┌─────────────────────────────────────────────────────────────────────────────┐
│ 数据获取层 (get_news.py) │
│ - 异步并发拉取多个 RSS 源 │
│ - 使用 trafilatura 提取文章正文 │
│ - 调用 db.save_news() 保存到数据库 │
└────────────────────────┬────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────────────────┐
│ 数据存储层 (db.py) │
│ - PostgreSQL 数据持久化 │
│ - URL 去重,防止重复存储 │
│ - 异步后台任务触发摘要生成 │
└────────────────────────┬────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────────────────┐
│ AI 处理层 │
│ ┌─────────────────────────────┐ ┌─────────────────────────────────────┐ │
│ │ 摘要生成 (summarize.py) │ │ 新闻筛选 (generate_news.py) │ │
│ │ - 调用通义千问 qwen-flash │ │ - 调用通义千问 qwen-plus │ │
│ │ - 生成 50-100 字中文摘要 │ │ - 从摘要中选出最重要的 10 条 │ │
│ │ - 异步批量处理 │ │ - 返回新闻 ID 列表 │ │
│ └─────────────────────────────┘ └─────────────────────────────────────┘ │
└────────────────────────┬────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────────────────┐
│ 消息推送层 (send.py) │
│ - 使用飞书开放平台 API │
│ - 发送 Markdown 格式的富文本消息 │
└─────────────────────────────────────────────────────────────────────────────┘
二、完整数据流程
系统的工作流程可以分为 6 个主要阶段:
阶段 1: 定时触发
┌─────────────┐
│ 定时触发 │
│ (每天 8:00) │
└──────┬──────┘
│
▼
┌────────────────────────────────────┐
│ main.py::job() │
│ - 这是整个系统的入口函数 │
│ - 由 APScheduler 每日定时调用 │
└────────────────────────────────────┘
技术细节:
- 使用
BlockingScheduler阻塞式调度器 - 配置为每天早上 8:00 执行
- Docker 容器启动后自动开始调度
阶段 2: 数据获取
┌─────────────────────────────────────────────────────────────┐
│ get_news.py::get_save_news() │
│ │
│ 1. 读取 config.toml 中的 RSS 源列表 │
│ 2. 使用 httpx.AsyncClient 并发请求所有 RSS 源 (异步并行) │
│ 3. feedparser 解析 RSS 响应,提取标题、链接、发布时间 │
│ 4. 按 URL 去重合并所有新闻 │
│ 5. 对每个新闻: │
│ - 检查是否已存在 (db.news_exists) │
│ - trafilatura 提取正文内容 │
│ - db.save_news() 保存到数据库 │
│ 6. 等待所有异步摘要生成任务完成 │
└─────────────────────────────────────────────────────────────┘
关键技术点:
-
异步并发抓取: 使用
httpx.AsyncClient+asyncio.gather()同时请求多个 RSS 源,大幅提升抓取速度 -
RSS 解析: 使用
feedparser库解析各种格式的 RSS/Atom 订阅源 -
正文提取: 使用
trafilatura库智能提取网页正文- 自动去除导航栏、广告、页脚等无关内容
- 保留文章主体内容
- 处理成功率比简单的 HTML 解析高很多
-
去重机制: 数据库层面使用
ON CONFLICT (url) DO NOTHING实现幂等性存储
阶段 3: 数据存储与摘要生成
┌─────────────────────────────────────────────────────────────┐
│ db.py::save_news() │
│ │
│ 1. INSERT INTO news (ON CONFLICT DO NOTHING) 实现去重 │
│ 2. 如果插入成功,创建后台异步任务 _generate_summary_task() │
│ - 调用 summarize_async() 生成摘要 │
│ - update_summary() 更新数据库中的 summary 字段 │
└─────────────────────────────────────────────────────────────┘
设计决策: 为什么摘要生成是异步的?
同步方案的问题:
获取10条新闻 → 生成摘要1 → 生成摘要2 → ... → 生成摘要10 → 继续
(等待2秒) (等待2秒) (等待2秒)
总共需要 20 秒
异步方案的优势:
获取10条新闻 → 同时启动10个摘要任务 → 等待全部完成 → 继续
(并发执行,总共只需 2-3 秒)
*用时仅供参考,实际用时会受到网络状况、上下文长度的影响
AI 摘要生成细节 (summarize.py):
# 使用 LangChain 调用通义千问
llm = ChatOpenAI(
model="qwen-flash", # 轻量级快速模型
temperature=0.3, # 低温度,输出更稳定
)
# Prompt 设计
System: "你是一个新闻摘要助手。请将以下新闻内容总结为一段 50-100 字的中文摘要。"
User: 新闻正文内容
阶段 4: 新闻筛选
获取并摘要完所有新闻后,系统需要从众多新闻中选出最重要的 10 条。
┌─────────────────────────────────────────────────────────────┐
│ generate_news.py::select_news() │
│ │
│ 输入: 所有新闻的 ID 和摘要列表 │
│ 格式: "1 摘要: xxx\n2 摘要: xxx\n..." │
│ │
│ 1. 使用 LangChain create_agent 创建结构化输出 Agent │
│ 2. 调用通义千问 qwen-plus 模型 │
│ 3. System Prompt: "请从以下新闻中选出对投资者最重要的 10 条" │
│ 4. 返回 SelectedNews 结构: {ids: [1, 2, 3, ...]} │
└─────────────────────────────────────────────────────────────┘
为什么使用 Agent + 结构化输出?
class SelectedNews(BaseModel):
"""新闻选择结果"""
ids: list[int] = Field(description="选中的新闻 ID 列表,最多10条")
# 创建结构化输出 Agent
structured_agent = create_structured_output_agent(
llm=llm,
output_schema=SelectedNews,
prompt=prompt
)
优势:
- 可靠性: 强制 AI 返回符合格式的 JSON,避免解析错误
- 类型安全: 使用 Pydantic 模型验证输出
- 可控性: 明确限制最多选择 10 条
阶段 5: 早报格式化
┌─────────────────────────────────────────────────────────────┐
│ generate_news.py::generate_morning_brief() │
│ │
│ 1. 遍历选中的新闻 ID 列表 │
│ 2. db.get_news_by_id() 获取每条新闻详情 │
│ 3. 格式化为 Markdown: │
│ │
│ ## [新闻标题](链接) │
│ 摘要内容... │
│ │
│ 4. 拼接成完整早报内容 │
└─────────────────────────────────────────────────────────────┘
生成的 Markdown 示例:
## [美联储宣布维持利率不变](https://...)
美联储在最新议息会议上决定维持基准利率不变,符合市场预期。美联储主席表示通胀数据仍需观察...
## [特斯拉发布新款Model 3](https://...)
特斯拉正式发布改款 Model 3,续航里程提升至 600 公里,售价保持不变...
阶段 6: 消息推送
┌─────────────────────────────────────────────────────────────┐
│ send.py::send_post() │
│ │
│ 1. 使用 lark_oapi 构建飞书客户端 │
│ 2. 构造富文本消息 (msg_type="post") │
│ 3. 发送 Markdown 内容到指定群聊 │
└─────────────────────────────────────────────────────────────┘
飞书消息格式:
使用 post 消息类型,支持 Markdown 格式的富文本展示:
- 标题链接可点击
- 支持多段落格式
- 手机端和桌面端都有良好的展示效果
三、数据库设计
表结构
CREATE TABLE news (
id SERIAL PRIMARY KEY, -- 自增主键
url VARCHAR(2048) UNIQUE NOT NULL, -- 新闻链接(唯一约束,去重用)
title VARCHAR(1024) NOT NULL, -- 新闻标题
published TIMESTAMP, -- 发布时间(RSS 提供)
summary TEXT, -- AI 生成的中文摘要
content TEXT, -- 完整正文内容
fetched_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP -- 抓取时间
);
-- 索引优化查询
CREATE INDEX idx_news_url ON news(url);
CREATE INDEX idx_news_fetched_at ON news(fetched_at DESC);
设计考量
- URL 作为唯一标识: 使用
UNIQUE约束确保同一篇新闻不会被重复存储 - summary 可为空: 摘要生成是异步的,刚保存时可能还没有摘要
- fetched_at 索引: 支持按时间范围快速查询(如"过去24小时的新闻")
四、关键技术决策
1. 双模型策略
| 任务 | 模型 | 理由 |
|---|---|---|
| 摘要生成 | qwen-flash | 大批量处理,追求速度和成本效益 |
| 新闻筛选 | qwen-plus | 需要推理能力,判断重要性 |
2. 异步架构
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ 主流程 │─────►│ 异步任务池 │◄─────│ 后台执行 │
│ │ │ │ │ │
│ 抓取新闻 │ │ 摘要任务1 │ │ 调用 AI API │
│ 保存数据库 │ │ 摘要任务2 │ │ 更新数据库 │
│ 等待完成 │◄─────│ ... │─────►│ │
└─────────────┘ └─────────────┘ └─────────────┘
优势:
- 多个摘要任务并发执行
- 主流程只需等待最慢的任务,而非所有任务之和
3. 容器化部署
# docker-compose.yml 架构
services:
postgres:
- 数据持久化卷
- 健康检查
news:
- 依赖 postgres
- 环境变量注入配置
- 启动后执行定时任务
优势:
- 一次配置,到处运行
- 数据库和应用一起编排
- 环境隔离,避免冲突
五、代码执行时序图
时间 ──────────────────────────────────────────────────────────────►
main.py get_news.py db.py summarize.py send.py
│ │ │ │ │
│ job() │ │ │ │
│─────────────────►│ │ │ │
│ │ │ │ │
│ │ 并发获取 RSS │ │ │
│ │ 提取正文 │ │ │
│ │ │ │ │
│ │─────────────►│ save_news() │ │
│ │ │───────────────►│ │
│ │ │ │ │
│ │ │ 创建异步摘要任务 │ │
│ │ │ 不等待,立即返回 │ │
│ │◄─────────────│ │ │
│ │ │ │ │
│◄─────────────────│ 等待所有摘要完成 │ │
│ (asyncio.gather) │ │ │
│ │ │ │ │
│ 获取24小时内新闻 │ │ │ │
│ 拼接摘要列表 │ │ │ │
│ │ │ │ │
│ select_news() │ │ │ │
│ AI 筛选10条重要新闻 │ │ │
│ │ │ │ │
│ generate_morning_brief() │ │ │
│ 格式化为 Markdown │ │ │
│ │ │ │ │
│──────────────────────────────────────────────────────────────► │
│ send_post() 发送到飞书 │
│ │
│◄──────────────────────────────────────────────────────────────│
│ 完成 ✅ │
六、异常处理与容错
1. 网络故障
# RSS 抓取失败不会影响其他源
async def fetch_feed(client, url):
try:
response = await client.get(url)
# 解析...
except Exception as e:
logger.error(f"获取 RSS 失败 {url}: {e}")
return [] # 返回空列表,不影响其他源
2. 正文提取失败
content = get_full_content(url)
if not content:
# 没有正文也能继续,跳过摘要生成
continue
3. 摘要生成失败
async def _generate_summary_task(news_id, content):
try:
summary = await summarize_async(content)
await update_summary(news_id, summary)
except Exception as e:
logger.error(f"生成摘要失败: {e}")
# 失败不影响主流程
4. 数据库重复
INSERT INTO news (...) VALUES (...)
ON CONFLICT (url) DO NOTHING
-- 重复 URL 自动忽略,不报错
七、扩展性考虑
1. 添加新的 RSS 源
只需在 config.toml 中添加 URL,无需修改代码。
2. 更换 AI 模型
所有 AI 调用都通过 LangChain 封装,更换模型只需修改配置:
# 当前使用通义千问
llm = ChatOpenAI(
model="qwen-flash",
base_url="https://dashscope.aliyuncs.com/compatible-mode/v1"
)
# 更换为 OpenAI 只需修改 base_url 和 api_key
3. 更换推送渠道
send.py 是独立的推送模块,更换推送渠道只需修改此文件。
八、总结
这个系统的设计哲学是:
- 简单可靠: 使用成熟的技术栈,避免过度设计
- 异步高效: 利用 asyncio 提升 I/O 密集型任务的性能
- 分层解耦: 数据获取、存储、AI 处理、推送各司其职
- 容错性强: 单点失败不影响整体流程
- 易于维护: 配置与代码分离,修改无需重新部署
通过 RSS 聚合、AI 摘要、智能筛选、定时推送的组合,系统实现了全自动的新闻早报生成,为投资者节省了大量信息筛选的时间。