feat(data): akshare事件类扩涨停池三件套——daily_stock_analysis调研落地第一步;双机实测(Mac+VPS生产IP)分流:push2ex集群(zt_pool/zbgc/dtgc)连打多轮全稳可用,push2 clist集群(行业/概念榜/板块资金流/主力资金)与push2his(cyq筹码/个股资金流)每IP长窗暗配额(首发200第二发起全reset,30s/镜像轮换/换头均不过)判敌对弃采;zt_pool 83只/zbgc 6只/dtgc 2只真实行数实证,dtgc空数据=当日真无跌停走既有写空+done语义;三类型进PER_DATE_TYPES家族(ak-events 19:30同窗+3单次调用<1min),零新schtask零新目录零新依赖;+4测试(注册/None契约/df契约/真空df落盘) [vps]
CI/CD / test (push) Successful in 4s
CI/CD / nas-deploy (push) Successful in 13s
CI/CD / nas-verify (push) Successful in 14s

This commit is contained in:
2026-09-02 08:07:31 +08:00
parent de9efec87e
commit db789d6e9c
3 changed files with 72 additions and 3 deletions
+3 -2
View File
@@ -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
@@ -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)
@@ -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()