Files
quanxiel/quantitative_data/数据批量导入.ipynb

21 KiB
Raw Permalink Blame History

量化投资数据批量导入

目标

将 Tushare 的日线行情数据及公司基本面数据批量导入 Docker PostgreSQL (192.168.27.11:5438)

数据库结构概览

表名 说明 Tushare 接口
stock_basic 股票基本信息 stock_basic
trade_cal 交易日历 trade_cal
daily 日线行情 daily
daily_basic 每日指标(估值/基本面) daily_basic
adj_factor 复权因子 adj_factor
income 利润表 income
balancesheet 资产负债表 balancesheet
cashflow 现金流量表 cashflow
fina_indicator 财务指标 fina_indicator
moneyflow 个股资金流向 moneyflow
index_daily 指数日线行情 index_daily

使用步骤

  1. 复制 .env.example.env,填入真实密码和 Token
  2. 逐 Cell 运行本 Notebook(敏感信息通过环境变量读取)

Step 0: 检查环境 & 连接测试

In [8]:
import sys
import os
os.chdir(r"d:\projects\quanxiel\quantitative_data")
print(f"工作目录: {os.getcwd()}")
print(f"Python 版本: {sys.version}")
In [9]:
# 检查依赖包
!pip list | findstr -i "tushare pandas psycopg2-binary sqlalchemy python-dotenv"
In [ ]:
# 如果需要安装依赖,取消注释下面这行
#!pip install -r requirements.txt
In [10]:
# 导入核心模块
from importer import (
    get_pg_connection,
    get_ts_pro,
    init_database,
    import_stock_basic,
    import_trade_cal,
    import_daily_for_stock,
    import_daily_by_date,
    import_daily_by_year,
    import_daily_basic,
    import_daily_basic_by_date,
    import_adj_factor,
    import_adj_factor_batch,
    import_financial_statements,
    import_index_daily,
    get_all_stock_codes,
    get_stock_codes_from_db,
    full_import,
    batch_insert,
    check_daily_progress,
    check_table_summary,
    get_missing_daily_dates,
    resume_daily_by_date,
    logger,
)
from config import DB_CONFIG, TUSHARE_TOKEN, START_DATE, END_DATE

print("模块导入成功!")
In [11]:
# 测试数据库连接
try:
    conn = get_pg_connection()
    cursor = conn.cursor()
    cursor.execute("SELECT version()")
    version = cursor.fetchone()[0]
    print(f"✓ PostgreSQL 连接成功!")
    print(f"  服务器版本: {version}")
    print(f"  连接信息: {DB_CONFIG['host']}:{DB_CONFIG['port']}/{DB_CONFIG['database']}")
    cursor.close()
    conn.close()
except Exception as e:
    print(f"✗ 连接失败: {e}")
    print("请检查:")
    print("  1. Docker 容器是否已启动: docker ps | findstr postgres")
    print("  2. 环境变量 (.env) 中的连接参数是否正确")
    print("  3. 防火墙是否开放 5438 端口")
In [ ]:
# 测试 Tushare API 连接
try:
    pro = get_ts_pro()
    # 简单测试: 获取一只股票信息
    df = pro.stock_basic(ts_code="000001.SZ", fields="ts_code,name,industry")
    print(f"✓ Tushare API 连接成功!")
    print(f"  测试查询: {df.iloc[0].to_dict()}")
except Exception as e:
    print(f"✗ Tushare API 连接失败: {e}")
    print("请检查环境变量 TUSHARE_TOKEN 是否正确(或 .env 文件)")

Step 1: 初始化数据库 Schema

In [ ]:
# 执行 DDL,创建所有表结构
init_database()

Step 2: 导入股票基本信息

In [ ]:
# 导入全量 A 股股票基本信息(含上市和退市)
import_stock_basic()
In [ ]:
# 验证:查看导入结果
conn = get_pg_connection()
cursor = conn.cursor()
cursor.execute("SELECT COUNT(*) FROM stock_basic")
print(f"stock_basic 总记录数: {cursor.fetchone()[0]}")
cursor.execute("SELECT list_status, COUNT(*) FROM stock_basic GROUP BY list_status")
for row in cursor.fetchall():
    print(f"  状态 '{row[0]}': {row[1]}")
cursor.close()
conn.close()

Step 3: 导入交易日历

In [ ]:
# 导入交易日历 (默认从环境变量 QUANT_START_DATE ~ QUANT_END_DATE)
import_trade_cal()
In [ ]:
# 或者指定日期范围
# import_trade_cal(start_date="2020-01-01", end_date="2025-12-31")
In [ ]:
# 验证
conn = get_pg_connection()
cursor = conn.cursor()
cursor.execute("""
    SELECT exchange, MIN(cal_date) AS first_date, MAX(cal_date) AS last_date, COUNT(*) AS total
    FROM trade_cal
    GROUP BY exchange
""")
for row in cursor.fetchall():
    print(f"  {row[0]}: {row[1]} ~ {row[2]}, 共 {row[3]}")
cursor.close()
conn.close()

Step 4: 导入日线行情 (核心表,最耗时)

新版改进: 按交易日循环拉取全市场数据 pro.daily_vip(trade_date='20180810')VIP接口), API 调用从 ~80,000次 (5000只×16年) 降至 ~4,000次 (250交易日×16年),速度提升约 20倍

4.1 按年批量导入 (推荐)

数据量估算: ~4000个交易日,每个交易日约5000条记录 ≈ 2000万条记录

注意: 新版 import_daily_by_year 不再需要 stock_list 参数,内部自动从 trade_cal 获取交易日列表后逐日拉取全市场数据。

In [ ]:
# 按年份逐批导入日线行情 (按交易日循环拉取全市场数据)
# 不再需要 stock_list 参数 — 自动查询 trade_cal 获取交易日
# 如果中断,修改年份范围从断点继续即可 (UPSERT 幂等,不会重复)
import_daily_by_year(
    start_year=2010,
    end_year=2025,
    sleep_interval=0.3,  # API 频率控制,免费版建议 0.3~0.5
)

4.1.1 查看日线导入进度

按年份统计 daily 表中已导入的记录数和独立股票数,了解导入到哪个年份了。

In [ ]:
# 查看每天的日线导入进度(按年份统计)
check_daily_progress()
In [ ]:
# 查看所有表的整体概览
check_table_summary()

4.1.2 按缺失日期断点续传 (推荐)

新版 resume_daily_by_date 自动对比 trade_caldaily 表,只导入缺失日期的全市场数据。 粒度精确到交易日级别,比旧的按年份续传更精细。

工作流程:

  1. 先调用 get_missing_daily_dates() 查看缺失的交易日列表
  2. 调用 resume_daily_by_date() 仅补缺缺失的交易日
In [ ]:
# 第1步:查看缺失的交易日 (仅查询,不导入)
missing = get_missing_daily_dates(
    start_date="2010-01-01",
    end_date="2025-12-31",
)
print(f"缺失交易日总数: {len(missing)}")
if len(missing) <= 30:
    print(f"缺失日期: {missing}")
else:
    # 按年份汇总显示
    from collections import Counter
    year_counts = Counter(d[:4] for d in missing)
    for y in sorted(year_counts):
        print(f"  {y} 年: {year_counts[y]} 个缺失交易日")
    print(f"  (前10个缺失日期): {missing[:10]}")
In [ ]:
# 第2步:按缺失日期断点续传,自动补充缺失的交易日数据
# 例如之前中断了,这里只会导入尚未导入的交易日数据
resume_daily_by_date(
    start_date="2010-01-01",
    end_date="2025-12-31",
    sleep_interval=0.3,
)

备用方案: 也可以手动指定断点年份重新调用 import_daily_by_year,因为 UPSERT 幂等,重复导入已存在的数据不会产生重复记录。

# 例如假设 2010~2020 已完成,从 2021 年继续
import_daily_by_year(start_year=2021, end_year=2025)

也可查看 import_data.log 文件获取最后一次成功的日志输出。

4.2 单只股票补充导入

如果某只股票数据缺失,可以单独补导(保留原接口兼容)。

In [ ]:
# 导入单只股票日线行情 (用于补充导入或测试)
conn = get_pg_connection()
n = import_daily_for_stock("000001.SZ", "2020-01-01", "2020-12-31", conn)
print(f"导入 000001.SZ 2020年数据: {n}")
conn.close()
In [ ]:
# 验证日线数据
conn = get_pg_connection()
cursor = conn.cursor()
cursor.execute("""
    SELECT 
        COUNT(*) AS total_records,
        COUNT(DISTINCT ts_code) AS stock_count,
        MIN(trade_date) AS first_date,
        MAX(trade_date) AS last_date
    FROM daily
""")
for row in cursor.fetchall():
    print(f"  总记录数: {row[0]:,}")
    print(f"  股票数量: {row[1]}")
    print(f"  日期范围: {row[2]} ~ {row[3]}")
cursor.close()
conn.close()

Step 5: 导入每日指标 (估值数据)

In [ ]:
# 按交易日导入每日指标 (PE/PB/PS/总市值/流通市值等)
import_daily_basic_by_date(
    start_date="2010-01-01",
    end_date="2025-12-31",
)
In [ ]:
# 验证
conn = get_pg_connection()
cursor = conn.cursor()
cursor.execute("SELECT COUNT(*), MIN(trade_date), MAX(trade_date) FROM daily_basic")
row = cursor.fetchone()
print(f"  daily_basic: {row[0]:,} 条, {row[1]} ~ {row[2]}")
cursor.close()
conn.close()

Step 6: 导入复权因子

In [ ]:
# 导入复权因子 (用于前复权/后复权价格计算)
# 需要先获取股票列表
stock_list = get_stock_codes_from_db()
print(f"{len(stock_list)} 只股票")

import_adj_factor_batch(
    stock_list,
    start_date="2010-01-01",
    end_date="2025-12-31",
)

Step 7: 导入财务数据 (三大报表 + 财务指标)

⚠ 此步骤耗时较长,约需数小时(取决于股票数量)

In [ ]:
# 导入利润表、资产负债表、现金流量表、财务指标
# stock_list 从上一步已获取,或重新获取
if 'stock_list' not in dir():
    stock_list = get_stock_codes_from_db()

import_financial_statements(
    stock_list,
    start_date="2010-01-01",
    end_date="2025-12-31",
)
In [ ]:
# 验证财务数据
conn = get_pg_connection()
cursor = conn.cursor()
for table in ["income", "balancesheet", "cashflow", "fina_indicator"]:
    cursor.execute(f"SELECT COUNT(*) FROM {table}")
    count = cursor.fetchone()[0]
    print(f"  {table}: {count:,}")
cursor.close()
conn.close()

Step 8: 导入指数日线行情

In [ ]:
# 导入主要指数日线行情
import_index_daily(
    index_codes=[
        "000001.SH",  # 上证指数
        "399001.SZ",  # 深证成指
        "000300.SH",  # 沪深300
        "000905.SH",  # 中证500
        "399006.SZ",  # 创业板指
        "000688.SH",  # 科创50
        "000016.SH",  # 上证50
        "399005.SZ",  # 中小100
        "000852.SH",  # 中证1000
    ],
    start_date="2010-01-01",
    end_date="2025-12-31",
)

一键全量导入 (可选)

如果不想逐步执行,可以运行下面这个 Cell 一键完成所有导入

In [ ]:
# 一键全量导入 (需数小时,请谨慎)
# full_import(
#     start_date="2010-01-01",
#     end_date="2025-12-31",
#     import_financials=True,  # 设为 False 跳过财务数据加快速度
# )

数据验证与查询示例

In [ ]:
import pandas as pd
import psycopg2

conn = get_pg_connection()

# 各表统计
tables = ["stock_basic", "daily", "daily_basic", "adj_factor",
          "income", "balancesheet", "cashflow", "fina_indicator",
          "trade_cal", "index_daily"]

print(f"{'表名':<20} {'记录数':>12} {'最早日期':>12} {'最晚日期':>12}")
print("-" * 60)
for table in tables:
    try:
        count_sql = f"SELECT COUNT(*) FROM {table}"
        count = pd.read_sql(count_sql, conn).iloc[0, 0]
        
        # 尝试获取日期范围
        date_col = None
        if table == "daily":
            date_col = "trade_date"
        elif table == "daily_basic":
            date_col = "trade_date"
        elif table in ["income", "balancesheet", "cashflow"]:
            date_col = "end_date"
        elif table == "fina_indicator":
            date_col = "end_date"
        elif table == "trade_cal":
            date_col = "cal_date"
        elif table == "index_daily":
            date_col = "trade_date"
        elif table == "adj_factor":
            date_col = "trade_date"
            
        if date_col:
            date_sql = f"SELECT MIN({date_col}), MAX({date_col}) FROM {table}"
            min_d, max_d = pd.read_sql(date_sql, conn).iloc[0]
            print(f"{table:<20} {count:>12,} {str(min_d)[:10]:>12} {str(max_d)[:10]:>12}")
        else:
            print(f"{table:<20} {count:>12,}")
    except Exception as e:
        print(f"{table:<20} {'错误':>12}: {str(e)[:40]}")

conn.close()
In [ ]:
# 示例查询 1: 查询某股票最近10个交易日数据
query1 = """
SELECT trade_date, open, high, low, close, vol, amount, pct_chg
FROM daily
WHERE ts_code = '000001.SZ'
ORDER BY trade_date DESC
LIMIT 10
"""
conn = get_pg_connection()
df1 = pd.read_sql(query1, conn)
print("平安银行(000001.SZ) 最近10个交易日:")
display(df1)
conn.close()
In [ ]:
# 示例查询 2: 日线行情 + 估值指标联合查询 (使用视图)
query2 = """
SELECT *
FROM v_daily_with_valuation
WHERE ts_code = '000001.SZ'
  AND trade_date >= '2024-01-01'
ORDER BY trade_date DESC
LIMIT 10
"""
conn = get_pg_connection()
df2 = pd.read_sql(query2, conn)
print("平安银行 - 日线+估值:")
display(df2[['trade_date', 'close', 'pct_chg', 'pe', 'pe_ttm', 'pb', 'total_mv']])
conn.close()
In [ ]:
# 示例查询 3: 最新财务指标 Top 20 (按 ROE 排序)
query3 = """
SELECT *
FROM v_latest_financials
WHERE roe IS NOT NULL
  AND roe > 0
ORDER BY roe DESC
LIMIT 20
"""
conn = get_pg_connection()
df3 = pd.read_sql(query3, conn)
print("ROE Top 20:")
display(df3[['ts_code', 'name', 'industry', 'roe', 'roa', 'eps', 'debt_to_assets']])
conn.close()