Files
MoFin/deploy/profile-scripts/mo_alphasift_bridge.py
T
hmo 309f074c51 feat: S6叙事一致性检查+统一入口+容量60硬顶+提拔先入观察
- candidate_filter S6: 消息(个股+行业)x资金x技术叙事矩阵
  利好出货/三重打击=硬否决(dropped=1,不可被高分抵消)
  共振做多+3/利空出尽+1/资金驱动+1/阴跌-1/行业利空-1
- promote: 信号一律先入'关注'(12维确认才升,消灭出生即买入); 容量60硬顶,RR<2的末位淘汰
- alphasift收编: 只写candidates不再直写自选(统一入口)
2026-07-24 12:09:30 +08:00

264 lines
9.3 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"))
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()