diff --git a/deploy/profile-scripts/candidate_filter.py b/deploy/profile-scripts/candidate_filter.py index f285f17b..be443cf4 100644 --- a/deploy/profile-scripts/candidate_filter.py +++ b/deploy/profile-scripts/candidate_filter.py @@ -377,6 +377,95 @@ def stage5_fundamental(code, name, price): return score >= 1, score, "; ".join(checks) if checks else "新标的" +# ── 单只处理(2026-08-24 并发化:每worker独立conn,逻辑=原主循环内嵌)── + +def process_one(code, name, reason, stages, stage_filter): + conn = get_conn() + try: + current_score = 0 + veto = False # S6硬否决标记(利好出货/三重打击) + print(f" {code} {name}", flush=True) + + # 获取K线(多关需要) + klines = None + + for stage_num, stage_fn, stage_name in stages: + if stage_filter and stage_num != stage_filter: + continue + + # 检查是否已通过此关 + col = f"pass_s{stage_num}" + existing = conn.execute(f"SELECT {col} FROM candidates WHERE code=?", (code,)).fetchone() + if existing and existing[0]: + continue + + if stage_num in (2, 3) and klines is None: + klines = fetch_daily_klines(code, conn) + + if stage_num == 2: + passed, sscore, detail = stage_fn(code, name, klines) + conn.execute("UPDATE candidates SET score_2nd=?, pass_s2=?, reason=? WHERE code=?", + (sscore, 1 if passed else 0, detail, code)) + log_candidate(conn, code, 2, passed, detail) + print(f" S2:{'✅' if passed else '❌'} {detail}", flush=True) + + elif stage_num == 3: + passed, sscore, detail = stage_fn(code, name, klines) + conn.execute("UPDATE candidates SET score_3rd=?, pass_s3=?, reason=? WHERE code=?", + (sscore, 1 if passed else 0, detail, code)) + log_candidate(conn, code, 3, passed, detail) + print(f" S3:{'✅' if passed else '❌'} {detail}", flush=True) + + elif stage_num == 4: + passed, sscore, detail = stage_fn(code, name) + conn.execute("UPDATE candidates SET score_4th=?, pass_s4=?, reason=? WHERE code=?", + (sscore, 1 if passed else 0, detail, code)) + log_candidate(conn, code, 4, passed, detail) + print(f" S4:{'✅' if passed else '❌'} {detail}", flush=True) + + elif stage_num == 5: + price = 0 # 从live_prices获取 + r = conn.execute("SELECT price FROM live_prices WHERE code=?", (code,)).fetchone() + if r: price = r[0] + passed, sscore, detail = stage_fn(code, name, price) + conn.execute("UPDATE candidates SET score_5th=?, pass_s5=?, reason=? WHERE code=?", + (sscore, 1 if passed else 0, detail, code)) + log_candidate(conn, code, 5, passed, detail) + print(f" S5:{'✅' if passed else '❌'} {detail}", flush=True) + + elif stage_num == 6: + passed, sscore, detail = stage_fn(code, name) + if sscore <= -99: + # 叙事硬否决:pass_s6=0, dropped=1, 综合分清零, 永不提拔 + veto = True + conn.execute("UPDATE candidates SET score_6th=0, pass_s6=0, dropped=1, drop_reason=? WHERE code=?", + (detail, code)) + log_candidate(conn, code, 6, False, detail) + print(f" S6:🚫否决 {detail}", flush=True) + else: + conn.execute("UPDATE candidates SET score_6th=?, pass_s6=?, reason=? WHERE code=?", + (sscore, 1 if passed else 0, detail, code)) + log_candidate(conn, code, 6, passed, detail) + print(f" S6:{'✅' if passed else '❌'} {detail}", flush=True) + + # 计算综合评分(S6否决则清零) + s2 = conn.execute("SELECT score_2nd FROM candidates WHERE code=?", (code,)).fetchone()[0] or 0 + s3 = conn.execute("SELECT score_3rd FROM candidates WHERE code=?", (code,)).fetchone()[0] or 0 + s4 = conn.execute("SELECT score_4th FROM candidates WHERE code=?", (code,)).fetchone()[0] or 0 + s5 = conn.execute("SELECT score_5th FROM candidates WHERE code=?", (code,)).fetchone()[0] or 0 + s6 = conn.execute("SELECT score_6th FROM candidates WHERE code=?", (code,)).fetchone()[0] or 0 + final = 0 if veto else (current_score + s2 + s3 + s4 + s5 + s6) + conn.execute("UPDATE candidates SET score_final=?, pass_final=? WHERE code=?", + (final, 0 if veto else 1, code)) + conn.commit() + return f"OK {code}" + except Exception as e: + print(f" ⚠️ {code} 处理异常: {e}", flush=True) + return f"ERR {code}: {e}" + finally: + conn.close() + + # ── 主流程 ── def main(): @@ -425,86 +514,16 @@ def main(): except Exception: pass - for code, name, reason in rows: - current_score = 0 - veto = False # S6硬否决标记(利好出货/三重打击) - print(f" {code} {name}", flush=True) - - # 获取K线(多关需要) - klines = None - - for stage_num, stage_fn, stage_name in stages: - if stage_filter and stage_num != stage_filter: - continue - - # 检查是否已通过此关 - col = f"pass_s{stage_num}" - existing = conn.execute(f"SELECT {col} FROM candidates WHERE code=?", (code,)).fetchone() - if existing and existing[0]: - continue - - if stage_num in (2, 3) and klines is None: - klines = fetch_daily_klines(code) - - if stage_num == 2: - passed, sscore, detail = stage_fn(code, name, klines) - conn.execute("UPDATE candidates SET score_2nd=?, pass_s2=?, reason=? WHERE code=?", - (sscore, 1 if passed else 0, detail, code)) - log_candidate(conn, code, 2, passed, detail) - print(f" S2:{'✅' if passed else '❌'} {detail}", flush=True) - - elif stage_num == 3: - passed, sscore, detail = stage_fn(code, name, klines) - conn.execute("UPDATE candidates SET score_3rd=?, pass_s3=?, reason=? WHERE code=?", - (sscore, 1 if passed else 0, detail, code)) - log_candidate(conn, code, 3, passed, detail) - print(f" S3:{'✅' if passed else '❌'} {detail}", flush=True) - - elif stage_num == 4: - passed, sscore, detail = stage_fn(code, name) - conn.execute("UPDATE candidates SET score_4th=?, pass_s4=?, reason=? WHERE code=?", - (sscore, 1 if passed else 0, detail, code)) - log_candidate(conn, code, 4, passed, detail) - print(f" S4:{'✅' if passed else '❌'} {detail}", flush=True) - - elif stage_num == 5: - price = 0 # 从live_prices获取 - r = conn.execute("SELECT price FROM live_prices WHERE code=?", (code,)).fetchone() - if r: price = r[0] - passed, sscore, detail = stage_fn(code, name, price) - conn.execute("UPDATE candidates SET score_5th=?, pass_s5=?, reason=? WHERE code=?", - (sscore, 1 if passed else 0, detail, code)) - log_candidate(conn, code, 5, passed, detail) - print(f" S5:{'✅' if passed else '❌'} {detail}", flush=True) - - elif stage_num == 6: - passed, sscore, detail = stage_fn(code, name) - if sscore <= -99: - # 叙事硬否决:pass_s6=0, dropped=1, 综合分清零, 永不提拔 - veto = True - conn.execute("UPDATE candidates SET score_6th=0, pass_s6=0, dropped=1, drop_reason=? WHERE code=?", - (detail, code)) - log_candidate(conn, code, 6, False, detail) - print(f" S6:🚫否决 {detail}", flush=True) - else: - conn.execute("UPDATE candidates SET score_6th=?, pass_s6=?, reason=? WHERE code=?", - (sscore, 1 if passed else 0, detail, code)) - log_candidate(conn, code, 6, passed, detail) - print(f" S6:{'✅' if passed else '❌'} {detail}", flush=True) - - # 计算综合评分(S6否决则清零) - s2 = conn.execute("SELECT score_2nd FROM candidates WHERE code=?", (code,)).fetchone()[0] or 0 - s3 = conn.execute("SELECT score_3rd FROM candidates WHERE code=?", (code,)).fetchone()[0] or 0 - s4 = conn.execute("SELECT score_4th FROM candidates WHERE code=?", (code,)).fetchone()[0] or 0 - s5 = conn.execute("SELECT score_5th FROM candidates WHERE code=?", (code,)).fetchone()[0] or 0 - s6 = conn.execute("SELECT score_6th FROM candidates WHERE code=?", (code,)).fetchone()[0] or 0 - final = 0 if veto else (current_score + s2 + s3 + s4 + s5 + s6) - conn.execute("UPDATE candidates SET score_final=?, pass_final=? WHERE code=?", - (final, 0 if veto else 1, code)) - - conn.commit() - conn.close() - print(f"[FILTER] 完成", flush=True) + # ── 2026-08-24 A+B改造:6路并发(老莫:紧急!)── + # 每只独立 conn(线程安全);S2读stock_daily/S4读capital_flow_cache后单只毫秒级, + # 瓶颈只剩 S6 新闻查询+写库,并发6轻松消化 278 只(原串行curl 20min+ → 目标<3min) + from concurrent.futures import ThreadPoolExecutor + with ThreadPoolExecutor(max_workers=6) as ex: + results = list(ex.map( + lambda rc: process_one(rc[0], rc[1], rc[2], stages, stage_filter), + rows)) + done = sum(1 for r in results if r and r.startswith("OK")) + print(f"[FILTER] 完成 {done}/{len(rows)}", flush=True) if __name__ == "__main__": main()