tradingSystem/app/services/tech_service.py

436 lines
20 KiB
Python
Raw Permalink 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.

# -*- coding: utf-8 -*-
"""技术面取数、落表、映射与对外状态 (2026-09-11 技术面接入)。
职责三件:
一, 分页拉决策系统 GET /api/v1/market/technical 的全市场读数, 宽松解析后落进
pms_tech_daily。**状态非 OK 不落表**; 分页里任一页失败不推进 data_date, 整轮算失败,
绝不落半份 (先记账后动作, 宁可这轮没有也不要落一份残缺的)。
二, 每早对相关票 (在持 + 计划主榜观察 + 待拍板) 用 core.tech_rules 合成技术面立场, 写进
运行参数 PMS_TECH_STATE_MAP。与 logic_state_service 的映射同一手法: 读不到留空带原因、
超 FRESH_DAYS 天按无读数、绝不折成看空 (设计原则二)。
三, 给页面出状态 (开关/新鲜度/行数/映射覆盖/错误) 与按票的研究面读数 (点击时调)。
取数纯函数 (fetch_page / fetch_all_pages / parse_item) 不读参数中心, 以便脱库单测。
总开关 PMS_TECH_ENABLED 关掉即整模块歇工 (不拉、不建映射、状态标注已停, 行为逐字如旧)。
"""
from __future__ import annotations
import json
import logging
from collections import Counter
from datetime import date, datetime
from app.core import tech_rules, tradedays
from app.core.command_spec import normalize_code
from app.repo import tech_repo
from app.services import param_store
logger = logging.getLogger("pms.tech")
MAP_KEY = "PMS_TECH_STATE_MAP"
FRESH_DAYS = 3 # 映射超过这么多自然日没刷新就当没读数 (调度断了不拿旧态拦人)
DEFAULT_PATH = "/api/v1/market/technical"
PAGE_SIZE = 1000 # 分页每页条数 (接口 limit 上限实测 1000, 设成一样避免误判最后一页)
MAX_PAGES = 20 # 防呆: 分页硬上限
class TechFeedError(RuntimeError):
"""拉技术面失败。调用方按无读数降级, 不产生任何拦截。"""
# ================================================================ 参数
def _getf(key: str, d: float) -> float:
v = param_store.get(key, None)
try:
return float(v) if v not in (None, "") else float(d)
except (TypeError, ValueError):
return float(d)
def _params() -> dict:
g, gi, gb = param_store.get, param_store.get_int, param_store.get_bool
base = (g("PMS_TECH_API_BASE", "") or "").strip().rstrip("/")
if not base: # 留空沿用研判接口根地址 (同一台 bionic)
base = (g("PMS_JUDGE_API_BASE", "") or "").strip().rstrip("/")
return {
"enabled": gb("PMS_TECH_ENABLED", True),
"base": base,
"path": (g("PMS_TECH_API_PATH", DEFAULT_PATH) or DEFAULT_PATH).strip(),
"timeout": gi("PMS_TECH_TIMEOUT", 20),
"keep_days": gi("PMS_TECH_KEEP_DAYS", 40),
"stale_tdays": gi("PMS_TECH_STALE_TDAYS", 2),
"rules": {
"squeeze_lookback": gi("PMS_TECH_SQUEEZE_LOOKBACK", 5),
"open_bw_growth": _getf("PMS_TECH_OPEN_BW_GROWTH", 0.20),
"choppy_flips": gi("PMS_TECH_CHOPPY_FLIPS", 4),
"flip_fresh_days": gi("PMS_TECH_FLIP_FRESH_DAYS", 2),
},
}
def enabled() -> bool:
return param_store.get_bool("PMS_TECH_ENABLED", True)
# ================================================================ 解析 (纯函数)
def _f(v):
try:
return None if v in (None, "") else float(v)
except (TypeError, ValueError):
return None
def _i(v):
try:
return None if v in (None, "") else int(float(v))
except (TypeError, ValueError):
return None
def _s(v):
if v is None:
return None
s = str(v).strip()
return s or None
def _tinyint(v):
"""接口布尔 → 1/0; None 保留 None (缺字段不误判成 0)。"""
if v is None:
return None
if isinstance(v, str):
return 1 if v.strip() in ("1", "true", "True", "") else 0
return 1 if v else 0
def parse_item(it: dict, data_date: int, algo_version=None):
"""接口一只票 → 落表行 dict。代码归一 (SH600000 → 600000.SH); 认不出返回 None。"""
if not isinstance(it, dict):
return None
code = normalize_code(str(it.get("stock_code") or it.get("code") or it.get("ts_code") or ""))
if not code or "." not in code:
return None
boll = it.get("boll") if isinstance(it.get("boll"), dict) else {}
bbi = it.get("bbiboll") if isinstance(it.get("bbiboll"), dict) else {}
sar = it.get("sar") if isinstance(it.get("sar"), dict) else {}
return {
"data_date": int(data_date), "ts_code": code,
"stock_name": _s(it.get("stock_name") or it.get("name")),
"boll_upper": _f(boll.get("upper")), "boll_mid": _f(boll.get("mid")),
"boll_lower": _f(boll.get("lower")), "boll_bw_pct": _f(boll.get("bandwidth_pct")),
"boll_pos": _f(boll.get("pos")), "boll_state": _s(boll.get("state")),
"boll_squeeze": _tinyint(boll.get("squeeze")),
"bbi": _f(bbi.get("bbi")), "bbi_upper": _f(bbi.get("upper")),
"bbi_lower": _f(bbi.get("lower")), "bbi_pos": _f(bbi.get("pos")),
"bbi_state": _s(bbi.get("state")), "bbi_dist_pct": _f(bbi.get("dist_pct")),
"sar_value": _f(sar.get("value")), "sar_side": _s(sar.get("side")),
"sar_flip_days": _i(sar.get("flip_days")), "sar_dist_pct": _f(sar.get("dist_pct")),
"reanchored": _tinyint(it.get("reanchored")), "in_pool": _tinyint(it.get("in_pool")),
"quality": _s(it.get("quality")), "bars_used": _i(it.get("bars_used")),
"last_bar_date": _i(it.get("last_bar_date")), "bars_lag": _i(it.get("bars_lag")),
"algo_version": _s(algo_version),
}
# ================================================================ 取数 (纯函数, 可注入 fetch)
def _http_get(url, params, timeout):
import requests
r = requests.get(url, params=(params or None), timeout=timeout)
r.raise_for_status()
return r.json()
def fetch_page(*, base, path, timeout, limit, offset, fetch=None) -> dict:
"""拉一页。fetch 可注入 (单测)。失败或应答非对象抛 TechFeedError。"""
base = (base or "").strip().rstrip("/")
if not base:
raise TechFeedError("技术面接口地址为空PMS_TECH_API_BASE 与 PMS_JUDGE_API_BASE 都没填)")
path = (path or DEFAULT_PATH).strip()
if not path.startswith("/"):
path = "/" + path
url = base + path
try:
payload = (fetch or _http_get)(url, {"limit": int(limit), "offset": int(offset)}, int(timeout))
except Exception as e:
logger.warning("拉技术面失败 %s offset=%s: %s: %s", url, offset, type(e).__name__, e)
raise TechFeedError(f"{type(e).__name__}: {e}") from e
if not isinstance(payload, dict):
raise TechFeedError(f"应答不是 JSON 对象: {type(payload).__name__}")
return payload
def fetch_all_pages(*, base=None, path=None, timeout=None, page_size=PAGE_SIZE, fetch=None) -> dict:
"""分页拉全市场。返回 {status, data_date, algo_version, rows, pages, total, matched}。
任一页 status 非 OK 或抛错 → 整轮失败 (抛 TechFeedError), 不返回半份。"""
if base is None or path is None or timeout is None:
p = _params()
base = p["base"] if base is None else base
path = p["path"] if path is None else path
timeout = p["timeout"] if timeout is None else timeout
rows, offset, pages, matched = [], 0, 0, None
data_date = algo = None
while pages < MAX_PAGES:
payload = fetch_page(base=base, path=path, timeout=timeout,
limit=page_size, offset=offset, fetch=fetch)
status = str(payload.get("status") or "")
if status != "OK":
raise TechFeedError(f"接口状态非 OK: {status or ''}")
if data_date is None:
data_date = _i(payload.get("data_date"))
algo = _s(payload.get("algo_version"))
matched = _i(payload.get("matched"))
items = payload.get("items") if isinstance(payload.get("items"), list) else []
for it in items:
row = parse_item(it, data_date, algo)
if row:
rows.append(row)
pages += 1
got = len(items)
offset += got
if got == 0: # 空页, 到底了
break
if matched is not None and offset >= matched: # 拉够 matched 只
break
# 不能用 got < page_size 判最后一页: 接口 limit 上限可能小于 page_size, 那样每页都返回
# 上限条数、got 恒小于 page_size, 第一页就误判成最后一页 (2026-09-11 真机: 上限 1000、
# matched 4991, 旧逻辑只落了 1000/4991)。matched 缺失时才退回用 got 判, 另有 MAX_PAGES 兜底。
if matched is None and got < page_size:
break
return {"status": "OK", "data_date": data_date, "algo_version": algo,
"rows": rows, "pages": pages, "total": len(rows), "matched": matched}
# ================================================================ 落表
def pull(*, fetch=None) -> dict:
"""拉一轮全市场 + 落表 + 删旧。返回读数供调度日志。失败不落表、不推进日期。"""
out = {"ok": True, "at": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
"data_date": None, "rows": 0, "pruned": 0, "error": None}
p = _params()
if not p["enabled"]:
out.update(ok=False, error="PMS_TECH_ENABLED 关着, 本轮不拉取")
return out
try:
got = fetch_all_pages(base=p["base"], path=p["path"], timeout=p["timeout"], fetch=fetch)
except TechFeedError as e:
out.update(ok=False, error=str(e)[:300])
return out
dd, rows = got["data_date"], got["rows"]
if not dd or not rows:
out.update(ok=False, error="接口未给 data_date 或无读数行")
return out
try:
out["rows"] = tech_repo.upsert_daily(rows)
except Exception as e: # noqa: BLE001
out.update(ok=False, error=f"落表失败: {type(e).__name__}: {e}")
return out
out["data_date"] = dd
out["pages"], out["matched"] = got["pages"], got["matched"]
# 删旧: 保留 keep_days 个交易日, 门槛 = data_date 往前第 keep_days 个交易日
try:
before = tradedays.ymd(tradedays.prev_trade_day(dd, int(p["keep_days"])))
out["pruned"] = tech_repo.prune(before)
except Exception as e: # noqa: BLE001
out["prune_error"] = f"{type(e).__name__}: {e}"
return out
# ================================================================ 映射 (照 logic_state_service)
def load_map() -> dict:
raw = param_store.get(MAP_KEY, "") or ""
if not raw:
return {}
try:
m = json.loads(raw)
except (TypeError, ValueError):
return {}
return m if isinstance(m, dict) else {}
def _save_map(m: dict) -> dict:
r = param_store.set_param(MAP_KEY, json.dumps(m, ensure_ascii=False), "system")
if not r.get("ok"):
logger.error("[技术面] 映射写入失败: %s —— 扫描层沿用上一份 (最多 %d 天后自动失效)",
r.get("error"), FRESH_DAYS)
return r
def state_map(m: dict | None = None) -> dict:
"""{ts_code: 立场 dict} —— 只在映射新鲜时给; 超 FRESH_DAYS 自然日一律按无读数返回空。"""
m = load_map() if m is None else (m or {})
at = str(m.get("at") or "")[:10]
try:
age = (date.today() - date.fromisoformat(at)).days
except ValueError:
return {}
if age > FRESH_DAYS:
return {}
states = m.get("states")
return dict(states) if isinstance(states, dict) else {}
def _compact(st: dict) -> dict:
"""synthesize 的结果压成映射里存的紧凑态 (方案附录丁「逐票紧凑状态」)。"""
return {"stance": st.get("stance"), "strength": st.get("strength"),
"phase": st.get("phase"), "confirm": st.get("confirm"),
"sar_side": st.get("sar_side"), "sar_value": st.get("sar_value"),
"sar_flip_days": st.get("sar_flip_days"), "choppy": st.get("choppy"),
"reason": st.get("reason"), "data_date": st.get("data_date")}
def _tdays_elapsed(from_ymd, to_ymd) -> int:
"""from_ymd 到 to_ymd 之间经过的交易日数 (不含 from, 含 to); to<=from 返回 0。"""
n = tradedays.trade_days_left(to_ymd, from_ymd) # 含两端的交易日数
return max(0, n - 1) if n else 0
def _relevant_codes() -> list:
"""映射要覆盖的票: 在持 + 计划主榜观察 + 待拍板提议。任一来源失败只跳过该来源、记日志。"""
from app.repo import pms_repo
from app.services import plan_feed
codes = []
try:
codes += [r.get("ts_code") for r in pms_repo.list_positions(only_open=True)]
except Exception as e: # noqa: BLE001
logger.warning("[技术面] 取持仓失败: %s", e)
try:
plan = plan_feed.get_plan()
codes += [r.get("ts_code") for r in (plan.get("main") or [])]
codes += [r.get("ts_code") for r in (plan.get("observe") or [])]
except Exception as e: # noqa: BLE001
logger.warning("[技术面] 取计划失败: %s", e)
try:
codes += [r.get("ts_code") for r in
pms_repo.list_proposals(statuses=("WAIT_USER",), limit=200)]
except Exception as e: # noqa: BLE001
logger.warning("[技术面] 取待拍板提议失败: %s", e)
return [c for c in dict.fromkeys(normalize_code(c) for c in codes if c) if c]
def build_map(*, now=None) -> dict:
"""对相关票算技术面立场并写映射。与 logic_state_service.pull_for_held 同构:
任何一步失败都不抛, 取不到写空映射带原因 (扫描层据此按无读数)。"""
now = now or datetime.now()
stamp = now.strftime("%Y-%m-%d %H:%M:%S")
out = {"ok": True, "at": stamp, "data_date": None, "codes": 0, "states": 0,
"by_stance": {}, "error": None}
p = _params()
if not p["enabled"]:
_save_map({"at": stamp, "data_date": None, "states": {}, "note": "PMS_TECH_ENABLED 关着"})
out.update(ok=False, error="PMS_TECH_ENABLED 关着")
return out
try:
dd = tech_repo.latest_date()
except Exception as e: # noqa: BLE001
_save_map({"at": stamp, "data_date": None, "states": {}, "error": f"读最新日失败: {e}"})
out.update(ok=False, error=f"读最新日失败: {type(e).__name__}: {e}")
return out
if not dd:
_save_map({"at": stamp, "data_date": None, "states": {}, "note": "库里还没有技术面读数"})
out.update(ok=False, error="库里还没有技术面读数")
return out
out["data_date"] = dd
# 读数陈旧 (最新日距今超 stale_tdays 个交易日) → 整份按无读数, 只记原因
if _tdays_elapsed(dd, tradedays.ymd(now.date())) > int(p["stale_tdays"]):
_save_map({"at": stamp, "data_date": dd, "stale": True, "states": {},
"note": f"技术面读数 {dd} 距今超过 {p['stale_tdays']} 个交易日, 按无读数"})
out.update(ok=False, error="读数陈旧")
return out
codes = _relevant_codes()
out["codes"] = len(codes)
if not codes:
_save_map({"at": stamp, "data_date": dd, "states": {}, "note": "没有需要技术面的票"})
return out
since = tradedays.ymd(tradedays.prev_trade_day(dd, tech_rules.CHOPPY_WINDOW + 5))
try:
hist = tech_repo.history_multi(codes, since=since)
except Exception as e: # noqa: BLE001
_save_map({"at": stamp, "data_date": dd, "states": {}, "error": f"取历史失败: {e}"})
out.update(ok=False, error=f"取历史失败: {type(e).__name__}: {e}")
return out
states = {}
for c in codes:
rows = hist.get(c) or []
if not rows or int(rows[-1].get("data_date") or 0) != dd:
continue # 今日没有这只票的读数 → 无读数, 不拿旧行硬合成
states[c] = _compact(tech_rules.synthesize(rows, params=p["rules"]))
r = _save_map({"at": stamp, "data_date": dd, "states": states})
if not r.get("ok"):
out["error"] = f"映射写入失败: {r.get('error')}"
out["states"] = len(states)
out["by_stance"] = dict(Counter(s.get("stance") or "" for s in states.values()))
return out
def pull_and_map(*, fetch=None) -> dict:
"""拉取 + 落表 + 重建映射, 一把梭 (手动拉取接口与 08:40 补拉用)。落表失败就不建映射。"""
pr = pull(fetch=fetch)
mp = build_map() if pr.get("ok") else {"ok": False, "error": "未落表, 跳过建映射"}
return {"ok": bool(pr.get("ok")), "pull": pr, "map": mp}
# ================================================================ 页面挂载与状态
def attach(rows, states=None):
"""把技术面立场挂到行上 (原地改 r['tech'])。无读数的票挂 None, 页面按无读数显示。"""
states = state_map() if states is None else (states or {})
for r in rows or []:
r["tech"] = states.get(normalize_code(r.get("ts_code") or "")) or None
return rows
def _row_view(r: dict) -> dict:
"""一行落表读数 → 页面友好的三块 (布林/多空布林/SAR)。"""
return {
"data_date": r.get("data_date"),
"boll": {"upper": r.get("boll_upper"), "mid": r.get("boll_mid"), "lower": r.get("boll_lower"),
"bw_pct": r.get("boll_bw_pct"), "pos": r.get("boll_pos"),
"state": r.get("boll_state"), "squeeze": bool(r.get("boll_squeeze"))},
"bbi": {"bbi": r.get("bbi"), "pos": r.get("bbi_pos"),
"state": r.get("bbi_state"), "dist_pct": r.get("bbi_dist_pct")},
"sar": {"value": r.get("sar_value"), "side": r.get("sar_side"),
"flip_days": r.get("sar_flip_days"), "dist_pct": r.get("sar_dist_pct")},
"quality": r.get("quality"), "reanchored": bool(r.get("reanchored")),
}
def research_feed(ts_code: str) -> dict:
"""单票研究面的技术面一块: 最新读数 + 合成立场 + 近日翻向次数。点击时直接读库算。"""
code = normalize_code(ts_code or "")
if not code:
return {"tech": None, "error": "代码认不出"}
if not enabled():
return {"tech": None, "note": "PMS_TECH_ENABLED 关着"}
p = _params()
try:
dd = tech_repo.latest_date()
since = tradedays.ymd(tradedays.prev_trade_day(dd, tech_rules.CHOPPY_WINDOW + 5)) if dd else 0
rows = tech_repo.history(code, since=since, limit=tech_rules.CHOPPY_WINDOW + 5)
except Exception as e: # noqa: BLE001
return {"tech": None, "error": f"{type(e).__name__}: {e}"}
if not rows:
return {"tech": None, "note": "这只票没有技术面读数"}
st = tech_rules.synthesize(rows, params=p["rules"])
return {"tech": {**_compact(st), "latest": _row_view(rows[-1]),
"flips": tech_rules.count_sar_flips(rows), "bars": len(rows)}}
def status() -> dict:
"""技术面模块状态: 开关/最新读数日/今日行数/映射覆盖与新鲜度/错误。页面顶栏芯片与运维看。"""
p = _params()
out = {"enabled": p["enabled"], "base_set": bool(p["base"]), "data_date": None,
"rows_today": 0, "map": {}, "error": None}
if not p["enabled"]:
out["note"] = "PMS_TECH_ENABLED 关着"
return out
try:
dd = tech_repo.latest_date()
out["data_date"] = dd
out["rows_today"] = tech_repo.count_on(dd) if dd else 0
except Exception as e: # noqa: BLE001
out["error"] = f"读库失败: {type(e).__name__}: {e}"
m = load_map()
if m:
out["map"] = {"at": m.get("at"), "data_date": m.get("data_date"),
"states": len(m.get("states") or {}), "fresh": bool(state_map(m)),
"note": m.get("note"), "error": m.get("error"), "stale": m.get("stale")}
return out