diff --git a/scripts/data_platform/ak_events_wrapper.ps1 b/scripts/data_platform/ak_events_wrapper.ps1 index 8d8dfcd..a8c217e 100644 --- a/scripts/data_platform/ak_events_wrapper.ps1 +++ b/scripts/data_platform/ak_events_wrapper.ps1 @@ -1,5 +1,6 @@ # ak_events_wrapper.ps1 — sanguo-ak-events schtask wrapper (daily 19:30 akshare 事件类 per-date 当日) -# dragon_tiger/block_trade/margin_sse/restricted --start today --end today; per-date 每类1 unit 快 +# dragon_tiger/block_trade/margin_sse/restricted/zt_pool 三件套 --start today --end today; +# per-date 每类1 unit 快; zt_pool 系 (2026-09-02 双机实测 push2ex 集群稳定) 供情绪/打板因子 $env:http_proxy = '' $env:https_proxy = '' $env:all_proxy = '' @@ -9,5 +10,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 --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 --start $today --end $today *>> $log exit $LASTEXITCODE diff --git a/scripts/data_platform/akshare_static_download.py b/scripts/data_platform/akshare_static_download.py index fcc49c5..45cbf36 100644 --- a/scripts/data_platform/akshare_static_download.py +++ b/scripts/data_platform/akshare_static_download.py @@ -143,12 +143,19 @@ PER_STOCK_TYPES = ( "financial_abstract", # stock_financial_abstract(symbol="600519") # top_holders 单独 (per-stock × per-period) ) -# 模式 B: per-date 类型 (4 类, margin_szse 跳过) +# 模式 B: per-date 类型 (7 类, 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) "margin_sse", # stock_margin_detail_sse(date) "restricted", # stock_restricted_release_detail_em(start_date, end_date) + # 涨停池三件套 (2026-09-02 双机实测 push2ex 集群稳定: 连打多次全 OK, 与 push2 + # clist 集群的 IP 长窗暗配额无关)。情绪周期/打板因子原料; dtgc 空数据=当日 + # 真无跌停 (写空+done 语义正确); 三端点均只支持最近 ~30 个交易日, 历史回补 + # 只能往前 30 日。 + "zt_pool", # stock_zt_pool_em(date) 涨停池 + "zt_pool_zbgc", # stock_zt_pool_zbgc_em(date) 炸板池 + "zt_pool_dtgc", # stock_zt_pool_dtgc_em(date) 跌停池 ) # 模式 C: per-period 类型 (2 类) PER_PERIOD_TYPES = ( @@ -611,6 +618,30 @@ def fetch_restricted(date: str) -> Optional[pd.DataFrame]: return df +def fetch_zt_pool(date: str) -> Optional[pd.DataFrame]: + """stock_zt_pool_em(date=date) — 涨停池 (连板数/封板资金/首次封板时间等)。""" + df, _status = call_ak_with_retry( + ak.stock_zt_pool_em, f"zt_pool/{date}", date=date, + ) + return df + + +def fetch_zt_pool_zbgc(date: str) -> Optional[pd.DataFrame]: + """stock_zt_pool_zbgc_em(date=date) — 炸板池 (涨停开板失败股)。""" + df, _status = call_ak_with_retry( + ak.stock_zt_pool_zbgc_em, f"zt_pool_zbgc/{date}", date=date, + ) + return df + + +def fetch_zt_pool_dtgc(date: str) -> Optional[pd.DataFrame]: + """stock_zt_pool_dtgc_em(date=date) — 跌停池 (空=当日真无跌停, 写空+done)。""" + df, _status = call_ak_with_retry( + ak.stock_zt_pool_dtgc_em, f"zt_pool_dtgc/{date}", date=date, + ) + return df + + # ======================== per-period fetch 函数 (2 类) ======================== def fetch_forecast(period: str) -> Optional[pd.DataFrame]: @@ -1051,6 +1082,9 @@ def run_type_dispatch( "block_trade": fetch_block_trade, "margin_sse": fetch_margin_sse, "restricted": fetch_restricted, + "zt_pool": fetch_zt_pool, + "zt_pool_zbgc": fetch_zt_pool_zbgc, + "zt_pool_dtgc": fetch_zt_pool_dtgc, } units = build_per_date_units(t, fetch_map[t], args) return run_one_type(t, units, args) diff --git a/tests/data_platform/test_akshare_static_download.py b/tests/data_platform/test_akshare_static_download.py index 1a16e73..992fd8d 100644 --- a/tests/data_platform/test_akshare_static_download.py +++ b/tests/data_platform/test_akshare_static_download.py @@ -99,3 +99,37 @@ class TestFailedNoOverwrite: lambda *a, **k: (df, "ok")) assert mod.fetch_income_sheet("SH600519") is df assert mod.fetch_express("20260630") is df + + +class TestZtPoolFamily: + """涨停池三件套 (2026-09-02 双机实测入列, push2ex 集群)。""" + + def test_registered_as_per_date_types(self): + for t in ("zt_pool", "zt_pool_zbgc", "zt_pool_dtgc"): + assert t in mod.PER_DATE_TYPES + assert t 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_zt_pool("20260901") is None + assert mod.fetch_zt_pool_zbgc("20260901") is None + assert mod.fetch_zt_pool_dtgc("20260901") is None + + def test_fetchers_df_on_ok(self, monkeypatch): + df = pd.DataFrame({"连板数": [3]}) + monkeypatch.setattr(mod, "call_ak_with_retry", + lambda *a, **k: (df, "ok")) + assert mod.fetch_zt_pool("20260901") is df + assert mod.fetch_zt_pool_zbgc("20260901") is df + assert mod.fetch_zt_pool_dtgc("20260901") is df + + def test_empty_dtgc_day_writes_empty_and_marks(self, tmp_path, monkeypatch): + """跌停池 0 跌停日 = 真空 df → 写空 + done (复跑不重拉)。""" + monkeypatch.setattr(mod, "OUT_DIR", tmp_path) + status, _ = mod.download_one_unit( + "zt_pool_dtgc", "20260901_zt_pool_dtgc", + lambda: pd.DataFrame(), force=True) + assert status == "empty" + assert (tmp_path / "zt_pool_dtgc" / "20260901_zt_pool_dtgc.parquet").exists() + assert (tmp_path / "zt_pool_dtgc" / ".20260901_zt_pool_dtgc.akshare").exists()