04-数据更新与调度策略

akshare 是数据聚合抓取工具,接口稳定性依赖上游网站(东方财富、新浪、同花顺、雪球、理杏仁等)。本策略围绕这一本质特征设计,重点处理接口失效、数据延迟、限流、数据漂移等问题。

核心认知:akshare 不是数据源

使用 akshare 必须建立以下认知:

  1. akshare 是抓取层,不是数据源:数据真实来源是上游财经网站,akshare 只是封装了抓取逻辑
  2. 上游接口随时可能变更:网站改版、反爬升级、接口下线都会导致 akshare 接口失效
  3. akshare 版本更新有滞后:上游变更后,akshare 维护者需要时间修复并发布新版本
  4. 返回格式可能漂移:即使接口不报错,返回的字段名、字段顺序、数据类型也可能悄悄变化
  5. 限流策略不可控:上游网站可能随时加强限流,批量抓取容易触发封禁

基于以上认知,更新策略的核心原则是:本地数据库才是可信数据源,akshare 只是获取渠道。所有决策分析基于本地库,不直接依赖 akshare 实时调用。

更新频率

按数据类型分级更新,避免无意义的重复抓取:

数据类型更新频率触发时机理由
指数日线每交易日收盘后 18:00当日行情需次日分析使用
指数估值每交易日收盘后 18:30依赖日线数据,稍后更新
基金日线每交易日收盘后 19:00ETF 与场外基金净值
股票日线每交易日收盘后 18:00同指数
行业日线每交易日收盘后 18:00同指数
宏观指标每月数据发布日次日CPI/PPI/PMI 等按月发布
资金流向每交易日收盘后 18:30北向/两融数据
基金持仓每季度季报披露截止后季报披露有 15 个工作日窗口
股票财报每季度财报披露季后年报 4/30 前,一季报 4/30 前,半年报 8/31 前,三季报 10/31 前
标的维度每月月初新上市/退市/更名

更新窗口错开,避免单次更新压力过大。18:00-19:30 完成日级数据,20:00 后可选触发低频数据校验。

调度方式

本项目环境为 Windows,采用 Windows 任务计划程序调度:

任务计划程序配置

创建 .bat 批处理文件作为任务计划程序的入口:

@echo off
cd /d "C:\Users\nobattle\Desktop\个人\趋势定投之路\05-数据库\02-数据脚本"
python update_all.py --mode daily >> logs\update_%date:~0,4%%date:~5,2%%date:~8,2%.log 2>&1

在任务计划程序中创建任务:

  • 触发器:每日 18:00
  • 操作:启动 update_daily.bat
  • 起始于:05-数据库\02-数据脚本\
  • 选项:不管用户是否登录都运行,失败后每 5 分钟重试,最多 3 次

调度模式

update_all.py 支持三种模式:

# 日常增量更新(默认)
python update_all.py --mode daily

# 全量重建(首次建库或数据修复)
python update_all.py --mode full

# 仅校验,不写入
python update_all.py --mode validate
  • daily:只更新缺失的日期,速度快,每日用
  • full:从零拉取全量历史,首次建库或数据污染后用
  • validate:跑校验规则,不修改数据,用于排查异常

容错与重试机制

akshare 接口调用失败是常态,fetcher 层统一封装重试逻辑:

import time
import logging
from functools import wraps

logger = logging.getLogger(__name__)

def fetch_with_retry(func, max_retries=3, base_delay=1.0,
                     backoff_factor=2.0, sleep_after=0.5):
    """
    带 retry 的 akshare 调用封装。
    - max_retries: 最大重试次数
    - base_delay: 首次重试延迟(秒)
    - backoff_factor: 退避因子,每次延迟乘以此系数
    - sleep_after: 成功后 sleep,避免触发限流
    """
    @wraps(func)
    def wrapper(*args, **kwargs):
        last_error = None
        for attempt in range(max_retries):
            try:
                result = func(*args, **kwargs)
                time.sleep(sleep_after)
                return result
            except Exception as e:
                last_error = e
                delay = base_delay * (backoff_factor ** attempt)
                logger.warning(
                    f"[{func.__name__}] 第 {attempt+1} 次失败: {e}, "
                    f"{delay}s 后重试"
                )
                time.sleep(delay)
        logger.error(f"[{func.__name__}] {max_retries} 次重试均失败: {last_error}")
        raise last_error
    return wrapper

关键设计:

  • 指数退避:1s → 2s → 4s,避免雪崩式重试
  • 失败记录:每次失败写入日志,便于排查
  • 限流保护:成功后 sleep 0.5s,降低被封风险
  • 不静默吞错:重试耗尽后抛出异常,由上层决定是否记录并继续下一个标的

失败记录与恢复

每次更新维护一张失败记录表,记录哪些标的、哪些接口失败了:

CREATE TABLE IF NOT EXISTS update_failure_log (
    id INTEGER PRIMARY KEY AUTOINCREMENT,
    run_id TEXT NOT NULL,               -- 本次更新批次 ID(时间戳)
    fetcher TEXT NOT NULL,              -- 失败的 fetcher 名
    target_code TEXT,                   -- 失败的标的代码
    api_name TEXT,                      -- akshare 接口名
    error_type TEXT,                    -- 错误类型:network/parse/empty/unknown
    error_message TEXT,                 -- 错误详情
    occurred_at TEXT NOT NULL           -- 发生时间
);

CREATE INDEX IF NOT EXISTS idx_failure_run
    ON update_failure_log(run_id);
CREATE INDEX IF NOT EXISTS idx_failure_target
    ON update_failure_log(target_code);

update_all.py 每次运行生成一个 run_id(如 20260714_180000),所有失败记录归属到该 ID。下次运行时:

  1. 先查询上次失败的标的列表
  2. 优先重试这些标的
  3. 仍然失败的,升级告警(写入 update_alert 表或日志文件)

失败记录保留 90 天,定期清理。

断点续传

数据获取过程可能因网络中断、程序崩溃、手动终止等原因中断。简单"从头重新拉取"在 full 模式下代价极高(全历史重新拉取可能需要数十分钟)。本项目通过三层机制实现高效断点续传。

第一层:数据驱动续传

每个标的更新前查询本地最新交易日,只拉取该日期之后的增量数据:

from fetcher_common import get_last_local_date, compute_start_date

local_latest = get_last_local_date(conn, "index_daily", "index_code", "000300")
# local_latest = "2026-07-13"(昨日已更新)
start_date = compute_start_date(local_latest, default_days=30, mode="daily")
# start_date = "2026-07-14"(只拉今天的数据)

compute_start_date 的逻辑:

场景本地最新日期返回 start_date
首次拉取(无本地数据)Nonedaily: 30 天前;full: 1990-01-01
增量更新2026-07-132026-07-14(+1 天)
已是最新今天今天(跳过拉取)

即使没有进度表,仅靠数据本身也能判断从哪里开始,保证幂等性。

第二层:进度表

update_progress 表记录每个 fetcher × 标的的处理状态,用于跳过已完成标的、重试失败标的:

CREATE TABLE update_progress (
    fetcher_name TEXT NOT NULL,
    target_code TEXT NOT NULL,
    last_trade_date TEXT,
    last_updated_at TEXT NOT NULL,
    status TEXT NOT NULL,           -- pending/in_progress/completed/failed
    run_id TEXT,
    PRIMARY KEY (fetcher_name, target_code)
);

状态流转:

pending -> in_progress -> completed(成功)
                     \-> failed(失败)
failed -> in_progress(下次重试)-> completed / failed

update_all.py 每次运行时:

  1. 查询 status='in_progress' 的标的(上次中断的),优先继续
  2. 查询 status='failed' 的标的,重试
  3. 跳过 status='completed' 且数据新鲜的标的
  4. 处理 status='pending' 或无记录的新标的
from fetcher_common import get_pending_targets

# 获取待更新的标的(跳过已完成的)
pending = get_pending_targets(conn, "fetch_index_daily", WATCH_INDEXES)
for index_code, symbol_em, name in pending:
    fetch_index_daily(conn, run_id, index_code, symbol_em, mode=mode)

第三层:事务原子性

每个标的的数据写入用事务包裹,保证不出现半成品:

try:
    conn.execute("BEGIN")
    conn.executemany("INSERT OR REPLACE INTO index_daily ...", records)
    conn.commit()
    mark_progress_done(conn, "fetch_index_daily", index_code, last_date, run_id)
except Exception:
    conn.execute("ROLLBACK")
    mark_progress_failed(conn, "fetch_index_daily", index_code, run_id)

中断时事务自动回滚,不会留下部分写入的数据。下次运行从同一标的继续。

中断恢复示例

假设 full 模式更新 6 个指数,更新到第 4 个时中断:

运行 1(中断):
  000300 -> completed(全历史数据写入成功)
  000905 -> completed
  000852 -> completed
  399006 -> in_progress(写入到 2023 年时中断)
  000016 -> pending(未开始)
  000688 -> pending

运行 2(恢复):
  000300 -> 跳过(completed,本地已有最新)
  000905 -> 跳过
  000852 -> 跳过
  399006 -> 从 2023 年继续(数据驱动续传)
  000016 -> 全量拉取(无本地数据)
  000688 -> 全量拉取

恢复运行只拉取必要部分,避免重复工作。

强制全量重拉

特殊情况下(如数据污染、接口格式变化导致历史数据错误)需要强制重拉,在 update_all.py 中传入 --force 参数:

python update_all.py --mode full --force

--force 会:

  1. 忽略 update_progress 中的 completed 状态
  2. 不使用本地最新日期做增量判断
  3. 从 1990-01-01 起全量拉取并覆盖

数据快照与回滚

每次 full 模式更新前自动创建数据库快照:

import shutil
from datetime import datetime
from pathlib import Path

def backup_database(db_path: str, backup_dir: str):
    """更新前备份当前数据库"""
    ts = datetime.now().strftime("%Y%m%d_%H%M%S")
    backup_path = Path(backup_dir) / f"trend_invest_{ts}.db"
    shutil.copy2(db_path, backup_path)
    # 清理超过 30 天的旧备份
    cutoff = datetime.now() - timedelta(days=30)
    for old in Path(backup_dir).glob("trend_invest_*.db"):
        if datetime.fromtimestamp(old.stat().st_mtime) < cutoff:
            old.unlink()
    return backup_path

快照策略:

  • 保留策略:最近 30 天的每日快照
  • 存储位置05-数据库/03-数据文件/backups/(加入 .gitignore)
  • 触发时机full 模式更新前、schema 迁移前、手动触发
  • 回滚方式:关闭数据库连接 → 用快照覆盖当前库 → 重启

daily 模式不创建快照(增量更新风险低),但可通过 --backup 参数强制创建。

接口失效应对

当 akshare 接口失效时,按以下流程处理:

第一步:识别失效

通过 update_failure_log 识别:

  • 同一接口在 3 次连续更新中均失败
  • 失败类型为 parse(解析失败,通常是返回格式变化)或 empty(返回空)

第二步:切换备用接口

每个关键接口在 03-数据获取接口与映射 中标注了备用接口。fetcher 层实现自动切换:

def fetch_index_daily(symbol: str, start_date: str, end_date: str):
    """主接口失败时自动切备用接口"""
    primary = lambda: ak.stock_zh_index_daily_em(symbol=symbol)
    fallback = lambda: ak.stock_zh_index_daily(symbol=symbol)

    for fetcher in [primary, fallback]:
        try:
            df = fetch_with_retry(fetcher)()
            if df is not None and len(df) > 0:
                return normalize_columns(df)  # 字段标准化
        except Exception as e:
            logger.warning(f"接口 {fetcher.__name__} 失败: {e}")
            continue

    raise RuntimeError(f"指数 {symbol} 日线获取全部失败")

第三步:升级 akshare

如果备用接口也失效,检查 akshare 是否有新版本:

pip install --upgrade akshare

升级后跑 validate 模式验证接口恢复。注意升级可能引入新接口签名变化,需同步检查 fetcher 层。

第四步:手工补数据

接口长时间失效时,从其他渠道手工获取数据(如网站导出 CSV),通过 import_csv.py 工具导入。手工补数据需在 update_failure_log 标注 manual_fix

数据漂移检测

akshare 返回的字段可能悄悄变化,fetcher 层做字段校验:

EXPECTED_COLUMNS = {
    "stock_zh_a_hist": ["日期", "开盘", "收盘", "最高", "最低",
                        "成交量", "成交额", "振幅", "涨跌幅", "涨跌额", "换手率"],
}

def validate_columns(df, api_name: str):
    """校验返回字段是否符合预期"""
    expected = EXPECTED_COLUMNS.get(api_name, [])
    actual = list(df.columns)
    missing = set(expected) - set(actual)
    extra = set(actual) - set(expected)
    if missing:
        logger.warning(
            f"[{api_name}] 字段缺失: {missing}; 实际字段: {actual}"
        )
        return False
    if extra:
        logger.info(f"[{api_name}] 新增字段: {extra}")
    return True

字段校验失败不中断更新,但写入 update_failure_log,便于发现接口变化。

更新流程总览

update_all.py 的完整流程:

1. 生成 run_id,记录开始时间
2. 创建数据库快照(full 模式)
3. 按优先级顺序更新:
   a. 标的维度表(dim_*)
   b. 行情数据(index_daily/fund_daily/stock_daily/industry_daily)
   c. 估值数据(index_valuation/stock_valuation)
   d. 基本面(fund_holdings/stock_financial)
   e. 宏观与资金(macro_indicator/money_flow)
4. 每个 fetcher 内部:
   - 读取上次失败的标的列表,优先重试
   - 对每个标的调用 fetch_with_retry
   - 字段校验 → 数据清洗 → 写入数据库
   - 失败记录到 update_failure_log
5. 更新完成后:
   - 写入 db_meta 的 last_updated
   - 运行数据质量校验(见 05-数据质量与校验规则)
   - 输出本次更新摘要(成功/失败数)
6. 清理 30 天前的快照与失败日志

完整脚本见 02-数据脚本/update_all.py

监控与告警

更新异常需要主动感知,不能等分析时才发现数据缺失:

日志文件

每次更新输出日志到 05-数据库/02-数据脚本/logs/update_YYYYMMDD.log,包含:

  • 每个接口的调用结果(成功/失败/重试次数)
  • 每个标的的更新行数
  • 耗时统计
  • 错误详情

简易告警

update_all.py 结束时检查失败率,超过阈值则输出告警:

failure_rate = failed_count / total_count
if failure_rate > 0.1:  # 超过 10% 失败
    logger.error(f"告警:本次更新失败率 {failure_rate:.1%},请检查")
    # 可选:写入 alert 文件或发送通知

告警方式保持简单(日志 + 文件标记),不引入复杂通知系统。如果未来需要邮件/微信通知,在告警模块扩展即可。