feat(data): A档全量落地——热度三件+股东户数收尾(用户批「按此方式实施,nas验证后推vps」);①xueqiu_hot雪球5633行+sina_sector新浪行业49行对账源进ak-events 19:30(同款快照语义)②hot_rank东财人气榜fetcher/注册先行,wrapper挂载留待恢复窗测量拍板(墙未定论不带不确定单元上线防缺日告警噪音)③gdhs股东户数进ak-weekly周六窗=per-period新形状:build_gdhs_units只枚举距今>90天季度末(披露守卫,防未披露期拉空走真空+done语义被marker锁死永不重拉的09-02设计坑),单次全市场5342行,marker去重无新期零网络;vintage缺日名单+2快照族(hot_rank留位);+6测试(注册/None/df契约/gdhs守卫两形态/缺日名单)215全绿 [vps]
CI/CD / test (push) Successful in 2s
CI/CD / nas-deploy (push) Successful in 9s
CI/CD / nas-verify (push) Successful in 16s

This commit is contained in:
2026-09-02 09:16:14 +08:00
parent 00cfd45d1a
commit 46a0569b9a
6 changed files with 145 additions and 6 deletions
+3 -2
View File
@@ -1,6 +1,7 @@
# ak_events_wrapper.ps1 — sanguo-ak-events schtask wrapper (daily 19:30 akshare 事件类 per-date 当日)
# dragon_tiger/block_trade/margin_sse/restricted/zt_pool 三件套/新浪资金流×2/同花顺行业概念 --start today --end today;
# dragon_tiger/block_trade/margin_sse/restricted/zt_pool 三件套/新浪资金流×2/同花顺行业概念/雪球热度/新浪行业对账 --start today --end today;
# per-date 每类1 unit 快; zt_pool 系 (2026-09-02 双机实测 push2ex 集群稳定) 供情绪/打板因子;
# hot_rank (东财人气榜) 墙待恢复窗测量, 拍板后加入 --types (fetcher/注册已就位);
# 新浪资金流+同花顺榜 (2026-09-02 用户批 A 档: 东财 push2 墙死后免费无墙替代, 秒级零限流)
$env:http_proxy = ''
$env:https_proxy = ''
@@ -11,5 +12,5 @@ $ts = Get-Date -Format 'yyyyMMdd_HHmmss'
$logDir = 'C:\sanguo_vnpy_v2\data\migration_logs'
if (-not (Test-Path $logDir)) { New-Item -ItemType Directory -Path $logDir -Force | Out-Null }
$log = Join-Path $logDir "ak_events_$ts.txt"
C:\Python310\python.exe -X utf8 C:\sanguo_vnpy_v2\scripts\data_platform\akshare_static_download.py --types dragon_tiger,block_trade,margin_sse,restricted,zt_pool,zt_pool_zbgc,zt_pool_dtgc,fund_flow_industry,fund_flow_concept,ths_industry,ths_concept --start $today --end $today *>> $log
C:\Python310\python.exe -X utf8 C:\sanguo_vnpy_v2\scripts\data_platform\akshare_static_download.py --types dragon_tiger,block_trade,margin_sse,restricted,zt_pool,zt_pool_zbgc,zt_pool_dtgc,fund_flow_industry,fund_flow_concept,ths_industry,ths_concept,xueqiu_hot,sina_sector --start $today --end $today *>> $log
exit $LASTEXITCODE
+4 -3
View File
@@ -1,5 +1,6 @@
# ak_stock_wrapper.ps1 — sanguo-ak-stock schtask wrapper (weekly 周六03:00 akshare per-stock 慢)
# northbound+share_capital+top_holders --force; per-stock 全量慢, 周末夜间
# ak_stock_wrapper.ps1 — sanguo-ak-weekly schtask wrapper (weekly 周六03:00 akshare per-stock 慢)
# northbound+share_capital+top_holders+gdhs --force; per-stock 全量慢, 周末夜间;
# gdhs=股东户数 per-period 全市场单次 (09-02 入列, 披露守卫>90天季度末, marker 去重无新期零网络)
$env:http_proxy = ''
$env:https_proxy = ''
$env:all_proxy = ''
@@ -8,5 +9,5 @@ $ts = Get-Date -Format 'yyyyMMdd_HHmmss'
$logDir = 'C:\sanguo_vnpy_v2\data\migration_logs'
if (-not (Test-Path $logDir)) { New-Item -ItemType Directory -Path $logDir -Force | Out-Null }
$log = Join-Path $logDir "ak_stock_$ts.txt"
C:\Python310\python.exe -X utf8 C:\sanguo_vnpy_v2\scripts\data_platform\akshare_static_download.py --types northbound,share_capital,top_holders --force *>> $log
C:\Python310\python.exe -X utf8 C:\sanguo_vnpy_v2\scripts\data_platform\akshare_static_download.py --types northbound,share_capital,top_holders,gdhs --force *>> $log
exit $LASTEXITCODE
@@ -143,7 +143,7 @@ PER_STOCK_TYPES = (
"financial_abstract", # stock_financial_abstract(symbol="600519")
# top_holders 单独 (per-stock × per-period)
)
# 模式 B: per-date 类型 (11 类, margin_szse 跳过)
# 模式 B: per-date 类型 (14 类, margin_szse 跳过)
PER_DATE_TYPES = (
"dragon_tiger", # stock_lhb_detail_em(start_date, end_date)
"block_trade", # stock_dzjy_mrmx(symbol="A股", start_date, end_date)
@@ -164,6 +164,12 @@ PER_DATE_TYPES = (
"fund_flow_concept", # 新浪 stock_fund_flow_concept('即时') 概念资金流 ~387 行
"ths_industry", # 同花顺 stock_board_industry_summary_ths() 行业一览 90 行×12 列
"ths_concept", # 同花顺 stock_board_concept_name_ths() 概念名单 ~375 行
# 热度/对账三件 (09-02 用户批全量落地; 快照语义同上). hot_rank=东财 emappdata
# 集群 (墙待测, fetcher/注册先行, wrapper 挂载等恢复窗测量拍板);
# xueqiu_hot/sina_sector 双机实测无墙.
"hot_rank", # 东财 stock_hot_rank_em() 人气榜 ~100 行
"xueqiu_hot", # 雪球 stock_hot_follow_xq('最热门') 关注热度 ~5633 行
"sina_sector", # 新浪 stock_sector_spot() 行业 49 行 (同花顺 90 行的对账源)
)
# 模式 C: per-period 类型 (2 类)
PER_PERIOD_TYPES = (
@@ -177,6 +183,12 @@ ONE_SHOT_TYPES = (
)
# top_holders 特殊: per-stock × per-period
TOP_HOLDERS = "top_holders"
# 股东户数 特殊: per-period (季度末) 单次全市场快照, 09-02 入列 ak-weekly。
# ⚠️ 披露守卫: 只枚举距今 >GDHS_DISCLOSE_DAYS 的季度末 —— 未披露期拉回空 df
# 会走「真空写空+done」语义永不重拉 (09-02 设计坑), 守卫保证枚举出来的期
# 数据必然已披露完, 空了才是真空。
GDHS = "gdhs"
GDHS_DISCLOSE_DAYS = 90 # 季度末后 90 天 = 披露期结束, 剩余未披露的属长尾
# 北交所代码段 (920新段 + 83/87/43历史段): akshare 东财 stock_gdfx_free_top_10_em
# 不支持北交所, build_top_holders_units 阶段直接跳过 (避免每只×20期×3retry 失败风暴).
BJ_PREFIXES = ("920", "83", "87", "43")
@@ -192,6 +204,7 @@ ALL_TYPES = (
+ PER_DATE_TYPES
+ PER_PERIOD_TYPES
+ ONE_SHOT_TYPES
+ (GDHS,)
)
@@ -682,6 +695,41 @@ def fetch_ths_concept(date: str) -> Optional[pd.DataFrame]:
return df
def fetch_hot_rank(date: str) -> Optional[pd.DataFrame]:
"""stock_hot_rank_em() — 东财人气榜快照 (~100 行)。
emappdata 集群墙待恢复窗测量定论; 失败走 None 语义 (不写不标), 无害。
"""
df, _status = call_ak_with_retry(
ak.stock_hot_rank_em, f"hot_rank/{date}",
)
return df
def fetch_xueqiu_hot(date: str) -> Optional[pd.DataFrame]:
"""stock_hot_follow_xq('最热门') — 雪球关注热度快照 (~5633 行)。"""
df, _status = call_ak_with_retry(
ak.stock_hot_follow_xq, f"xueqiu_hot/{date}", symbol="最热门",
)
return df
def fetch_sina_sector(date: str) -> Optional[pd.DataFrame]:
"""stock_sector_spot() — 新浪行业快照 (49 行, 同花顺行业榜对账源)。"""
df, _status = call_ak_with_retry(
ak.stock_sector_spot, f"sina_sector/{date}",
)
return df
def fetch_gdhs(period: str) -> Optional[pd.DataFrame]:
"""stock_zh_a_gdhs(symbol=period) — 股东户数全市场按期单次 (~5342 行)。"""
df, _status = call_ak_with_retry(
ak.stock_zh_a_gdhs, f"gdhs/{period}", symbol=period,
)
return df
# ======================== per-period fetch 函数 (2 类) ========================
def fetch_forecast(period: str) -> Optional[pd.DataFrame]:
@@ -790,6 +838,32 @@ def download_one_unit(
# ======================== 主循环 (通用, 适用所有四种模式) ========================
def build_gdhs_units(
now: Optional[datetime.date] = None,
) -> List[Tuple[str, Callable[[], pd.DataFrame]]]:
"""构造 gdhs units: 近 3 年季度末 × 披露守卫, 每期 1 unit。
只枚举距今 >GDHS_DISCLOSE_DAYS 的季度末 (披露完毕), 防止未披露期拉回
空 df 走「真空写空+done」语义后被 marker 锁死永不重拉。marker 去重:
已抓期 skip, 无新期时零网络调用 (ak-weekly 每周六空转)。
"""
now = now or datetime.date.today()
cutoff = now - datetime.timedelta(days=GDHS_DISCLOSE_DAYS)
periods = []
for y in range(now.year - 2, now.year + 1):
for (m, d) in [(3, 31), (6, 30), (9, 30), (12, 31)]:
p = datetime.date(y, m, d)
if p <= cutoff:
periods.append(f"{y}{m:02d}{d:02d}")
units = [(f"{p}_{GDHS}", partial(fetch_gdhs, p)) for p in periods]
logger.info(
"[%s] 季度末 %d 个 (披露守卫 >%d 天, %s..%s)",
GDHS, len(units), GDHS_DISCLOSE_DAYS,
periods[0] if periods else "-", periods[-1] if periods else "-",
)
return units
def run_one_type(
data_type: str,
units: List[Tuple[str, Callable[[], pd.DataFrame]]],
@@ -1129,10 +1203,17 @@ def run_type_dispatch(
"fund_flow_concept": fetch_fund_flow_concept,
"ths_industry": fetch_ths_industry,
"ths_concept": fetch_ths_concept,
"hot_rank": fetch_hot_rank,
"xueqiu_hot": fetch_xueqiu_hot,
"sina_sector": fetch_sina_sector,
}
units = build_per_date_units(t, fetch_map[t], args)
return run_one_type(t, units, args)
if t == GDHS:
units = build_gdhs_units()
return run_one_type(t, units, args)
if t in PER_PERIOD_TYPES:
fetch_map = {
"forecast": fetch_forecast,
@@ -38,6 +38,8 @@ PANEL_TYPES = (
"zt_pool", "zt_pool_zbgc", "zt_pool_dtgc", # 可回补族
"fund_flow_industry", "fund_flow_concept", # 快照族 (永久)
"ths_industry", "ths_concept",
"xueqiu_hot", "sina_sector", # 快照族 (09-02 全量落地)
# hot_rank 待恢复窗测量拍板挂载后再入列 (避免未上线就缺日告警)
)
BACKFILLABLE_PANEL_TYPES = frozenset({"zt_pool", "zt_pool_zbgc", "zt_pool_dtgc"})
PANEL_BACKFILL_CALENDAR_DAYS = 42 # ≈30 交易日回补窗的日历日近似
@@ -158,3 +158,50 @@ class TestSinaThsFamily:
for fetch in (mod.fetch_fund_flow_industry, mod.fetch_fund_flow_concept,
mod.fetch_ths_industry, mod.fetch_ths_concept):
assert fetch("20260902") is df
class TestHotAndGdhsFamily:
"""热度三件 + 股东户数 (09-02 用户批全量落地)。"""
def test_registry(self):
for t in ("hot_rank", "xueqiu_hot", "sina_sector"):
assert t in mod.PER_DATE_TYPES
assert mod.GDHS in mod.ALL_TYPES
def test_fetchers_none_on_retry_exhausted(self, monkeypatch):
monkeypatch.setattr(
mod, "call_ak_with_retry", lambda *a, **k: (None, "failed"))
assert mod.fetch_hot_rank("20260902") is None
assert mod.fetch_xueqiu_hot("20260902") is None
assert mod.fetch_sina_sector("20260902") is None
assert mod.fetch_gdhs("20250630") is None
def test_fetchers_df_on_ok(self, monkeypatch):
df = pd.DataFrame({"a": [1]})
monkeypatch.setattr(mod, "call_ak_with_retry",
lambda *a, **k: (df, "ok"))
assert mod.fetch_hot_rank("20260902") is df
assert mod.fetch_xueqiu_hot("20260902") is df
assert mod.fetch_sina_sector("20260902") is df
assert mod.fetch_gdhs("20250630") is df
def test_gdhs_units_disclosure_guard(self):
"""只枚举距今 >90 天的季度末: 90 天内的季度末绝不入列 (防空档标 done 陷阱)。"""
import datetime as _dt
now = _dt.date(2026, 9, 2)
units = mod.build_gdhs_units(now=now)
ids = [u[0] for u in units]
# cutoff = 2026-06-04: 20260630 在 90 天内必须缺席; 20250630 必须在列
assert "20250630_gdhs" in ids
assert all(not i.startswith(("20260630", "20260930")) for i in ids)
assert all(i[4:8] in ("0331", "0630", "0930", "1231") for i in ids)
# 近 3 年窗口
assert ids[0].startswith(("2024", "2025"))
def test_gdhs_units_empty_when_nothing_disclosed(self):
import datetime as _dt
# 年初 1 月: 最新的已披露季度末是去年 9/30 (1231 未满 90 天)
units = mod.build_gdhs_units(now=_dt.date(2026, 1, 15))
ids = [u[0] for u in units]
assert "20250930_gdhs" in ids
assert "20251231_gdhs" not in ids
@@ -200,6 +200,13 @@ class TestPanelGapCheck:
assert hole.isoformat() in (
status_panel["fund_flow_industry"]["holes_permanent"])
def test_panel_registry_snapshot_family(self):
"""09-02 全量落地批次注册正确: 快照族在列且不落可回补档; hot_rank 未挂载不入列。"""
for t in ("xueqiu_hot", "sina_sector", "fund_flow_concept", "ths_concept"):
assert t in svc.PANEL_TYPES
assert t not in svc.BACKFILLABLE_PANEL_TYPES
assert "hot_rank" not in svc.PANEL_TYPES # 拍板挂载后才入列
# ======================== main / JSON 契约 ========================