一个完整的A股数据收集系统,采用ClickHouse存储,TuShare API数据源,并支持自动化编排。
git clone <repository-url>
cd stock_datasource
uv sync
cp .env.example .env
# 编辑 .env 文件,填入 TuShare Token 和 ClickHouse 配置
uv run cli.py init-db
# 发现所有插件
uv run python -m stock_datasource.cli_plugins discover
# 列出所有插件
uv run python -m stock_datasource.cli_plugins list
# 查看插件详情
uv run python -m stock_datasource.cli_plugins info tushare_daily
# 测试插件数据提取
uv run python -m stock_datasource.cli_plugins test --date 20251024
# 获取特定日期数据
uv run cli.py ingest-daily --date 20251024
# 批量回填数据
uv run cli.py backfill --start-date 20250101 --end-date 20251024
# 查看摄入状态
uv run cli.py status --date 20251024
# 运行质量检查
uv run cli.py quality-check --date 20251024
# 生成日报告
uv run cli.py report --date 20251024
# 优化表去除重复数据
uv run python -c "from src.stock_datasource.models.database import db_client; db_client.execute('OPTIMIZE TABLE ods_daily FINAL')"
# 检查重复数据情况
uv run python -c "
from src.stock_datasource.models.database import db_client
total = db_client.execute('SELECT COUNT(*) FROM ods_daily')[0][0]
unique = db_client.execute('SELECT COUNT(DISTINCT (ts_code, trade_date)) FROM ods_daily')[0][0]
print(f'总记录: {total:,}, 唯一: {unique:,}, 重复: {total-unique:,}')
"
stock_datasource/
├── 📄 核心文档
│ ├── README.md # 项目概览(本文件)
│ ├── DEVELOPMENT_GUIDE.md # 详细开发指导
│ ├── PLUGIN_QUICK_START.md # 新建插件快速参考
│ ├── README_SUMMARY.md # 项目总结和导航
│ └── BASEPLUGIN_QUICK_REFERENCE.md # BasePlugin API 参考
│
├── 🔧 核心代码
│ ├── cli.py # CLI 命令入口
│ ├── pyproject.toml # 项目配置和依赖
│ ├── uv.lock # 依赖锁定文件
│ └── src/stock_datasource/
│ ├── core/ # 核心模块
│ │ ├── base_plugin.py # BasePlugin 基类
│ │ ├── database.py # 数据库连接
│ │ └── plugin_manager.py # 插件管理器
│ ├── plugins/ # 数据插件(7 个)
│ │ ├── tushare_daily/
│ │ ├── tushare_adj_factor/
│ │ ├── tushare_daily_basic/
│ │ ├── tushare_stock_basic/
│ │ ├── tushare_stk_limit/
│ │ ├── tushare_suspend_d/
│ │ └── tushare_trade_calendar/
│ └── utils/ # 工具函数
│
├── 📊 数据目录
│ ├── data/ # 数据存储
│ │ └── exports/ # 导出数据
│ └── logs/ # 运行日志
│
├── 🧪 测试目录
│ └── tests/ # 单元测试
│
├── 🛠️ 脚本工具
│ └── scripts/
│ └── optimize_tables.py # 表优化脚本(去重复数据)
│
├── 📚 配置文件
│ ├── .env.example # 环境变量示例
│ ├── .gitignore # Git 忽略规则
│ ├── .python-version # Python 版本
│ └── LICENSE # 许可证
│
└── 📖 文档目录
└── docs/ # 其他文档
TuShare API
↓
插件 (提取 → 验证 → 转换 → 加载)
↓
ODS 层 (原始数据,ReplacingMergeTree 幂等存储)
├─ ods_daily (日线数据,按月分区)
├─ ods_adj_factor (复权因子,version 去重)
├─ ods_daily_basic (日线基础指标,自动合并)
├─ ods_stock_basic (股票基础信息)
├─ ods_stk_limit (涨跌停数据)
├─ ods_suspend_d (停复牌数据)
└─ ods_trade_calendar (交易日历)
↓
DM/Fact 层 (清洗数据,稳定业务表)
├─ fact_daily_bar (事实表)
└─ dim_security (维度表)
↓
元数据层 (审计日志)
├─ ingestion_logs (摄入日志)
├─ quality_checks (质量检查)
└─ schema_evolution (Schema 演变)
ReplicingMergeTree 引擎特性:
version 字段(时间戳)toYYYYMM(trade_date) 提升性能| 插件 | 表名 | 描述 | 参数 |
|---|---|---|---|
tushare_daily | ods_daily | 日线数据 | trade_date |
tushare_adj_factor | ods_adj_factor | 复权因子 | trade_date |
tushare_daily_basic | ods_daily_basic | 日线基础指标 | trade_date |
tushare_stock_basic | ods_stock_basic | 股票基础信息 | 无 |
tushare_stk_limit | ods_stk_limit | 涨跌停数据 | trade_date |
tushare_suspend_d | ods_suspend_d | 停复牌数据 | trade_date |
tushare_trade_calendar | ods_trade_calendar | 交易日历 | start_date, end_date |
mkdir -p src/stock_datasource/plugins/my_plugin
cd src/stock_datasource/plugins/my_plugin
touch __init__.py plugin.py extractor.py service.py config.json schema.json
from stock_datasource.core.base_plugin import BasePlugin
import pandas as pd
from .extractor import extractor
class MyPlugin(BasePlugin):
@property
def name(self) -> str:
return "my_plugin"
@property
def description(self) -> str:
return "我的自定义插件"
def extract_data(self, **kwargs) -> pd.DataFrame:
"""从数据源获取原始数据"""
trade_date = kwargs.get('trade_date')
data = extractor.extract(trade_date)
return data
def validate_data(self, data: pd.DataFrame) -> bool:
"""验证数据的完整性和正确性"""
if data.empty:
return False
return True
def load_data(self, data: pd.DataFrame) -> dict:
"""将清洗后的数据加载到数据库"""
if not self.db:
return {"status": "failed", "error": "数据库未初始化"}
self.db.insert_dataframe('ods_my_table', data)
return {"status": "success", "records": len(data)}
import pandas as pd
from stock_datasource.config.settings import settings
import tushare as ts
class Extractor:
def __init__(self):
self.pro = ts.pro_api(settings.TUSHARE_TOKEN)
def extract(self, trade_date: str) -> pd.DataFrame:
"""从 TuShare API 获取数据"""
data = self.pro.daily(trade_date=trade_date)
return data
extractor = Extractor()
from typing import List, Dict, Any
from stock_datasource.core.base_service import BaseService, query_method, QueryParam
class MyPluginService(BaseService):
"""我的插件数据查询服务。"""
def __init__(self):
super().__init__("my_plugin")
@query_method(
description="根据代码和日期范围查询数据",
params=[
QueryParam(name="code", type="str", description="股票代码", required=True),
QueryParam(name="start_date", type="str", description="起始日期 YYYYMMDD", required=True),
QueryParam(name="end_date", type="str", description="结束日期 YYYYMMDD", required=True),
]
)
def get_data(self, code: str, start_date: str, end_date: str) -> List[Dict[str, Any]]:
"""从数据库查询数据。"""
query = f"""
SELECT * FROM ods_my_table
WHERE ts_code = '{code}'
AND trade_date >= '{start_date}'
AND trade_date <= '{end_date}'
ORDER BY trade_date ASC
"""
df = self.db.execute_query(query)
return df.to_dict('records')
{
"enabled": true,
"rate_limit": 120,
"timeout": 30,
"retry_attempts": 3,
"description": "我的自定义插件",
"parameters_schema": {
"trade_date": {
"type": "string",
"format": "date",
"required": true,
"description": "交易日期,格式为 YYYYMMDD"
}
}
}
{
"table_name": "ods_my_table",
"table_type": "ods",
"columns": [
{"name": "ts_code", "data_type": "String", "nullable": false},
{"name": "trade_date", "data_type": "Date", "nullable": false},
{"name": "version", "data_type": "UInt64", "nullable": false},
{"name": "_ingested_at", "data_type": "DateTime", "nullable": false}
],
"partition_by": "toYYYYMM(trade_date)",
"order_by": ["ts_code", "trade_date"],
"engine": "ReplacingMergeTree"
}
from .plugin import MyPlugin
from .service import MyPluginService
__all__ = ["MyPlugin", "MyPluginService"]
# 发现插件
uv run python -m stock_datasource.cli_plugins discover
# 查看插件详情
uv run python -m stock_datasource.cli_plugins info my_plugin
# 测试数据提取
uv run python -m stock_datasource.cli_plugins test --plugin my_plugin --date 20251024
extract_data()、validate_data()、load_data()extract_data() 必须返回 pd.DataFrameversion 和 _ingested_atself.logger 记录日志extract() 中进行数据转换应在 transform() 中进行(中间层)None 或其他非 DataFrame 类型| 文档 | 目的 | 适合人群 |
|---|---|---|
README.md | 项目概述和快速开始 | 所有人 |
DEVELOPMENT_GUIDE.md | 详细的开发指导 | 开发者 |
PLUGIN_QUICK_START.md | 新插件快速参考 | 新手开发者 |
README_SUMMARY.md | 项目总结和导航 | 项目经理 |
BASEPLUGIN_QUICK_REFERENCE.md | BasePlugin API 参考 | API 开发者 |
| 变量 | 描述 | 默认值 |
|---|---|---|
TUSHARE_TOKEN | TuShare API Token | 必需 |
CLICKHOUSE_HOST | ClickHouse 服务器地址 | localhost |
CLICKHOUSE_PORT | ClickHouse 服务器端口 | 9000 |
CLICKHOUSE_DATABASE | ClickHouse 数据库名称 | stock_datasource |
LOG_LEVEL | 日志级别 | INFO |
该项目采用了统一服务层设计,通过每个插件的 service.py 统一管理所有数据查询逻辑:
plugins/tushare_daily/
├── plugin.py (数据采集:提取 → 验证 → 加载)
├── extractor.py (API 调用逻辑)
├── service.py (查询接口:定义 @query_method)
├── config.json (插件配置)
└── schema.json (表结构定义)
│
└─→ ServiceGenerator (自动生成)
├── 生成 HTTP 路由 (FastAPI)
└── 生成 MCP 工具定义
│
├─→ HTTP Server (http_server.py)
│ └── 暴露 REST API 端点
│
└─→ MCP Server (mcp_server.py)
└── 暴露 MCP 工具接口
关键特性:
service.py 中间定义一次@query_method 装饰器定义参数和描述客户端请求
│
├─→ HTTP 请求 → HTTP Server → ServiceGenerator → TuShareDailyService.get_daily_data() → ClickHouse
│
└─→ MCP 请求 → MCP Server → ServiceGenerator → TuShareDailyService.get_daily_data() → ClickHouse
# 方法 1:使用 uv
uv run python -m stock_datasource.services.http_server
# 方法 2:使用 uvicorn
uvicorn stock_datasource.services.http_server:app --host 0.0.0.0 --port 8000
# 查询特定股票的日线数据
curl -X POST http://localhost:8000/api/tushare_daily/get_daily_data \
-H "Content-Type: application/json" \
-d '{
"code": "000001.SZ",
"start_date": "20250101",
"end_date": "20251024"
}'
# 查询最新日线数据
curl -X POST http://localhost:8000/api/tushare_daily/get_latest_daily \
-H "Content-Type: application/json" \
-d '{
"codes": ["000001.SZ", "000