Files
MoFin/deploy/profile-scripts/mo_alphasift_bridge.py
T

265 lines
9.4 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/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"))
conn.execute("PRAGMA busy_timeout=30000") # 2026-08-28 防并发写锁(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
# 写入 DB2026-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()