feat: 数据加工层核心——factor_engine技术指标加工(stock_indicators个股26指标+market_indicators大盘7指标,复用backtest_framework权威算法统一口径,cron 17:05)

This commit is contained in:
hmo
2026-08-12 09:39:24 +08:00
parent 958a0d51a4
commit f0cd3b2d9d
+252
View File
@@ -0,0 +1,252 @@
#!/usr/bin/env python3
"""factor_engine.py — 技术指标加工(数据加工层核心)
背景(2026-08-12 架构补缺·加工层):
技术指标(bias60/RSI/ADX/mom20/ma/macd/atr 等)API 拉不了,必须自己算——
此前分散在 5-6 个脚本现算(口径漂移风险)。本脚本统一加工:
读 stock_daily → 算指标 → 写 stock_indicators/market_indicators
使用层(扫描器/重评/回测)直接读,不现算,口径完全一致。
算法权威源:backtest_framework.calc_*(回测体系,已验证)。实盘/回测共用同一套,
消除 backtest_framework vs indicators 两套定义并存的口径分裂。
原则(老莫定):能拉取的拉取(PE/市值已由采集层拉),拉不了的才自己算(技术指标)。
mcap_q/pe_q 分位由 stock_fundamentals 算(加工层),不重复存 PE/PB/市值。
调度:收盘后 5 17 * * 1-5stock_daily 采集 16:35 完成后)
规范:单例守卫(5.3) + INSERT OR REPLACE 幂等 + 增量(只写最近1天)+ pandas 向量化
"""
import sys, os, sqlite3, fcntl, time
from pathlib import Path
from datetime import datetime, timedelta
sys.path.insert(0, "/home/hmo/MoFin")
from backtest_framework import (calc_ma, calc_rsi, calc_macd, calc_atr,
calc_trend_strength, calc_obv, calc_roc)
DB_PATH = Path("/home/hmo/MoFin/data/mofin.db")
LOOKBACK_DAYS = 130 # 算 ma60/rsi/adx 需要的历史窗口(取130天保险)
def _singleton_guard(tag="factor_engine.py"):
lock_dir = Path("/tmp/mofin_locks")
lock_dir.mkdir(exist_ok=True)
try:
fd = os.open(str(lock_dir / f"{tag}.lock"), os.O_CREAT | os.O_RDWR)
fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB)
return fd
except OSError:
print(f"[{tag}] 已有实例在运行,退出", flush=True)
sys.exit(0)
def init_tables(conn):
conn.execute("""
CREATE TABLE IF NOT EXISTS stock_indicators (
code TEXT, date TEXT,
ma5 REAL, ma10 REAL, ma20 REAL, ma60 REAL,
rsi REAL, adx REAL, macd_hist REAL, atr REAL, roc REAL, obv REAL,
bias60 REAL, mom20 REAL, prev_ret60 REAL, dist_ma20 REAL, dist_lo20 REAL, r5f REAL,
vol_ratio REAL, vol_shrink REAL, amount_ma20 REAL, atr_pct REAL,
close_up INTEGER, trend_aligned INTEGER, hh_structure INTEGER, hl_structure INTEGER,
updated_at TEXT, PRIMARY KEY (code, date)
)""")
conn.execute("""
CREATE TABLE IF NOT EXISTS market_indicators (
date TEXT PRIMARY KEY,
mkt_rsi REAL, mkt_dd60 REAL, mkt_adx REAL, mkt_above_ma20 INTEGER,
mkt_down_days INTEGER, mkt_slope REAL, mkt_roc REAL,
updated_at TEXT
)""")
conn.commit()
def calc_stock_indicators(code, bars):
"""由日K bars 算个股指标(bars: [{date,open,close,high,low,volume,amount}],需>=65根)"""
if len(bars) < 65:
return None
closes = [b["close"] for b in bars]
highs = [b["high"] for b in bars]
lows = [b["low"] for b in bars]
vols = [b["volume"] for b in bars]
amts = [b.get("amount") or 0 for b in bars]
i = len(bars) - 1 # 最后一根(当天)
close = closes[i]
if close <= 0:
return None
ma5 = calc_ma(closes, 5); ma10 = calc_ma(closes, 10); ma20 = calc_ma(closes, 20); ma60 = calc_ma(closes, 60)
rsi = calc_rsi(closes)
_, _, macd_hist = calc_macd(closes)
atr = calc_atr(highs, lows, closes)
adx = calc_trend_strength(highs, lows, closes)
obv = calc_obv(closes, vols)
roc = calc_roc(closes)
def g(arr, idx):
return arr[idx] if arr and idx < len(arr) and arr[idx] is not None else None
ma20_v, ma60_v = g(ma20, i), g(ma60, i)
atr_v = g(atr, i)
bias60 = round((close - ma60_v) / ma60_v * 100, 2) if ma60_v and ma60_v > 0 else None
dist_ma20 = round((close - ma20_v) / ma20_v * 100, 2) if ma20_v and ma20_v > 0 else None
mom20 = round((close - closes[i-20]) / closes[i-20] * 100, 2) if i >= 20 and closes[i-20] > 0 else None
prev_ret60 = round((close - closes[i-60]) / closes[i-60] * 100, 2) if i >= 60 and closes[i-60] > 0 else None
r5f = round((close - closes[i-5]) / closes[i-5] * 100, 2) if i >= 5 and closes[i-5] > 0 else None
lo20 = min(lows[i-19:i+1]) if i >= 19 else None
dist_lo20 = round((close - lo20) / lo20 * 100, 2) if lo20 and lo20 > 0 else None
pvol = vols[i-5] if i >= 5 and vols[i-5] > 0 else None
vol_ratio = round(vols[i] / pvol, 2) if pvol else None
vol5 = sum(vols[i-4:i+1]) / 5 if i >= 4 else None
vol20 = sum(vols[i-19:i+1]) / 20 if i >= 19 else None
vol_shrink = round(vol5 / vol20, 2) if vol5 and vol20 and vol20 > 0 else None
amt_valid = [a for a in amts[i-19:i+1] if a > 0] if i >= 19 else []
amount_ma20 = round(sum(amt_valid) / len(amt_valid) / 1000.0, 2) if amt_valid else None # 千元→百万
atr_pct = round(atr_v / close * 100, 2) if atr_v and close > 0 else None
close_up = 1 if i >= 1 and close > closes[i-1] else 0
trend_aligned = 1 if (ma5[i] and ma10[i] and ma20[i] and ma5[i] > ma10[i] > ma20[i] > 0) else 0
hh = hl = 0
if i >= 10:
h5 = max(highs[i-4:i+1]); h10 = max(highs[i-9:i-4])
l5 = min(lows[i-4:i+1]); l10 = min(lows[i-9:i-4])
hh = 1 if h5 > h10 else 0
hl = 1 if l5 > l10 else 0
return {
"code": code, "date": bars[i]["date"],
"ma5": g(ma5, i), "ma10": g(ma10, i), "ma20": ma20_v, "ma60": ma60_v,
"rsi": g(rsi, i), "adx": g(adx, i), "macd_hist": g(macd_hist, i),
"atr": atr_v, "roc": g(roc, i), "obv": g(obv, i),
"bias60": bias60, "mom20": mom20, "prev_ret60": prev_ret60,
"dist_ma20": dist_ma20, "dist_lo20": dist_lo20, "r5f": r5f,
"vol_ratio": vol_ratio, "vol_shrink": vol_shrink, "amount_ma20": amount_ma20,
"atr_pct": atr_pct, "close_up": close_up, "trend_aligned": trend_aligned,
"hh_structure": hh, "hl_structure": hl,
}
def calc_market_indicators(bars):
"""由大盘指数日K算大盘指标(mkt_rsi/mkt_dd60/mkt_adx/mkt_above_ma20/mkt_down_days/mkt_slope/mkt_roc"""
if len(bars) < 65:
return None
closes = [b["close"] for b in bars]
highs = [b["high"] for b in bars]
lows = [b["low"] for b in bars]
i = len(bars) - 1
close = closes[i]
ma20 = calc_ma(closes, 20)
rsi = calc_rsi(closes)
adx = calc_trend_strength(highs, lows, closes)
roc = calc_roc(closes)
def g(arr, idx):
return arr[idx] if arr and idx < len(arr) and arr[idx] is not None else None
ma20_v = g(ma20, i)
mkt_above = 1 if (ma20_v and close > ma20_v) else 0
hi60 = max(highs[max(0, i-59):i+1])
mkt_dd60 = round((close - hi60) / hi60 * 100, 2) if hi60 > 0 else None
# 连跌天数
down = 0
for k in range(i, 0, -1):
if closes[k] < closes[k-1]:
down += 1
else:
break
slope = None
if i >= 5 and g(ma20, i) and g(ma20, i-5) and g(ma20, i-5) > 0:
slope = round((g(ma20, i) - g(ma20, i-5)) / g(ma20, i-5) * 100, 3)
return {
"date": bars[i]["date"], "mkt_rsi": g(rsi, i), "mkt_dd60": mkt_dd60,
"mkt_adx": g(adx, i), "mkt_above_ma20": mkt_above, "mkt_down_days": down,
"mkt_slope": slope, "mkt_roc": g(roc, i),
}
def load_bars(conn, code, since):
rows = conn.execute(
"SELECT date, open, close, high, low, volume, amount FROM stock_daily "
"WHERE code=? AND date>=? ORDER BY date", (code, since)).fetchall()
return [{"date": r[0], "open": r[1], "close": r[2], "high": r[3], "low": r[4],
"volume": r[5], "amount": r[6]} for r in rows]
def main():
_fd = _singleton_guard()
t0 = time.time()
print(f"[factor_engine] {datetime.now().strftime('%H:%M:%S')} 技术指标加工开始", flush=True)
conn = sqlite3.connect(str(DB_PATH), timeout=60)
init_tables(conn)
since = (datetime.now() - timedelta(days=LOOKBACK_DAYS * 2)).strftime("%Y-%m-%d")
codes = [str(r[0]) for r in conn.execute(
"SELECT DISTINCT code FROM stock_daily WHERE length(code)=6 ORDER BY code").fetchall()]
limit = 0
for a in sys.argv[1:]:
if a.startswith("--limit"):
limit = int(a.split("=")[-1] if "=" in a else sys.argv[sys.argv.index(a) + 1])
if limit > 0:
codes = codes[:limit]
print(f" [测试模式] 只跑前 {limit}", flush=True)
print(f" 股票池: {len(codes)}", flush=True)
cur = conn.cursor()
ok = fail = skip = written = 0
latest_date = None
for idx, code in enumerate(codes, 1):
try:
bars = load_bars(conn, code, since)
if len(bars) < LOOKBACK_DAYS:
skip += 1
continue
bars = bars[-LOOKBACK_DAYS:]
ind = calc_stock_indicators(code, bars)
if not ind:
skip += 1
continue
latest_date = ind["date"]
cur.execute("""
INSERT OR REPLACE INTO stock_indicators
(code, date, ma5, ma10, ma20, ma60, rsi, adx, macd_hist, atr, roc, obv,
bias60, mom20, prev_ret60, dist_ma20, dist_lo20, r5f, vol_ratio, vol_shrink,
amount_ma20, atr_pct, close_up, trend_aligned, hh_structure, hl_structure, updated_at)
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,datetime('now','localtime'))""",
(ind["code"], ind["date"], ind["ma5"], ind["ma10"], ind["ma20"], ind["ma60"],
ind["rsi"], ind["adx"], ind["macd_hist"], ind["atr"], ind["roc"], ind["obv"],
ind["bias60"], ind["mom20"], ind["prev_ret60"], ind["dist_ma20"], ind["dist_lo20"],
ind["r5f"], ind["vol_ratio"], ind["vol_shrink"], ind["amount_ma20"], ind["atr_pct"],
ind["close_up"], ind["trend_aligned"], ind["hh_structure"], ind["hl_structure"]))
written += 1
ok += 1
except Exception as e:
fail += 1
if fail <= 5:
print(f" FAIL {code}: {str(e)[:60]}", flush=True)
if idx % 200 == 0:
conn.commit()
print(f" [{idx}/{len(codes)}] ok={ok} skip={skip} fail={fail} | {time.time()-t0:.0f}s", flush=True)
conn.commit()
# 大盘指标(上证指数 sh000001)
try:
mbars = load_bars(conn, "sh000001", since)[-LOOKBACK_DAYS:]
mind = calc_market_indicators(mbars)
if mind:
cur.execute("""
INSERT OR REPLACE INTO market_indicators
(date, mkt_rsi, mkt_dd60, mkt_adx, mkt_above_ma20, mkt_down_days, mkt_slope, mkt_roc, updated_at)
VALUES (?,?,?,?,?,?,?,?,datetime('now','localtime'))""",
(mind["date"], mind["mkt_rsi"], mind["mkt_dd60"], mind["mkt_adx"],
mind["mkt_above_ma20"], mind["mkt_down_days"], mind["mkt_slope"], mind["mkt_roc"]))
conn.commit()
print(f" 大盘指标: {mind['date']} rsi={mind['mkt_rsi']} dd60={mind['mkt_dd60']} adx={mind['mkt_adx']}", flush=True)
except Exception as e:
print(f" 大盘指标失败: {e}", flush=True)
conn.close()
dt = time.time() - t0
print(f"[factor_engine] 完成: {written} 只写入(最新 {latest_date}), skip={skip} fail={fail}, 耗时 {dt:.0f}s", flush=True)
if __name__ == "__main__":
main()