feat: candidate_filter A+B改造——S2读stock_daily/S4读capital_flow_cache主力净流入(数据分层铁律,砍掉sina/腾讯curl) + 主循环6路并发(process_one独立conn,278只从20min+→目标3min)

This commit is contained in:
xxm
2026-08-24 20:00:47 +08:00
parent 48a20adaa3
commit d2f4ae8460
+99 -80
View File
@@ -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()