- candidate_filter S6: 消息(个股+行业)x资金x技术叙事矩阵 利好出货/三重打击=硬否决(dropped=1,不可被高分抵消) 共振做多+3/利空出尽+1/资金驱动+1/阴跌-1/行业利空-1 - promote: 信号一律先入'关注'(12维确认才升,消灭出生即买入); 容量60硬顶,RR<2的末位淘汰 - alphasift收编: 只写candidates不再直写自选(统一入口)
264 lines
9.3 KiB
Python
264 lines
9.3 KiB
Python
#!/usr/bin/env python3
|
||
"""
|
||
mo_alphasift_bridge.py — AlphaSift 多策略并行选股 → MoFin 自选池
|
||
|
||
支持同时跑多个策略,合并去重后写入自选池。
|
||
默认三策略: balanced_alpha + dual_low + quality_value
|
||
|
||
用法:
|
||
python3 mo_alphasift_bridge.py # 默认三策略
|
||
python3 mo_alphasift_bridge.py --strategy dual_low # 单策略
|
||
python3 mo_alphasift_bridge.py --dry-run # 只看不写
|
||
python3 mo_alphasift_bridge.py --strategy list # 列出所有策略
|
||
"""
|
||
|
||
import sys, os, json, argparse, urllib.request, time
|
||
from datetime import datetime
|
||
from pathlib import Path
|
||
|
||
DSA_API = "http://127.0.0.1:8001"
|
||
MOFIN_DATA = Path("/home/hmo/web-dashboard/data")
|
||
# 已迁移到 DB — watchlist.json / portfolio.json 不再使用
|
||
# 保留路径仅用于兼容 import,实际数据从 mofin.db 读取
|
||
|
||
DEFAULT_STRATEGIES = "balanced_alpha,dual_low,quality_value"
|
||
DEFAULT_MARKET = "cn"
|
||
DEFAULT_MAX = 15
|
||
MIN_SCORE = 5
|
||
MAX_ADD = 5
|
||
POLL_INTERVAL = 5
|
||
POLL_TIMEOUT = 300
|
||
|
||
# 全局开关: ALPHASIFT_ENABLED=true 才执行选股
|
||
ALPHASIFT_ENABLED = os.environ.get("ALPHASIFT_ENABLED", "false").lower() == "true"
|
||
|
||
|
||
def load_json(path):
|
||
try: return json.loads(path.read_text(encoding="utf-8"))
|
||
except: return None
|
||
|
||
def save_json(path, data):
|
||
path.parent.mkdir(parents=True, exist_ok=True)
|
||
path.write_text(json.dumps(data, ensure_ascii=False, indent=2), encoding="utf-8")
|
||
|
||
def api(endpoint, method="GET", body=None):
|
||
url = f"{DSA_API}{endpoint}"
|
||
data = json.dumps(body).encode() if body else None
|
||
req = urllib.request.Request(url, data=data, method=method,
|
||
headers={"Content-Type": "application/json"})
|
||
try:
|
||
with urllib.request.urlopen(req, timeout=30) as r:
|
||
return json.loads(r.read())
|
||
except Exception as e:
|
||
print(f" API错误: {e}")
|
||
return None
|
||
|
||
def get_existing_codes():
|
||
"""从 DB 读取已有持仓+自选。JSON 已废弃。"""
|
||
codes = set()
|
||
try:
|
||
import sqlite3
|
||
db = sqlite3.connect(str(MOFIN_DATA / "mofin.db"))
|
||
for row in db.execute("SELECT code FROM watchlist_stocks WHERE is_active=1"):
|
||
codes.add(str(row[0]).strip())
|
||
for row in db.execute("SELECT code FROM holdings"):
|
||
codes.add(str(row[0]).strip())
|
||
db.close()
|
||
except Exception as e:
|
||
print(f"WARN: DB读取失败: {e}")
|
||
return codes
|
||
|
||
|
||
def run_one_strategy(strategy, market, max_results):
|
||
"""跑单个策略,返回候选股列表"""
|
||
print(f"\n{'='*50}")
|
||
print(f"策略: {strategy}")
|
||
print(f"{'='*50}")
|
||
|
||
task = api("/api/v1/alphasift/screen/tasks", "POST", {
|
||
"strategy": strategy, "market": market, "max_results": max_results
|
||
})
|
||
if not task or not task.get("task_id"):
|
||
print(f" FAIL: 提交失败")
|
||
return []
|
||
|
||
task_id = task["task_id"]
|
||
print(f" 任务: {task_id[:12]}...", flush=True)
|
||
|
||
waited = 0
|
||
while waited < POLL_TIMEOUT:
|
||
time.sleep(POLL_INTERVAL)
|
||
waited += POLL_INTERVAL
|
||
status = api(f"/api/v1/alphasift/screen/tasks/{task_id}")
|
||
if not status: continue
|
||
s = status.get("status", "")
|
||
if s == "completed":
|
||
print(f" 完成 ({waited}s)")
|
||
result = status.get("result", {})
|
||
candidates = result.get("candidates", [])
|
||
print(f" 候选: {len(candidates)} 只")
|
||
if candidates:
|
||
for c in candidates[:3]:
|
||
print(f" {c.get('code','?')} {c.get('name',c.get('title','?'))} 评分{c.get('score','?'):.1f}")
|
||
return candidates
|
||
elif s == "failed":
|
||
print(f" FAIL: {status.get('error','')}")
|
||
return []
|
||
else:
|
||
if waited % 60 == 0:
|
||
print(f" ...{s} ({status.get('progress',0)}%)", flush=True)
|
||
|
||
print(f" FAIL: 超时")
|
||
return []
|
||
|
||
|
||
def run_all(strategies_str, market, max_results, dry_run=False):
|
||
"""多策略并行 → 合并去重 → MoFin 自选池"""
|
||
if not ALPHASIFT_ENABLED and not dry_run:
|
||
print("AlphaSift 已禁用。设置 ALPHASIFT_ENABLED=true 或 --enable 启用。")
|
||
return
|
||
strategies = [s.strip() for s in strategies_str.split(",") if s.strip()]
|
||
now = datetime.now()
|
||
date_str = now.strftime("%Y-%m-%d")
|
||
time_str = now.strftime("%Y-%m-%d %H:%M")
|
||
|
||
print(f"AlphaSift 多策略选股: {', '.join(strategies)}")
|
||
print(f"开始: {time_str}")
|
||
|
||
# 逐个跑策略,汇总
|
||
all_candidates = []
|
||
seen = set()
|
||
for strategy in strategies:
|
||
candidates = run_one_strategy(strategy, market, max_results)
|
||
for c in candidates:
|
||
code = str(c.get("code", "")).strip()
|
||
if code in seen: continue
|
||
seen.add(code)
|
||
c["_strategy"] = strategy
|
||
all_candidates.append(c)
|
||
|
||
if not all_candidates:
|
||
print("\n无候选股")
|
||
return
|
||
|
||
print(f"\n汇总: {len(all_candidates)} 只候选股 (去重后)")
|
||
for s in strategies:
|
||
cnt = sum(1 for c in all_candidates if c.get("_strategy") == s)
|
||
print(f" {s}: {cnt} 只")
|
||
|
||
# 过滤
|
||
existing = get_existing_codes()
|
||
new_stocks = []
|
||
skipped_score = 0
|
||
skipped_dup = 0
|
||
|
||
for c in all_candidates:
|
||
code = str(c.get("code", "")).strip()
|
||
score = c.get("score", 0) or c.get("llm_score", 0) or 0
|
||
if score < MIN_SCORE:
|
||
skipped_score += 1; continue
|
||
if code in existing:
|
||
skipped_dup += 1; continue
|
||
|
||
name = c.get("name", "") or c.get("title", "") or code
|
||
reason = c.get("reason", "") or c.get("llm_thesis", "")
|
||
src = c.get("_strategy", "unknown")
|
||
|
||
factors = c.get("factor_scores", {})
|
||
factor_note = ", ".join(f"{k}={v:.0f}" for k,v in list(factors.items())[:3]) if factors else ""
|
||
|
||
notes = f"AlphaSift/{src} 评分{score:.0f}"
|
||
if factor_note: notes += f" [{factor_note}]"
|
||
if reason: notes += f" | {reason[:120]}"
|
||
|
||
new_stocks.append({
|
||
"code": code,
|
||
"name": name,
|
||
"price": c.get("price", 0),
|
||
"source": "alpha_sift",
|
||
"source_detail": {
|
||
"strategy": src,
|
||
"strategies_run": strategies,
|
||
"score": score,
|
||
"factor_scores": factors,
|
||
"date": date_str,
|
||
"reason": reason[:300],
|
||
},
|
||
"notes": notes,
|
||
"added_at": time_str,
|
||
"added_by": "AlphaSift",
|
||
"analysis": {},
|
||
})
|
||
existing.add(code)
|
||
|
||
print(f"\n过滤: {len(new_stocks)} 新标的 (评分不足{skipped_score} + 重复{skipped_dup})")
|
||
|
||
if not new_stocks:
|
||
print("无符合条件的新标的")
|
||
return
|
||
|
||
new_stocks.sort(key=lambda s: s["source_detail"]["score"], reverse=True)
|
||
new_stocks = new_stocks[:MAX_ADD]
|
||
|
||
print(f"\n新增 {len(new_stocks)} 只到自选池:")
|
||
for s in new_stocks:
|
||
sd = s["source_detail"]
|
||
print(f" {s['code']} {s['name']} ({sd['strategy']} 评分{sd['score']:.0f})")
|
||
|
||
if dry_run:
|
||
print("\n[DRY RUN] 未写入")
|
||
return
|
||
|
||
# 写入 DB(2026-07-24 老爸"统一入口":只写 candidates,不许直写自选。
|
||
# 后续由 promote_candidates 的 score>=7+RR>=2+S6叙事闸 统一提拔)
|
||
try:
|
||
sys.path.insert(0, str(MOFIN_DATA.parent))
|
||
from mofin_db import get_conn
|
||
conn = get_conn()
|
||
for s in new_stocks:
|
||
sd = s.get("source_detail", {})
|
||
conn.execute(
|
||
"INSERT INTO candidates (code, name, sector, reason, entry_range, stop_loss, target, created_at) "
|
||
"VALUES (?,?,?,?,?,?,?,datetime('now','localtime')) "
|
||
"ON CONFLICT(code) DO UPDATE SET "
|
||
"name=excluded.name, sector=excluded.sector, reason=excluded.reason, "
|
||
"entry_range=excluded.entry_range, stop_loss=excluded.stop_loss, target=excluded.target",
|
||
(s["code"], s["name"], "alpha_sift",
|
||
f"{s.get('notes','')}",
|
||
"", 0, 0))
|
||
conn.commit()
|
||
conn.close()
|
||
print(f"\n已写入 {len(new_stocks)} 只到 candidates(等待 candidate_filter 多级过滤+promote 统一提拔)")
|
||
except Exception as e:
|
||
print(f"WARN: DB写入失败: {e}")
|
||
return
|
||
|
||
|
||
def list_strategies():
|
||
r = api("/api/v1/alphasift/strategies")
|
||
if r and r.get("strategies"):
|
||
for s in r["strategies"]:
|
||
print(f" {s['id']:22s} {s.get('name','?'):10s} {s.get('description','')[:60]}")
|
||
|
||
|
||
def main():
|
||
global MIN_SCORE
|
||
p = argparse.ArgumentParser(description="AlphaSift → MoFin")
|
||
p.add_argument("--strategy", default=DEFAULT_STRATEGIES)
|
||
p.add_argument("--market", default=DEFAULT_MARKET)
|
||
p.add_argument("--max", type=int, default=DEFAULT_MAX)
|
||
p.add_argument("--min-score", type=int, default=MIN_SCORE)
|
||
p.add_argument("--enable", action="store_true", help="覆盖 ALPHASIFT_ENABLED 开关")
|
||
p.add_argument("--dry-run", action="store_true")
|
||
args = p.parse_args()
|
||
MIN_SCORE = args.min_score
|
||
if args.enable:
|
||
global ALPHASIFT_ENABLED
|
||
ALPHASIFT_ENABLED = True
|
||
if args.strategy == "list": list_strategies()
|
||
else: run_all(args.strategy, args.market, args.max, args.dry_run)
|
||
|
||
|
||
if __name__ == "__main__":
|
||
main()
|