返回市场
股票数据源

股票数据源

作者:Yourdaylight3 星标更新:2025-11-23

项目介绍

股票数据源 - A股金融数据库

一个完整的A股数据收集系统,采用ClickHouse存储,TuShare API数据源,并支持自动化编排。

📊 项目特性

  • 完整的A股数据 日线图、复权因子、基本指标、涨跌停、停牌复牌等
  • 7个预装插件 开箱即用的数据收集插件
  • 高性能存储 ClickHouse列式数据库,支持PB级数据
  • 自动化编排 Airflow DAG支持定时任务
  • 多层数据质量 ODS → DM/Fact → 元数据三层架构
  • 幂等保证 ReplacingMergeTree引擎确保数据一致性
  • CSV快照功能 每次数据提取后自动保存CSV文件,便于定义模式和数据检查
  • 多模式支持 一个插件可以定义多个表(ODS+FACT/BIM)
  • 可扩展架构 容易添加新的数据源和插件

🚀 快速开始

预备条件

  • Python 3.11+
  • ClickHouse 服务器
  • TuShare API Token
  • UV 包管理工具

安装步骤

  1. 克隆仓库
git clone <repository-url>
cd stock_datasource
  1. 安装依赖
uv sync
  1. 配置环境
cp .env.example .env
# 编辑 .env 文件,填入 TuShare Token 和 ClickHouse 配置
  1. 初始化数据库
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) 提升性能

📋 7个预装插件

插件表名描述参数
tushare_dailyods_daily日线数据trade_date
tushare_adj_factorods_adj_factor复权因子trade_date
tushare_daily_basicods_daily_basic日线基础指标trade_date
tushare_stock_basicods_stock_basic股票基础信息
tushare_stk_limitods_stk_limit涨跌停数据trade_date
tushare_suspend_dods_suspend_d停复牌数据trade_date
tushare_trade_calendarods_trade_calendar交易日历start_date, end_date

🔧 创建新插件的完整步骤

1. 创建插件目录结构

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

2. 实现 plugin.py (数据收集)

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)}

3. 实现 extractor.py (API调用)

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()

4. 实现 service.py (查询接口)

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')

5. 编写 config.json (配置)

{
  "enabled": true,
  "rate_limit": 120,
  "timeout": 30,
  "retry_attempts": 3,
  "description": "我的自定义插件",
  "parameters_schema": {
    "trade_date": {
      "type": "string",
      "format": "date",
      "required": true,
      "description": "交易日期,格式为 YYYYMMDD"
    }
  }
}

6. 编写 schema.json (表结构)

{
  "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"
}

7. 注册插件(init.py)

from .plugin import MyPlugin
from .service import MyPluginService

__all__ = ["MyPlugin", "MyPluginService"]

8. 测试验证

# 发现插件
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.DataFrame
  • 添加系统字段version_ingested_at
  • 使用日志:使用 self.logger 记录日志
  • 处理异常:捕获并处理所有异常情况

❌ 常见错误

  • ❌ 在 extract() 中进行数据转换应在 transform() 中进行(中间层)
  • ❌ 忽视异常并让它们传播到上层
  • ❌ 忽略数据验证
  • ❌ 返回 None 或其他非 DataFrame 类型
  • ❌ 硬编码配置值

📚 文档导航

文档目的适合人群
README.md项目概述和快速开始所有人
DEVELOPMENT_GUIDE.md详细的开发指导开发者
PLUGIN_QUICK_START.md新插件快速参考新手开发者
README_SUMMARY.md项目总结和导航项目经理
BASEPLUGIN_QUICK_REFERENCE.mdBasePlugin API 参考API 开发者

🔐 环境配置

环境变量

变量描述默认值
TUSHARE_TOKENTuShare API Token必需
CLICKHOUSE_HOSTClickHouse 服务器地址localhost
CLICKHOUSE_PORTClickHouse 服务器端口9000
CLICKHOUSE_DATABASEClickHouse 数据库名称stock_datasource
LOG_LEVEL日志级别INFO

📊 数据统计

  • 时间范围 从2025年1月1日至2025年10月24日(195个交易日)
  • 股票数量 ~5400只A股
  • 数据表 7个ODS表 + 2个Fact表 + 1个Dim表
  • 总记录数:约1.2亿条记录(每天约6百万条记录)

🌐 自动生成HTTP服务器与MCP服务器之间的接口

架构描述

该项目采用了统一服务层设计,通过每个插件的 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 中间定义一次
  • 自动生成 HTTP路由和MCP工具从Service方法自动生成
  • 元数据驱动 通过 @query_method 装饰器定义参数和描述
  • 代码复用 HTTP和MCP共享相同的业务逻辑
  • 动态发现 自动发现所有插件的服务类

数据流

客户端请求
    │
    ├─→ HTTP 请求 → HTTP Server → ServiceGenerator → TuShareDailyService.get_daily_data() → ClickHouse
    │
    └─→ MCP 请求 → MCP Server → ServiceGenerator → TuShareDailyService.get_daily_data() → ClickHouse

🚀 HTTP服务器使用

启动HTTP服务器

# 方法 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

HTTP请求示例

# 查询特定股票的日线数据
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