From a2fe51b4b063d88e833d7976af28bc776c44c6db Mon Sep 17 00:00:00 2001 From: claude_dev Date: Tue, 25 Aug 2026 21:32:28 +0800 Subject: [PATCH] =?UTF-8?q?fix(nas):=20bs=5F5m=E5=A4=B1=E8=B4=A5=E8=82=A1?= =?UTF-8?q?=E5=B0=BE=E9=83=A8=E8=A1=A5=E6=89=AB=E4=B8=80=E8=BD=AE=E2=80=94?= =?UTF-8?q?=E2=80=9403:0x=E5=90=8C=E6=AD=A5=E9=93=BE=E6=8C=81=E5=86=99?= =?UTF-8?q?=E9=94=8160s=20busy=5Ftimeout=E8=80=97=E5=B0=BDdatabase=20is=20?= =?UTF-8?q?locked(=E6=AF=8F=E6=99=9A1-5=E5=8F=AA)+=E5=81=B6=E5=8F=91?= =?UTF-8?q?=E7=BD=91=E7=BB=9C=E9=94=99(600247),=E5=8E=9F=E5=BD=A2=E6=80=81?= =?UTF-8?q?=3D=E8=AF=A5=E7=89=87=E7=AA=97=E5=8F=A3=E6=B0=B8=E4=B9=85?= =?UTF-8?q?=E5=B0=8F=E6=B4=9E;=E4=BF=AE=3D=5Fsweep=E8=BF=94failed=E6=B8=85?= =?UTF-8?q?=E5=8D=95,run=E5=B0=BE=E9=83=A8(07:40+=E5=90=8C=E6=AD=A5?= =?UTF-8?q?=E9=93=BE=E6=97=A9=E5=B7=B2=E6=94=B6=E5=B7=A5)chunk/inc?= =?UTF-8?q?=E5=90=84=E8=87=AA=E5=A4=B1=E8=B4=A5=E9=9B=86=E5=90=84=E8=A1=A5?= =?UTF-8?q?=E6=89=AB=E4=B8=80=E8=BD=AE=E8=87=AA=E6=84=88,=E8=A1=A5?= =?UTF-8?q?=E6=89=ABquery=E8=AE=A1=E5=85=A5swept=E9=85=8D=E9=A2=9D?= =?UTF-8?q?=E3=80=81rlr=E5=B9=B6=E5=85=A5exit=E8=AF=AD=E4=B9=89;+2?= =?UTF-8?q?=E6=B5=8B=E8=AF=95(=E5=8F=8C=E7=AA=97=E5=8F=A3=E5=A4=B1?= =?UTF-8?q?=E8=B4=A5=E2=86=92=E5=8F=8Cretry=3D6=E8=B0=83=E7=94=A8=E8=87=AA?= =?UTF-8?q?=E6=84=88/=E9=9B=B6=E5=A4=B1=E8=B4=A5=E6=97=A0retry=E5=9B=9E?= =?UTF-8?q?=E5=BD=92);venv310=20644=E7=BB=BF=20[nas]?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- scripts/data_platform/bs_5m_eod.py | 27 ++++++++++++++++++----- tests/data_platform/test_bs_5m_eod.py | 31 +++++++++++++++++++++++++++ 2 files changed, 53 insertions(+), 5 deletions(-) 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