汇总三轮未提交的开发(每轮均在本机 MariaDB + 真实浏览器上验证):
1) 股息率案例(全市场股息率最高 n 只,默认 20,每 m 月择股)
- 新增日频估值表 daily_basic + 迁移;股息率因子(dv_ratio / dividend_yield / TTM)
- 名称历史表 stock_name_history:剔除 ST 按**择股日当时名称**判定,消除
「曾高股息后 ST」的股息陷阱(实测 3.70pp 偏差)
- 区间择股/调仓双周期(m 择股 / y 调仓)、指数成分与白名单、停牌近似剔除
- 复权因子口径核对(4,164,742 行、缺失 0.0%)、收盘价成交与涨跌停拦单
- 案例实测:2020-01-01~2026-09-04 总收益 +24.86%(年化 3.52%、回撤 -28.58%)
2) 策略库与前端统一
- strategy 表 + CRUD/PUT 原地更新 + `describe_strategy` 按 spec 真实推导
「一句话说明 + 计算公式 + 执行步骤 + 注意事项」(与引擎实执行规则同源)
- 任何出现股票代码处都成对显示名称且可点击进个股页
- 全站图表基座统一 TradingView Lightweight Charts(ECharts 依赖、
锁文件、组件与文档标注一并清除),买卖点标记只落在真实交易日上
3) 回测存档完整化(可往复查看)
- 同步端点(POST /api/backtests、/api/factor-tests)此前完全不落库 → 现在同样归档,
归档 id 经响应头 X-Experiment-Id 返回(不破坏 response_model)
- data_version 首次真实写入(数据快照指纹:最新交易日 + 各表规模)
- 个股收益曲线默认**全量保存**(此前硬截断 60 只);超出体积预算才裁剪,
并写 archive_meta(机器可读)+ unimplemented(人可读)如实标注
- 列表 kind/q 过滤 + X-Total-Count(此前 limit=50 静默截断)、DELETE 归档
- 只读归档页 /experiments/{id}(Server Component,SSR 直出**选股条件**与
**交易执行依据**);结果视图按 kind 分发(backtest/factor_test/selection),
非回测归档不套用回测口径
- 新增 CLI:prune_experiments(保留策略,默认 dry-run)、
restore_experiment_from_job(从 Job 副本按原 id 重建被删的历史归档,默认 dry-run)
门禁:pytest 388 passed、ruff All checks passed、tsc 0 错误、图表单测 7 passed、
next build 成功、契约脚本 verify_strategy_workspace 59/59(含按 kind 逐类验证归档页)。
643 lines
26 KiB
Python
643 lines
26 KiB
Python
"""Phase 1 数据同步 CLI(Tushare 首选 → SQLite,新浪校验兜底)。
|
||
|
||
用法(cd backend):
|
||
uv run python -m app.cli.sync basic
|
||
uv run python -m app.cli.sync calendar --start 20240101 --end 20241231
|
||
uv run python -m app.cli.sync daily --symbols 600519.SH,000001.SZ --start 20240101
|
||
uv run python -m app.cli.sync daily --all --start 20240101 # 全市场
|
||
uv run python -m app.cli.sync financial --all # 财务指标(增量)
|
||
uv run python -m app.cli.sync financial --all --full # 财务指标(强制全量重拉)
|
||
uv run python -m app.cli.sync verify --symbol 600519.SH # 新浪交叉验证
|
||
uv run python -m app.cli.sync daily_basic --start 20200101 # 每日指标(股息率等)
|
||
|
||
增量与兜底:
|
||
- daily --resume:从本地最新交易日续传(已有);Tushare 失败时走新浪校验兜底,
|
||
只有「两源重叠历史一致」才用新浪补本地缺失交易日(source=sina/前复权)。
|
||
- financial:默认增量——本地已含最新应披露报告期则跳过;Tushare 失败时新浪
|
||
数据须通过「两边一致」校验(重叠报告期 eps/销售毛利率逐期一致)才允许补入
|
||
本地缺失键(source=sina)。失败股票留待下轮重跑补齐,不会静默导入未核验数据。
|
||
- 每次拉取写入 sync_log 审计(来源 / 成功与否 / 行数 / 区间),禁止静默切源。
|
||
|
||
本模块是组装层(composition root):在此装配 Provider / Repository / Session,
|
||
业务逻辑在 application.services.data_sync,业务层仍只依赖抽象。
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import argparse
|
||
import sys
|
||
import time
|
||
from datetime import date, datetime, timedelta
|
||
|
||
from sqlalchemy import select
|
||
|
||
from app.application.services.data_sync import (
|
||
DailyBasicSyncer,
|
||
DailySymbolResult,
|
||
FinancialSymbolResult,
|
||
NameHistorySyncer,
|
||
VerifiedDailySyncer,
|
||
VerifiedFinancialSyncer,
|
||
)
|
||
from app.core.config import get_settings
|
||
from app.infrastructure.data_sources.errors import DataSourceError
|
||
from app.infrastructure.data_sources.sina import SinaProvider
|
||
from app.infrastructure.data_sources.tushare import TushareProvider
|
||
from app.infrastructure.persistence.sqlalchemy.models.market import StockModel
|
||
from app.infrastructure.persistence.sqlalchemy.repositories.index_impl import (
|
||
SqlAlchemyIndexConstituentRepository,
|
||
)
|
||
from app.infrastructure.persistence.sqlalchemy.repositories.market_impl import (
|
||
SqlAlchemyAdjustFactorRepository,
|
||
SqlAlchemyDailyBarRepository,
|
||
SqlAlchemyDailyBasicRepository,
|
||
SqlAlchemyFinancialRepository,
|
||
SqlAlchemyStockNameHistoryRepository,
|
||
SqlAlchemyStockRepository,
|
||
SqlAlchemySyncLogRepository,
|
||
SqlAlchemyTradingCalendarRepository,
|
||
)
|
||
from app.infrastructure.persistence.sqlalchemy.session import SessionLocal
|
||
|
||
_DATE_FMT = "%Y%m%d"
|
||
|
||
|
||
def _parse_day(text: str) -> date:
|
||
return datetime.strptime(text, _DATE_FMT).date()
|
||
|
||
|
||
def _failover_provider(session):
|
||
"""Tushare 首选 + 新浪兜底(basic/calendar 用;daily/financial 走校验兜底服务)。
|
||
|
||
FailoverProvider 每次尝试写 sync_log(AGENT.md §7)。能力矩阵:新浪仅提供
|
||
日线/财务,basic/calendar 新浪不支持 → 抛错保留单源语义,日志可见。
|
||
"""
|
||
from app.infrastructure.data_sources.failover import FailoverProvider
|
||
from app.infrastructure.data_sources.sina import SinaProvider
|
||
|
||
audit_repo = SqlAlchemySyncLogRepository(session)
|
||
primary = TushareProvider(token=get_settings().tushare_token)
|
||
return FailoverProvider(primary, fallback=SinaProvider(), audit=audit_repo.add)
|
||
|
||
|
||
def _session_ctx():
|
||
return SessionLocal()
|
||
|
||
|
||
def cmd_basic(args) -> int:
|
||
"""股票基础信息。--include-delisted 同时拉取已退市/暂停上市(幸存者偏差修正)。"""
|
||
statuses = ["L", "P", "D"] if getattr(args, "include_delisted", False) else ["L"]
|
||
with _session_ctx() as session:
|
||
provider = _failover_provider(session)
|
||
repo = SqlAlchemyStockRepository(session)
|
||
total = 0
|
||
failed: list[str] = []
|
||
for st in statuses:
|
||
try:
|
||
stocks = provider.get_stock_basic(st)
|
||
except DataSourceError as exc:
|
||
# failover 会把备用源(新浪)的 NotSupported 包装成 DataSourceError,
|
||
# 因此这里必须捕获 DataSourceError 而不是 DataSourceNotSupported,
|
||
# 否则 --include-delisted 在主源抖动时会整体抛栈退出(L 已写、P/D 静默缺失)。
|
||
print(f"[error] list_status={st} 拉取失败:{exc}", file=sys.stderr)
|
||
failed.append(st)
|
||
continue
|
||
touched = repo.upsert_many(stocks)
|
||
session.commit()
|
||
total += touched
|
||
n_delisted = sum(1 for s in stocks if s.delist_date is not None)
|
||
print(
|
||
f"[basic] list_status={st} 拉取 {len(stocks)} 只(含 delist_date {n_delisted} 只),"
|
||
f"落库 {touched} 条"
|
||
)
|
||
if len(statuses) > 1:
|
||
print(
|
||
"[basic] 已退市股票已入库;其历史行情需另行同步(否则回测仍无法使用):\n"
|
||
" 请执行 sync daily --symbols <退市代码逗号分隔> --start 20200101"
|
||
)
|
||
print(f"[basic] 合计落库 {total} 条")
|
||
if failed:
|
||
print(
|
||
f"[error] 以下 list_status 拉取失败:{','.join(failed)};"
|
||
"退市股缺失会让回测重新出现幸存者偏差,请重跑",
|
||
file=sys.stderr,
|
||
)
|
||
return 1
|
||
return 0
|
||
|
||
|
||
def cmd_calendar(args) -> int:
|
||
start = _parse_day(args.start)
|
||
end = _parse_day(args.end)
|
||
with _session_ctx() as session:
|
||
provider = _failover_provider(session)
|
||
days = provider.get_trade_cal(start, end)
|
||
repo = SqlAlchemyTradingCalendarRepository(session)
|
||
touched = repo.upsert_many(days)
|
||
session.commit()
|
||
open_days = sum(1 for d in days if d.is_open)
|
||
print(f"[calendar] {start}~{end} 共 {len(days)} 条(交易日 {open_days}),落库 {touched} 条")
|
||
return 0
|
||
|
||
|
||
def cmd_namechange(args) -> int:
|
||
"""股票名称变更历史同步(时点 ST 判定依据;按自然年分片)。"""
|
||
start = _parse_day(args.start) if args.start else date(1990, 1, 1)
|
||
end = _parse_day(args.end) if args.end else date.today()
|
||
if start > end:
|
||
print("[error] --start 不能晚于 --end", file=sys.stderr)
|
||
return 2
|
||
|
||
with _session_ctx() as session:
|
||
audit_repo = SqlAlchemySyncLogRepository(session)
|
||
primary = TushareProvider(token=get_settings().tushare_token)
|
||
repo = SqlAlchemyStockNameHistoryRepository(session)
|
||
syncer = NameHistorySyncer(primary=primary, repo=repo, audit=audit_repo.add)
|
||
|
||
total_span = (end.year - start.year) + 1
|
||
print(f"[namechange] {start}~{end} 共 {total_span} 个年度分片")
|
||
ok = failed = written = fetched = 0
|
||
t0 = time.time()
|
||
|
||
def _on(idx: int, total: int, chunk_start: date) -> None:
|
||
_progress(idx, total, t0, f"{chunk_start.year} 年")
|
||
|
||
results = syncer.sync_range(start, end, on_progress=_on)
|
||
for res in results:
|
||
if res.status == "ok":
|
||
ok += 1
|
||
written += res.rows_written
|
||
fetched += res.rows_fetched
|
||
else:
|
||
failed += 1
|
||
for note in res.notes:
|
||
print(f"\n[warn] {res.start}~{res.end}: {note}", file=sys.stderr)
|
||
session.commit()
|
||
lo, hi = repo.namechange_dates()
|
||
print(
|
||
f"\n[namechange] 分片成功 {ok} / 失败 {failed};"
|
||
f"拉取 {fetched} 行、写入 {written} 行;"
|
||
f"本地生效起点 {lo} ~ {hi}"
|
||
)
|
||
return 1 if failed else 0
|
||
|
||
|
||
def _symbols_of(args) -> list[str]:
|
||
if getattr(args, "all", False):
|
||
with _session_ctx() as session:
|
||
symbols = list(session.scalars(select(StockModel.symbol).order_by(StockModel.symbol)))
|
||
if not symbols:
|
||
print("[error] stock 表为空,请先运行:python -m app.cli.sync basic")
|
||
sys.exit(2)
|
||
return symbols
|
||
return [s.strip() for s in args.symbols.split(",") if s.strip()]
|
||
|
||
|
||
def _stock_names(session, symbols: list[str]) -> dict[str, str]:
|
||
"""一次性取出股票名称(进度描述用);批量查询避开 SQLite 变量上限。"""
|
||
names: dict[str, str] = {}
|
||
for i in range(0, len(symbols), 500):
|
||
chunk = symbols[i : i + 500]
|
||
rows = session.execute(
|
||
select(StockModel.symbol, StockModel.name).where(StockModel.symbol.in_(chunk))
|
||
)
|
||
names.update({sym: nm for sym, nm in rows})
|
||
return names
|
||
|
||
|
||
def _warn_notes(notes: list[str]) -> None:
|
||
for note in notes:
|
||
print(f" [warn] {note}", file=sys.stderr)
|
||
|
||
|
||
def cmd_daily(args) -> int:
|
||
from sqlalchemy import func
|
||
|
||
from app.infrastructure.persistence.sqlalchemy.models.market import StockDailyModel
|
||
|
||
symbols = _symbols_of(args)
|
||
start = _parse_day(args.start) if args.start else date(2005, 1, 1)
|
||
end = _parse_day(args.end) if args.end else date.today()
|
||
started = time.monotonic()
|
||
n_ok = n_sina = n_failed = n_skip = 0
|
||
rows_tushare = rows_sina = 0
|
||
with _session_ctx() as session:
|
||
names = _stock_names(session, symbols)
|
||
audit = SqlAlchemySyncLogRepository(session).add
|
||
syncer = VerifiedDailySyncer(
|
||
primary=TushareProvider(token=get_settings().tushare_token),
|
||
fallback=SinaProvider(),
|
||
bars=SqlAlchemyDailyBarRepository(session),
|
||
factors=SqlAlchemyAdjustFactorRepository(session),
|
||
audit=audit,
|
||
)
|
||
bar_repo = SqlAlchemyDailyBarRepository(session)
|
||
# 增量基准:本地数据已到该日期即视为「已最新」,resume 时不再调 API
|
||
global_latest = (
|
||
session.scalar(select(func.max(StockDailyModel.trade_date))) if args.resume else None
|
||
)
|
||
for i, symbol in enumerate(symbols, start=1):
|
||
begin = start
|
||
if args.resume:
|
||
latest = bar_repo.latest_date(symbol)
|
||
if latest is not None:
|
||
if global_latest is not None and latest >= global_latest:
|
||
n_skip += 1 # 已同步到本地最新交易日,无需续拉
|
||
continue
|
||
begin = max(begin, latest + timedelta(days=1))
|
||
if begin > end:
|
||
n_skip += 1 # 无待拉区间(如区间已含在本地)
|
||
continue
|
||
if getattr(args, "sleep", 0) > 0:
|
||
time.sleep(args.sleep)
|
||
res: DailySymbolResult = syncer.sync_symbol(symbol, begin, end)
|
||
session.commit() # 逐只落库:中断/报错只丢当前一只,重跑增量续传
|
||
if res.status == "ok":
|
||
n_ok += 1
|
||
rows_tushare += res.bars_written
|
||
elif res.status == "sina":
|
||
n_sina += 1
|
||
rows_sina += res.bars_written
|
||
elif res.status == "failed":
|
||
n_failed += 1
|
||
_warn_notes(res.notes)
|
||
if i % 100 == 0:
|
||
name = names.get(symbol, "")
|
||
print(
|
||
f" ... {i}/{len(symbols)} {symbol} {name}: "
|
||
f"累计 tushare {rows_tushare} 根 + 新浪补缺 {rows_sina} 根;"
|
||
f"成功 {n_ok} / 新浪 {n_sina} / 失败待重试 {n_failed}"
|
||
)
|
||
elapsed = time.monotonic() - started
|
||
detail = (
|
||
f"[daily] {len(symbols)} 只股票:成功 {n_ok} / 新浪校验补缺 {n_sina} / "
|
||
f"已最新跳过 {n_skip} / 失败待重试 {n_failed}"
|
||
)
|
||
if args.resume:
|
||
detail += f"(本地最新 {global_latest})"
|
||
detail += f";写入 {rows_tushare} 根(tushare 不复权)+ {rows_sina} 根(sina 前复权),耗时 {elapsed:.0f}s"
|
||
print(detail)
|
||
return 0
|
||
|
||
|
||
def _fin_progress_line(i: int, n: int, symbol: str, name: str, res: FinancialSymbolResult) -> str:
|
||
"""financial 逐只进度行:结果 + 导入内容简单描述(报告期/公告区间、来源)。"""
|
||
head = f"[financial {i}/{n}] {symbol} {name or ''}".rstrip()
|
||
if res.status == "skip":
|
||
return f"{head}:已最新,跳过(增量)"
|
||
if res.status == "failed":
|
||
return f"{head}:失败待重试(tushare 失败;新浪源 {'未通过校验' if res.source == 'sina' else '不可用'})"
|
||
if res.status == "sina":
|
||
return (
|
||
f"{head}:tushare 失败 → 新浪校验通过,补入 {res.written} 行(source=sina)"
|
||
+ _fin_span(res)
|
||
)
|
||
# status == ok(tushare 成功)
|
||
if res.written:
|
||
updated = f",覆盖更新 {res.updated} 行" if res.updated else ""
|
||
return f"{head}:tushare 返回 {res.fetched} 行 → 新增 {res.written} 行{updated}" + _fin_span(res)
|
||
return f"{head}:tushare 返回 {res.fetched} 行,均已在库,无新增"
|
||
|
||
|
||
def _fin_span(res: FinancialSymbolResult) -> str:
|
||
if not res.written or res.report_first is None:
|
||
return ""
|
||
if res.announce_first is None or res.announce_last is None:
|
||
return ""
|
||
return (
|
||
f";报告期 {res.report_first.isoformat()}~{res.report_last.isoformat()}"
|
||
f"(公告 {res.announce_first.isoformat()}~{res.announce_last.isoformat()})"
|
||
)
|
||
|
||
|
||
def cmd_financial(args) -> int:
|
||
symbols = _symbols_of(args)
|
||
started = time.monotonic()
|
||
n_ok = n_sina = n_failed = n_skip = 0
|
||
rows_tushare = rows_sina = 0
|
||
with _session_ctx() as session:
|
||
names = _stock_names(session, symbols)
|
||
audit = SqlAlchemySyncLogRepository(session).add
|
||
syncer = VerifiedFinancialSyncer(
|
||
primary=TushareProvider(token=get_settings().tushare_token),
|
||
fallback=SinaProvider(),
|
||
repo=SqlAlchemyFinancialRepository(session),
|
||
audit=audit,
|
||
)
|
||
for i, symbol in enumerate(symbols, start=1):
|
||
if getattr(args, "sleep", 0) > 0:
|
||
time.sleep(args.sleep)
|
||
res: FinancialSymbolResult = syncer.sync_symbol(symbol, force_full=args.full)
|
||
session.commit() # 逐只落库:中断只丢当前一只,重跑增量续传
|
||
print(_fin_progress_line(i, len(symbols), symbol, names.get(symbol, ""), res))
|
||
_warn_notes(res.notes)
|
||
if res.status == "ok":
|
||
n_ok += 1
|
||
rows_tushare += res.written
|
||
elif res.status == "sina":
|
||
n_sina += 1
|
||
rows_sina += res.written
|
||
elif res.status == "failed":
|
||
n_failed += 1
|
||
elif res.status == "skip":
|
||
n_skip += 1
|
||
elapsed = time.monotonic() - started
|
||
mode = "全量重拉(--full)" if args.full else "增量"
|
||
print(
|
||
f"[financial] 共 {len(symbols)} 只({mode}):成功 {n_ok} / 新浪校验兜底 {n_sina} / "
|
||
f"已最新跳过 {n_skip} / 失败待重试 {n_failed};"
|
||
f"合计写入 {rows_tushare + rows_sina} 行(tushare {rows_tushare} + sina {rows_sina}),"
|
||
f"耗时 {elapsed:.0f}s"
|
||
)
|
||
if n_failed:
|
||
print(
|
||
" [tip] 失败股票未写入未核验数据,重跑本命令即可续传补齐;"
|
||
"若因频率超限,可用 --sleep 加大间隔(如 --sleep 60)分多次跑。",
|
||
file=sys.stderr,
|
||
)
|
||
return 0
|
||
|
||
|
||
def cmd_index_weight(args) -> int:
|
||
"""同步指数历史成分(Tushare index_weight;新浪不支持 → failover 审计留痕)。"""
|
||
|
||
code = args.code
|
||
with _session_ctx() as session:
|
||
provider = _failover_provider(session)
|
||
try:
|
||
rows = provider.get_index_weight(code)
|
||
except DataSourceError as exc:
|
||
print(f"[index_weight] {code} 失败:{exc}")
|
||
return 1
|
||
repo = SqlAlchemyIndexConstituentRepository(session)
|
||
touched = repo.upsert_many(rows)
|
||
latest = repo.latest_date(code)
|
||
session.commit()
|
||
print(
|
||
f"[index_weight] {code} 拉取 {len(rows)} 期成分行,落库 {touched} 条"
|
||
f",最新快照 {latest}(as_of 查询见 Universe.index_code)"
|
||
)
|
||
return 0
|
||
|
||
|
||
def cmd_verify(args) -> int:
|
||
"""新浪交叉验证:取新浪最新前复权收盘,与本地最新交易日对照。
|
||
|
||
注意:新浪为前复权口径,数值不直接等于本地不复权收盘,
|
||
本命令仅用于确认新浪可用性 / 最新交易日,不把新浪数据并入主库。
|
||
"""
|
||
from app.infrastructure.persistence.sqlalchemy.repositories.market_impl import (
|
||
SqlAlchemyDailyBarRepository,
|
||
)
|
||
|
||
sina = SinaProvider()
|
||
end = date.today()
|
||
start = end - timedelta(days=20)
|
||
try:
|
||
bars = sina.get_daily(args.symbol, start, end)
|
||
except DataSourceError as exc:
|
||
print(f"[verify] 新浪不可用: {exc}", file=sys.stderr)
|
||
return 1
|
||
if not bars:
|
||
print(f"[verify] 新浪最近无数据({args.symbol})")
|
||
return 1
|
||
latest = max(bars, key=lambda b: b.trade_date)
|
||
with _session_ctx() as session:
|
||
local = SqlAlchemyDailyBarRepository(session).latest_date(args.symbol)
|
||
print(
|
||
f"[verify] {args.symbol}: 新浪最新 {latest.trade_date} 收盘(前复权) {latest.close};"
|
||
f"本地最新交易日 {local}"
|
||
)
|
||
return 0
|
||
|
||
|
||
def cmd_export(args) -> int:
|
||
"""把 SQLite 日线按年导出为 Parquet(data/parquet/stock_daily/<year>.parquet)。"""
|
||
from pathlib import Path
|
||
|
||
import pandas as pd
|
||
|
||
from app.infrastructure.persistence.sqlalchemy.models.market import StockDailyModel
|
||
|
||
settings = get_settings()
|
||
out_root = settings.storage.get("parquet_dir") or Path("data/parquet")
|
||
out_root.mkdir(parents=True, exist_ok=True)
|
||
total = 0
|
||
years = [int(y) for y in (args.years or "").split(",") if y.strip()] or None
|
||
with _session_ctx() as session:
|
||
all_bars = session.execute(
|
||
select(StockDailyModel).order_by(StockDailyModel.trade_date)
|
||
).scalars()
|
||
frame = pd.DataFrame(
|
||
[
|
||
{
|
||
"symbol": b.symbol,
|
||
"trade_date": b.trade_date,
|
||
"open": float(b.open) if b.open is not None else None,
|
||
"high": float(b.high) if b.high is not None else None,
|
||
"low": float(b.low) if b.low is not None else None,
|
||
"close": float(b.close) if b.close is not None else None,
|
||
"volume": float(b.volume) if b.volume is not None else None,
|
||
"amount": float(b.amount) if b.amount is not None else None,
|
||
}
|
||
for b in all_bars
|
||
]
|
||
)
|
||
if frame.empty:
|
||
print("[export] 无日线数据,请先运行 sync daily")
|
||
return 0
|
||
frame["trade_date"] = pd.to_datetime(frame["trade_date"])
|
||
out_dir = out_root / "stock_daily"
|
||
out_dir.mkdir(parents=True, exist_ok=True)
|
||
for year, group in frame.groupby(frame["trade_date"].dt.year):
|
||
if years and int(year) not in years:
|
||
continue
|
||
path = out_dir / f"{year}.parquet"
|
||
group.sort_values(["symbol", "trade_date"]).to_parquet(path, index=False)
|
||
total += len(group)
|
||
print(f"[export] {year} → {path}({len(group)} 行)")
|
||
print(f"[export] 合计 {total} 行 → {out_dir}")
|
||
return 0
|
||
|
||
|
||
_DAILY_BASIC_DEFAULT_START = date(2020, 1, 1)
|
||
|
||
|
||
def cmd_daily_basic(args) -> int:
|
||
"""每日指标同步(按交易日整表;新浪不支持 → 失败如实记录,不留静默缺口)。"""
|
||
start = _parse_day(args.start) if args.start else _DAILY_BASIC_DEFAULT_START
|
||
end = _parse_day(args.end) if args.end else date.today()
|
||
if start > end:
|
||
print("[error] --start 不能晚于 --end", file=sys.stderr)
|
||
return 2
|
||
|
||
with _session_ctx() as session:
|
||
audit_repo = SqlAlchemySyncLogRepository(session)
|
||
primary = TushareProvider(token=get_settings().tushare_token)
|
||
repo = SqlAlchemyDailyBasicRepository(session)
|
||
syncer = DailyBasicSyncer(primary=primary, repo=repo, audit=audit_repo.add)
|
||
|
||
if args.full:
|
||
cal = SqlAlchemyTradingCalendarRepository(session)
|
||
days = [d.calendar_date for d in cal.list_range(start, end) if d.is_open]
|
||
else:
|
||
days = repo.missing_dates(start, end)
|
||
total = len(days)
|
||
if total == 0:
|
||
print(f"[daily_basic] {start}~{end} 无待补交易日(本地已完整)")
|
||
return 0
|
||
print(f"[daily_basic] {start}~{end} 待同步 {total} 个交易日")
|
||
|
||
ok = failed = written = 0
|
||
failures: list[date] = []
|
||
t0 = time.time()
|
||
for idx, day in enumerate(days, start=1):
|
||
try:
|
||
res = syncer.sync_day(day)
|
||
except DataSourceError as exc:
|
||
print(f"\n[error] 第 {idx}/{total} 日 {day} 权限/凭证故障,中止:{exc}", file=sys.stderr)
|
||
session.commit()
|
||
return 1
|
||
if res.status == "ok":
|
||
ok += 1
|
||
written += res.rows_written
|
||
else:
|
||
failed += 1
|
||
failures.append(day)
|
||
_progress(idx, total, t0, f"{day} 行数={res.rows_fetched}")
|
||
if idx % 20 == 0 or idx == total:
|
||
session.commit() # 分批提交:中断时已完成的部分保持有效
|
||
if args.sleep:
|
||
time.sleep(args.sleep)
|
||
session.commit()
|
||
|
||
span = time.time() - t0
|
||
print(
|
||
f"\n[daily_basic] 完成:成功 {ok} 日 / 失败 {failed} 日,累计写入 {written} 行,"
|
||
f"耗时 {span:.1f}s"
|
||
)
|
||
if failures:
|
||
print(
|
||
"[daily_basic] 失败交易日(可重跑本命令补齐):"
|
||
+ ", ".join(d.isoformat() for d in failures[:20])
|
||
+ (" …" if len(failures) > 20 else ""),
|
||
file=sys.stderr,
|
||
)
|
||
return 1
|
||
return 0
|
||
|
||
|
||
def _progress(idx: int, total: int, t0: float, extra: str = "") -> None:
|
||
"""单行进度条(同步长任务可读性)。"""
|
||
elapsed = time.time() - t0
|
||
rate = idx / elapsed if elapsed > 0 else 0.0
|
||
eta = (total - idx) / rate if rate > 0 else 0.0
|
||
pct = idx / total * 100 if total else 100.0
|
||
end = "\n" if idx >= total else "\r"
|
||
print(
|
||
f" 进度 {idx}/{total} ({pct:5.1f}%) 已用 {elapsed:6.1f}s 预计剩余 {eta:6.1f}s {extra}",
|
||
end=end,
|
||
flush=True,
|
||
)
|
||
|
||
|
||
def build_parser() -> argparse.ArgumentParser:
|
||
parser = argparse.ArgumentParser(prog="app.cli.sync", description="Tushare 数据同步 CLI")
|
||
sub = parser.add_subparsers(dest="command", required=True)
|
||
|
||
p_nc = sub.add_parser("namechange", help="同步股票名称变更历史(时点 ST 判定)")
|
||
p_nc.add_argument("--start", help="起始日 YYYYMMDD(默认 19900101,覆盖全历史)")
|
||
p_nc.add_argument("--end", help="结束日 YYYYMMDD(默认今天)")
|
||
p_nc.set_defaults(func=cmd_namechange)
|
||
|
||
p_basic = sub.add_parser("basic", help="同步股票基础信息")
|
||
p_basic.add_argument(
|
||
"--include-delisted",
|
||
action="store_true",
|
||
help="同时拉取已退市(D)与暂停上市(P),填充 delist_date(幸存者偏差修正)",
|
||
)
|
||
p_basic.set_defaults(func=cmd_basic)
|
||
|
||
p_cal = sub.add_parser("calendar", help="同步交易日历")
|
||
p_cal.add_argument("--start", required=True, help="YYYYMMDD")
|
||
p_cal.add_argument("--end", required=True, help="YYYYMMDD")
|
||
p_cal.set_defaults(func=cmd_calendar)
|
||
|
||
p_db = sub.add_parser(
|
||
"daily_basic",
|
||
help="同步每日指标(估值/股息率/市值,按交易日整表;新浪不支持本接口)",
|
||
)
|
||
p_db.add_argument("--start", default="20200101", help="YYYYMMDD(默认 20200101)")
|
||
p_db.add_argument("--end", default="", help="YYYYMMDD(默认今天)")
|
||
p_db.add_argument(
|
||
"--full",
|
||
action="store_true",
|
||
help="忽略本地已有日期,重拉区间内全部开市日(默认只补缺失日)",
|
||
)
|
||
p_db.add_argument(
|
||
"--sleep",
|
||
type=float,
|
||
default=0,
|
||
help="每个交易日请求间隔秒数(限速时加大,如 0.2 或 1)",
|
||
)
|
||
p_db.set_defaults(func=cmd_daily_basic)
|
||
|
||
p_daily = sub.add_parser("daily", help="同步日线与复权因子(Tushare 失败 → 新浪校验兜底补缺)")
|
||
p_daily.add_argument("--symbols", default="", help="600519.SH,000001.SZ")
|
||
p_daily.add_argument("--all", action="store_true", help="遍历 stock 表全部股票")
|
||
p_daily.add_argument("--start", default="", help="YYYYMMDD(默认 20050101)")
|
||
p_daily.add_argument("--end", default="", help="YYYYMMDD(默认今天)")
|
||
p_daily.add_argument("--resume", action="store_true", help="从本地最新交易日续传(增量)")
|
||
p_daily.add_argument(
|
||
"--sleep",
|
||
type=float,
|
||
default=0,
|
||
help="每只股票请求间隔秒数(限速时加大,如 1 或 60)",
|
||
)
|
||
p_daily.set_defaults(func=cmd_daily)
|
||
|
||
p_fin = sub.add_parser("financial", help="同步财务指标快照(默认增量;Tushare 失败 → 新浪校验兜底)")
|
||
p_fin.add_argument("--symbols", default="")
|
||
p_fin.add_argument("--all", action="store_true", help="遍历 stock 表全部股票")
|
||
p_fin.add_argument(
|
||
"--full",
|
||
action="store_true",
|
||
help="强制全量重拉并覆盖既有行(默认只补本地缺失/更新的报告期,已最新跳过)",
|
||
)
|
||
p_fin.add_argument(
|
||
"--sleep",
|
||
type=float,
|
||
default=0,
|
||
help="每只股票请求间隔秒数(限速时加大,如 1 或 60)",
|
||
)
|
||
p_fin.set_defaults(func=cmd_financial)
|
||
|
||
p_idx = sub.add_parser("index_weight", help="同步指数历史成分(如沪深300 000300.SH)")
|
||
p_idx.add_argument("--code", required=True, help="指数代码,如 000300.SH / 000905.SH")
|
||
p_idx.set_defaults(func=cmd_index_weight)
|
||
|
||
p_verify = sub.add_parser("verify", help="新浪交叉验证最新行情")
|
||
p_verify.add_argument("--symbol", required=True)
|
||
p_verify.set_defaults(func=cmd_verify)
|
||
|
||
p_export = sub.add_parser("export", help="日线按年导出 Parquet(data/parquet)")
|
||
p_export.add_argument("--years", default="", help="逗号分隔年份,留空导出全部")
|
||
p_export.set_defaults(func=cmd_export)
|
||
return parser
|
||
|
||
|
||
def main(argv: list[str] | None = None) -> int:
|
||
args = build_parser().parse_args(argv)
|
||
try:
|
||
return args.func(args)
|
||
except DataSourceError as exc:
|
||
print(f"[error] {exc}", file=sys.stderr)
|
||
return 1
|
||
except KeyboardInterrupt:
|
||
print("\n[interrupt] 已中止", file=sys.stderr)
|
||
return 130
|
||
|
||
|
||
if __name__ == "__main__":
|
||
raise SystemExit(main())
|