fix(nas): bs_5m失败股尾部补扫一轮——03:0x同步链持写锁60s busy_timeout耗尽database is locked(每晚1-5只)+偶发网络错(600247),原形态=该片窗口永久小洞;修=_sweep返failed清单,run尾部(07:40+同步链早已收工)chunk/inc各自失败集各补扫一轮自愈,补扫query计入swept配额、rlr并入exit语义;+2测试(双窗口失败→双retry=6调用自愈/零失败无retry回归);venv310 644绿 [nas]
This commit is contained in:
@@ -211,10 +211,11 @@ def _sweep(conn, stocks, start, end, swept, label):
|
||||
"""扫全 A 一个窗口: per-stock 短事务 + 周期 relogin + 间隔。
|
||||
|
||||
swept 为跨窗口累计股数(=query 数, 每股恰 1 次调用), 由调用方持有并回填;
|
||||
返 (stats, limit_reached)。
|
||||
返 (stats, limit_reached, failed); failed=[(code,prefix)...] 供 run 尾部补扫。
|
||||
"""
|
||||
stats = {"ok": 0, "empty": 0, "failed": 0, "db_rows": 0}
|
||||
limit_reached = False
|
||||
failed = []
|
||||
t0 = time.time()
|
||||
for i, (code, prefix) in enumerate(stocks):
|
||||
if swept[0] >= DAILY_LIMIT:
|
||||
@@ -232,6 +233,7 @@ def _sweep(conn, stocks, start, end, swept, label):
|
||||
stats["empty"] += 1
|
||||
except Exception as e:
|
||||
stats["failed"] += 1
|
||||
failed.append((code, prefix))
|
||||
if stats["failed"] <= 5 or stats["failed"] % 100 == 0:
|
||||
log.warning("[%s] %s err: %s", label, code, e)
|
||||
if not relogin():
|
||||
@@ -244,7 +246,7 @@ def _sweep(conn, stocks, start, end, swept, label):
|
||||
log.warning("[%s] 周期 relogin 失败, 继续跑", label)
|
||||
if i < len(stocks) - 1:
|
||||
time.sleep(BS_INTERVAL)
|
||||
return stats, limit_reached
|
||||
return stats, limit_reached, failed
|
||||
|
||||
|
||||
def _run(args):
|
||||
@@ -276,9 +278,10 @@ def _run(args):
|
||||
limit_reached = False
|
||||
t0 = time.time()
|
||||
try:
|
||||
failed_chunk = []
|
||||
if chunk:
|
||||
cs, lr = _sweep(conn, stocks, chunk[0], chunk[1], swept,
|
||||
"chunk %s~%s" % chunk)
|
||||
cs, lr, failed_chunk = _sweep(conn, stocks, chunk[0], chunk[1], swept,
|
||||
"chunk %s~%s" % chunk)
|
||||
limit_reached = limit_reached or lr
|
||||
log.info("[CHUNK] %s~%s ok=%d empty=%d failed=%d rows=%d",
|
||||
chunk[0], chunk[1], cs["ok"], cs["empty"], cs["failed"],
|
||||
@@ -289,11 +292,25 @@ def _run(args):
|
||||
# 实锤多推一片跳过半年窗口)
|
||||
save_chunk_state({"next_end": chunk[0]})
|
||||
# 每日增量照跑(回灌期也跑, 近端 7 天始终新鲜)
|
||||
istats, lr = _sweep(conn, stocks, inc_start, inc_end, swept, "inc")
|
||||
istats, lr, failed_inc = _sweep(conn, stocks, inc_start, inc_end, swept, "inc")
|
||||
limit_reached = limit_reached or lr
|
||||
log.info("[INC] ok=%d empty=%d failed=%d rows=%d",
|
||||
istats["ok"], istats["empty"], istats["failed"],
|
||||
istats["db_rows"])
|
||||
# 失败股尾部补扫一轮(2026-08-25): 03:0x 同步链持写锁→60s busy_timeout
|
||||
# 耗尽 database is locked(每晚 1-5 只)+偶发网络错, 原形态=该片窗口永久
|
||||
# 小洞; run 尾部(07:40+)同步链早已收工, 补扫一轮自愈两类失败
|
||||
retries = []
|
||||
if chunk and failed_chunk:
|
||||
retries.append(("chunk-retry", failed_chunk, chunk[0], chunk[1]))
|
||||
if failed_inc:
|
||||
retries.append(("inc-retry", failed_inc, inc_start, inc_end))
|
||||
for rlabel, rcodes, rstart, rend in retries:
|
||||
log.info("[%s] 补扫 %d 只失败股", rlabel, len(rcodes))
|
||||
rs, rlr, _rf = _sweep(conn, rcodes, rstart, rend, swept, rlabel)
|
||||
limit_reached = limit_reached or rlr
|
||||
log.info("[%s] ok=%d empty=%d failed=%d rows=%d",
|
||||
rlabel, rs["ok"], rs["empty"], rs["failed"], rs["db_rows"])
|
||||
if next_chunk(today) is None and not limit_reached:
|
||||
save_chunk_state({"next_end": None}) # 回灌完结盖章(幂等)
|
||||
finally:
|
||||
|
||||
@@ -267,3 +267,34 @@ def test_main_second_instance_skips(main_env, monkeypatch):
|
||||
with pytest.raises(SystemExit) as e:
|
||||
b5.main()
|
||||
assert e.value.code == 0
|
||||
|
||||
|
||||
def test_main_failed_stocks_retry_pass(main_env, monkeypatch):
|
||||
"""失败股尾部补扫: 首扫抛错的股在 retry 轮重拉, 锁错/网络错自愈不留洞。"""
|
||||
monkeypatch.setattr(sys, "argv", ["bs_5m_eod.py"])
|
||||
seen = set()
|
||||
|
||||
def fake_process(conn, code, prefix, start, end, timeout):
|
||||
key = (code, start)
|
||||
first = key not in seen
|
||||
seen.add(key)
|
||||
if code == "600001" and first:
|
||||
raise RuntimeError("database is locked")
|
||||
return 48
|
||||
|
||||
with patch.object(b5, "_process_one", side_effect=fake_process) as m:
|
||||
with pytest.raises(SystemExit) as e:
|
||||
b5.main()
|
||||
assert e.value.code == 0
|
||||
# A 在 chunk/inc 两窗口各自首调抛错 → chunk-retry + inc-retry 各补 1 次:
|
||||
# chunk(A,B) + inc(A,B) + chunk-retry(A) + inc-retry(A) = 6
|
||||
assert m.call_count == 6
|
||||
|
||||
|
||||
def test_main_retry_not_triggered_when_no_failures(main_env):
|
||||
"""零失败 → 无补扫轮(回归: 原 4 次调用不变)。"""
|
||||
with patch.object(b5, "_process_one", return_value=48) as m:
|
||||
with pytest.raises(SystemExit) as e:
|
||||
b5.main()
|
||||
assert e.value.code == 0
|
||||
assert m.call_count == 4 # 2 股 × (chunk+inc), 无 retry
|
||||
|
||||
Reference in New Issue
Block a user