MCP服务器+股价摄入管道
它的作用: 通过Airflow每5分钟获取5个股票价格,将其存储在Postgres中,并公开实时MCP工具+历史REST API。
我的学习过程: MCP对我来说是新的,所以我开始构建一个更简单的版本 本教程 并阅读 MCP官方文件 了解工具注册、验证和错误处理。一旦我理解了这个模式,我就将其应用于生产考虑。
我是如何构建它的:
- 从Postgres(不是专门的时间序列数据库)开始,因为在这种规模下操作起来更简单
- 每一层都内置了可观察性:结构化日志、指标计数器和清晰的错误消息
- 使气流重试安全
ON CONFLICT DO NOTHING所以重复运行不会破坏任何东西 - 使用模拟测试了所有内容,因此测试运行速度很快,不会影响外部API
是什么让这个生产准备就绪: 每个设计选择都有一个“为什么”记录在下面,代码被键入并测试,如果凌晨3点有什么东西坏了,日志会准确地告诉你失败的原因。
┌────────────┐ yfinance / CoinGecko
│ MCP tools │──────────────┐
└────────────┘ │
┌────────────────┴───────┐
│ FastAPI (/mcp,/prices)│
└──────┬─────────────────┘
│
▼
┌──────────────────────┐
│ Postgres / Timescale │ ◄── Airflow DAG (5mins)
└──────────────────────┘______________________________________________________________________
架构概述
数据流:
- 气流调度器 每5分钟触发一次
- 摄入CLI 通过yfinance获取5家股票(AAPL、MSFT、GOOGL、AMZN、TSLA)的价格
- Postgres 使用幂等写存储价格(
ON CONFLICT DO NOTHING) - 快速API 暴露出两个表面:
- MCP工具 (/mcp/tools/*)-实时股票/加密货币价格,绕过数据库 - REST API (/prices/*)-来自数据库的历史查询(最新、范围、24小时返回)
关键部件:
app/services/market_data.py-从yfinance/CoinGecko获取,具有3层回退、5秒超时、60秒缓存app/services/ingestion.py-按指数回退(0s、1s、2s)重试每个股票dags/stock_ingestion_dag.py-气流DAG运行docker exec price-api-1 python3 -m app.cli.ingest_once
______________________________________________________________________
设计合理性
存储:带时间刻度选项的Postgres
选择Postgres而不是专门的时间序列数据库(InfluxDB、QuestDB),因为:
- 操作更简单:每个人都知道Postgres,标准工具有效
- 足够好的性能:5个5分钟间隔的自动收报机不需要专门的数据库
- 轻松升级路径:
CREATE EXTENSION timescaledb如果我们长大了香草Postgres - ACID保证:对财务数据完整性很重要
方案设计:
CREATE TABLE prices (
ticker TEXT NOT NULL,
price NUMERIC(18, 8) NOT NULL,
bucket_start TIMESTAMPTZ NOT NULL, -- 5-minute aligned UTC timestamp
source_ts TIMESTAMPTZ NULL, -- when upstream reported the price
PRIMARY KEY (ticker, bucket_start)
);
CREATE INDEX idx_prices_source_ts ON prices (source_ts);- 复合PK
(ticker, bucket_start)强制唯一性并启用幂等写入 - source_ts索引 允许追踪雅虎财经实际发送数据的时间
- 数字(18,8) for price避免浮点错误
Idempotent写道: 使用 ON CONFLICT DO NOTHING 使气流重试安全。如果一次运行执行两次(网络光点),第二次尝试会自动跳过现有行。没有违反约束,没有重复数据。
错误处理:
- 按股票代码try/except:如果TSLA超时,其他4个仍然会被写入
- 指数回退:重试之间为0s、1s、2s
- 所有yfnance调用超时5秒(API在开市期间可以挂起20秒以上)
可观察性:
- 结构化日志:插入/跳过/失败的每次运行日志计数
- Metrics注册表:跟踪缓存命中率,为未来的Prometheus集成进行摄取运行
- 清除错误消息:404表示未知符号,429表示速率限制,503表示上游故障
______________________________________________________________________
评估
我所测量的:
生成了约10万行(5个股票代码的1年5分钟数据),并测试了所有端点:
| 端点 | 有效负载 | 平均延迟 | 备注 |
|---|---|---|---|
/prices/latest | 1行 | 3.2毫秒 | 单PK查找 |
/prices/range (24小时) | 288行 | 6.8毫秒 | 查询+JSON序列化 |
/prices/range (1年) | 105k行 | 41.5ms | 全扫描,仍然可以接受 |
/prices/24h-return | 2行 | 4.1毫秒 | 导出度量计算 |
在SQLite上使用FastAPI测试客户端进行测试。Postgres应该匹配或击败这些数字。
查询优化:
/latest使用复合PK进行O(1)查找/range连续bucketstart值的好处(适用于时间序列扫描)source_ts索引支持无需全表扫描的跟踪查询
缩放注释:
- 当前瓶颈:金融API调用(顺序调用,每次约500ms)
- 在100个自动收报机时:切换到异步批取或并行任务分片
- 每隔1分钟:考虑流式馈送(websocket/Kafka)而不是轮询
- 重读:添加Postgres读副本或Redis缓存
/latest
______________________________________________________________________
下一步(还有更多时间)
演出
- 启用压缩旧数据的时间尺度超表(
chunk_time_interval => '1 day') - 每月分区一次,每天摄入数百万行
- 90天后将旧数据归档到S3/GCS中的Parquet
可靠性:
- 添加SLA监控:如果
ingest_last_success_epoch不前进超过10分钟 - 在停机期间对错过的间隔实施回填逻辑
- 为金融添加断路器(连续N次故障后自动禁用自动收报机)
操作:
- 切换到气流的LocalExecutor+Postgres元数据(生产就绪)
- 为仪表板添加Prometheus指标导出
- 设置自动模式迁移测试
特征:
- 通过API支持自定义股票行情表(
POST /config/tickers) - 添加更多MCP工具:历史统计数据、波动率计算、收益日历
- 用于向前端实时传输价格的WebSocket端点
______________________________________________________________________
安装说明
先决条件:
- Docker&Docker编写
- Python 3.11+(用于本地测试)
快速启动:
# 1. Start all services
docker compose up -d
# 2. Wait ~60s for Airflow to initialize, then check http://localhost:8080
# (username: admin, password: auto-generated in logs)
# 3. Verify ingestion is running
docker exec price-airflow-1 airflow dags list
# 4. Test MCP tools
curl http://localhost:8000/mcp/tools | jq
curl -X POST http://localhost:8000/mcp/tools/get_stock_price \
-H "Content-Type: application/json" \
-d '{"ticker":"AAPL"}' | jq
# 5. Test REST API (after first ingestion completes)
curl "http://localhost:8000/prices/latest?ticker=AAPL" | jq运行测试:
python3 -m venv .venv
source .venv/bin/activate
pip install -r requirements/base.txt -r requirements/dev.txt
pytest -v项目结构:
app/
├── api/ # FastAPI routers for MCP + /prices
├── cli/ # CLI script for ingestion (called by Airflow)
├── core/ # Settings, logging, metrics
├── db/ # SQLAlchemy models + queries
├── services/ # Market data client + ingestion logic
└── main.py # FastAPI app factory
dags/ # Airflow DAG (5-minute schedule)
migrations/ # SQL DDL (idempotent)
tests/ # Pytest suite (mocked externals)
docs/evidence/ # Proof of working system (logs, API responses)配置: 通过环境变量进行所有设置(请参见 .env.example):
DB_URL:Postgres连接字符串TICKERS:逗号分隔的股票代码列表(默认:AAPL、MSFT、GOOGL、AMZN、TSLA)CACHE_TTL_SECONDS:MCP工具缓存持续时间(默认值:60)YFINANCE_TIMEOUT_SECONDS:API调用超时(默认值:5)
______________________________________________________________________
证据
气流5分钟时间表证明: docs/evidence/airflow-5min-schedule-proof.txt 显示了UTC时间02:15、02:20、02:25、02:30、02:35、02:40连续6次成功运行,所有运行均为 state=success.
API回应:
docs/evidence/mcp-tools.json-3个MCP工具列表docs/evidence/mcp-stock-price.json-通过MCP的股票价格docs/evidence/mcp-crypto-price.json-MCP加密货币价格docs/evidence/prices-latest.json-DB的最新价格docs/evidence/prices-range.json-历史价格区间docs/evidence/pytest-output.txt-全部10项测试通过
