diff --git a/scripts/data_platform/bs_5m_eod.py b/scripts/data_platform/bs_5m_eod.py index b7ce8e7..151bf3c 100644 --- a/scripts/data_platform/bs_5m_eod.py +++ b/scripts/data_platform/bs_5m_eod.py @@ -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: diff --git a/tests/data_platform/test_bs_5m_eod.py b/tests/data_platform/test_bs_5m_eod.py index 9190f49..f4884e5 100644 --- a/tests/data_platform/test_bs_5m_eod.py +++ b/tests/data_platform/test_bs_5m_eod.py @@ -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