04-数据更新与调度策略
akshare 是数据聚合抓取工具,接口稳定性依赖上游网站(东方财富、新浪、同花顺、雪球、理杏仁等)。本策略围绕这一本质特征设计,重点处理接口失效、数据延迟、限流、数据漂移等问题。
核心认知:akshare 不是数据源
使用 akshare 必须建立以下认知:
- akshare 是抓取层,不是数据源:数据真实来源是上游财经网站,akshare 只是封装了抓取逻辑
- 上游接口随时可能变更:网站改版、反爬升级、接口下线都会导致 akshare 接口失效
- akshare 版本更新有滞后:上游变更后,akshare 维护者需要时间修复并发布新版本
- 返回格式可能漂移:即使接口不报错,返回的字段名、字段顺序、数据类型也可能悄悄变化
- 限流策略不可控:上游网站可能随时加强限流,批量抓取容易触发封禁
基于以上认知,更新策略的核心原则是:本地数据库才是可信数据源,akshare 只是获取渠道。所有决策分析基于本地库,不直接依赖 akshare 实时调用。
更新频率
按数据类型分级更新,避免无意义的重复抓取:
更新窗口错开,避免单次更新压力过大。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。下次运行时:
- 先查询上次失败的标的列表
- 优先重试这些标的
- 仍然失败的,升级告警(写入
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 的逻辑:
即使没有进度表,仅靠数据本身也能判断从哪里开始,保证幂等性。
第二层:进度表
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 每次运行时:
- 查询
status='in_progress' 的标的(上次中断的),优先继续
- 查询
status='failed' 的标的,重试
- 跳过
status='completed' 且数据新鲜的标的
- 处理
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 会:
- 忽略
update_progress 中的 completed 状态
- 不使用本地最新日期做增量判断
- 从 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 文件或发送通知
告警方式保持简单(日志 + 文件标记),不引入复杂通知系统。如果未来需要邮件/微信通知,在告警模块扩展即可。