This commit is contained in:
2026-09-16 09:09:10 +08:00
parent 71a0f6e404
commit a490fdc110
8 changed files with 487 additions and 196 deletions

View File

@@ -1,8 +1,9 @@
"""全市场数据同步(未复权,写入 candles 全量底座)。
设计trade_cal 取近 N 个交易日 -> 逐日 pro.daily(trade_date=...) 一次返回全市场当日数据
-> upsert 进 candles不复权底座ON CONFLICT 幂等daily_basic 同步最新交易日到
DailySnapshot市值/PE/PB/换手率等截面字段)
-> upsert 进 candles不复权底座ON CONFLICT 幂等daily_basic 同步最新交易日到
DailySnapshot市值/PE/PB 等截面字段),并复用该次调用把换手率回写 candles.turnover
(历史缺漏日由自愈循环补,见 _run_sync 第 3.5 步)。
同步为进程内后台任务MVP 不引入任务队列),前端轮询 /api/screener/sync/status。
daily 与 daily_basic 分步独立落库daily_basic 积分不足时快照仍可用,错误写入状态不中断任务。
@@ -14,7 +15,7 @@ import logging
import time
from datetime import datetime, timedelta
from sqlalchemy import delete, func, insert, select
from sqlalchemy import delete, func, insert, select, text
from sqlalchemy.dialects.postgresql import insert as pg_insert
from sqlalchemy.ext.asyncio import AsyncSession
@@ -244,7 +245,7 @@ async def _upsert_candle_day(session: AsyncSession, rows: list[dict], listed: se
"open": r["open"], "high": r["high"], "low": r["low"], "close": r["close"],
"volume": r["vol"] * 100.0, # 手 -> 股
"amount": (r["amount"] * 1000.0) if r["amount"] is not None else None, # 千元 -> 元
"turnover": None, # 换手率由 daily_basic 快照维护
"turnover": None, # 换手率由 _run_sync 第 3/3.5 步从 daily_basic 回写
}
for r in rows
if plain_code(r["ts_code"]) in listed
@@ -269,10 +270,65 @@ async def _upsert_candle_day(session: AsyncSession, rows: list[dict], listed: se
await session.commit()
async def _run_sync(days: int, force: bool) -> None:
"""后台任务主体stock_basic -> 逐日日线 -> 最新交易日快照。异常写状态。
_TURNOVER_FLOOR = "20000104" # daily_basic 最早覆盖日,更早的交易日拉了也是空
daily_basic 只拉最新交易日(快照条件仅作用于最新截面,且低积分 token 限频 1 次/分钟)。
async def _backfill_turnover_day(session: AsyncSession, basic_rows: list[dict], d_str: str) -> int:
"""把 daily_basic 的 turnover_rate 回写 candles.turnover只动该列幂等
basic_rows 复用 _fetch_basic 的返回(零额外 API 调用);单条 UPDATE...FROM
unnest 批量写回ETF 等不在 daily_basic 的行不会命中。
"""
syms = [plain_code(r["ts_code"]) for r in basic_rows if r["turnover_rate"] is not None]
trs = [r["turnover_rate"] for r in basic_rows if r["turnover_rate"] is not None]
if not syms:
return 0
res = await session.execute(
text("UPDATE candles AS c SET turnover = v.t "
"FROM unnest(CAST(:syms AS text[]), CAST(:trs AS float8[])) AS v(sym, t) "
"WHERE c.symbol = v.sym AND c.timeframe = '1d' AND c.ts = :ts"),
{"syms": syms, "trs": trs, "ts": _parse_d(d_str)},
)
await session.commit()
return res.rowcount or 0
async def _backfill_turnover_gaps(pro) -> int:
"""换手率全范围自愈:按日聚合在市股票的换手覆盖,过半缺失的交易日逐日拉
daily_basic 补齐,返回处理的缺口天数。
夜间同步与手动同步共用本函数(唯一入口,幂等可断点续跑——补完的日子下轮
不再命中),正常无缺口时零 API 调用。只统计在市股票(与 _recent_day_counts
同口径ETF/DEMO 行 daily_basic 天然不覆盖,混进来会把健康日误判成缺换手。
"""
from ..db import async_session # 延迟导入避免循环
async with async_session() as session:
rows = (await session.execute(
select(func.date(Candle.ts), func.count(), func.count(Candle.turnover))
.where(Candle.timeframe == "1d", Candle.ts >= _parse_d(_TURNOVER_FLOOR),
Candle.symbol.in_(select(StockBasic.symbol).where(StockBasic.list_status == "L")))
.group_by(func.date(Candle.ts))
.order_by(func.date(Candle.ts))
)).all()
gaps = [d.strftime("%Y%m%d") for d, total, done in rows if total and done < total // 2]
for i, d_str in enumerate(gaps, 1):
_sync_state["step"] = f"正在回补 {d_str} 换手率({i}/{len(gaps)}"
try:
basic_rows = await asyncio.to_thread(_fetch_basic, pro, d_str)
if basic_rows:
async with async_session() as session:
await _backfill_turnover_day(session, basic_rows, d_str)
except Exception: # noqa: BLE001 —— 单日失败不中断,下次同步再试
log.warning("换手率回补 %s 失败(下次同步再试)", d_str, exc_info=True)
return len(gaps)
async def _run_sync(days: int, force: bool) -> None:
"""后台任务主体stock_basic -> 逐日日线 -> 最新交易日快照 + 换手率回写/自愈。异常写状态。
daily_basic 只拉最新交易日(快照条件仅作用于最新截面,且低积分 token 限频 1 次/分钟);
历史缺口的换手率由 _backfill_turnover_gaps 统一补齐,夜间/手动同步共用同一管道。
"""
from ..db import async_session # 延迟导入避免循环
@@ -343,6 +399,7 @@ async def _run_sync(days: int, force: bool) -> None:
async with async_session() as session:
latest_dt = await session.scalar(select(func.max(Candle.ts)))
latest = latest_dt.strftime("%Y%m%d") if latest_dt else None
basic_rows: list[dict] = []
if latest:
async with async_session() as session:
have_snap = force or latest not in await _existing_dates(session, DailySnapshot)
@@ -353,6 +410,23 @@ async def _run_sync(days: int, force: bool) -> None:
async with async_session() as session:
await _replace_day(session, DailySnapshot, basic_rows, latest)
# 3.5) 换手率回写:日线同步不写 turnoverdaily_basic 才有)——快照那次调用
# 顺手回写最新日(零额外 API 调用);历史缺口统一由 _backfill_turnover_gaps
# 全范围扫补,夜间/手动同步共用同一管道
if latest and basic_rows:
_sync_state["step"] = f"正在回写 {latest} 换手率"
try:
async with async_session() as session:
await _backfill_turnover_day(session, basic_rows, latest)
except Exception: # noqa: BLE001 —— 回写失败不影响快照,缺口由自愈兜底
log.warning("换手率回写 %s 失败(下次同步自愈)", latest, exc_info=True)
try:
n_gap = await _backfill_turnover_gaps(pro)
if n_gap:
log.info("换手率自愈补齐 %d 个交易日", n_gap)
except Exception: # noqa: BLE001 —— 自愈失败不阻断同步收尾,下次再试
log.warning("换手率自愈失败(下次同步再试)", exc_info=True)
# candles/复权因子已更新:作废旧 K 线预览缓存(键含版本号,自增即全体失效)
await cache.bump_version("candles")
# 预热统计缓存:同步任务自己付一次重聚合(>10s。SWR 下轮询方不等待——