From 890ff364854841f0d5c770b926a47dd0b984b7ba Mon Sep 17 00:00:00 2001 From: lookt Date: Tue, 15 Sep 2026 21:36:26 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20K=E7=BA=BFworker=E9=99=90=E9=80=9F+?= =?UTF-8?q?=E6=96=B0=E6=B5=AA=E6=97=A5=E7=BA=BF=E5=A4=87=E9=80=89+?= =?UTF-8?q?=E5=A4=B1=E8=B4=A5=E9=80=80=E9=81=BF=EF=BC=88=E5=BA=94=E5=AF=B9?= =?UTF-8?q?=E4=B8=9C=E8=B4=A2=E9=99=90=E9=80=9F=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/fetcher/kline_fetcher.py | 13 +++++++++++-- src/quant/engine.py | 23 ++++++++++++++++++++++- 2 files changed, 33 insertions(+), 3 deletions(-) diff --git a/src/fetcher/kline_fetcher.py b/src/fetcher/kline_fetcher.py index eb3ceb2..c5d18c8 100644 --- a/src/fetcher/kline_fetcher.py +++ b/src/fetcher/kline_fetcher.py @@ -48,8 +48,17 @@ class KlineFetcher: end = datetime.now().strftime('%Y%m%d') start = (datetime.now() - timedelta(days=days + 5)).strftime('%Y%m%d') if cfg['ak_fn'] == 'hist': - df = ak.stock_zh_a_hist(symbol=code, period='daily', - start_date=start, end_date=end, adjust='qfq') + try: + df = ak.stock_zh_a_hist(symbol=code, period='daily', + start_date=start, end_date=end, adjust='qfq') + except Exception as e: + # EM 被限速/拦截时回退新浪日线(前复权) + sym = ('sh' if code[:1] == '6' else 'sz') + code + try: + df = ak.stock_zh_a_daily(symbol=sym, start_date=start, end_date=end, + adjust='qfq') + except Exception as e2: + raise e2 from e else: df = ak.stock_zh_a_hist_min_em(symbol=code, period=cfg['ak_period'], start_date=start + ' 09:30:00', diff --git a/src/quant/engine.py b/src/quant/engine.py index 9ce6e66..4873af9 100644 --- a/src/quant/engine.py +++ b/src/quant/engine.py @@ -75,6 +75,17 @@ class QuantEngine: finally: pass + def _kline_count(self, tf): + import sqlite3 + try: + conn = sqlite3.connect(self.db_path, check_same_thread=False) + n = conn.execute("SELECT COUNT(*) FROM kline WHERE tf=? AND bar_time >= datetime('now','-2 days')", + (tf,)).fetchone()[0] + conn.close() + return n + except Exception: + return 0 + def _kline_worker(self): """后台轮询:持续刷新宇宙内 K 线(day 全量 + 30/5 分钟)""" while not self._stop.is_set(): @@ -85,12 +96,22 @@ class QuantEngine: continue with self.state_lock: self.names = names + fails = 0 for code in codes: if self._stop.is_set(): return for tf in ('day', '30'): + before = self._kline_count(tf) self.fetcher.sync(code, tf) - time.sleep(1.0 / KLINE_QPS) + if self._kline_count(tf) == before: + fails += 1 + else: + fails = max(0, fails - 1) + # 限速:连续失败加退避,防触发源封禁 + if fails and fails % 5 == 0: + time.sleep(min(30, 5 + fails)) + time.sleep(1.2) + print('[quant] 本轮K线刷新: {} 只, 连续未新增 {}'.format(len(codes), fails), flush=True) with self.state_lock: self.kline_ready = True print('[quant] K线刷新完成一轮: {} 只'.format(len(codes)), flush=True)