Files
kis_bot/kis_trader/engine/feed_fallback.py
Your Name 94daf529f0 fix(정합성·뱃지): 틱 로더 recv_ts 전달 · _lag_sec 폴백 · use_rust 결과JSON 기록
Q1) 쓰레기 스킵 22.8% 여전한 이유 → 로더가 recv_ts 를 tick dict 에 안 넣음
- breakout_tick_loader._ingest_tick_rows: recv_ts 문자열 필드 추가
  (bar_is_garbage 가 wall-clock 정합 판정에 사용)
- feed_fallback.bar_is_garbage: _lag_sec 사전계산 필드 최우선 사용
  Rust Parquet 로더는 recv_ts 없이 _lag_sec 만 실어옴 → 이 폴백으로 커버
- 우선순위: _lag_sec > recv_ts > 봉끝(레거시)

Q2) 자동 등록 잡 (mode_refine phase1/phase2) 뱃지  뜨는 이유
- register_result_json_as_job 이 원본 JSON 의 use_rust 를 읽는데 없어서 None
- optuna_web_jobs: inherit_use_rust 인자 신설 → 부모 실행 잡 값 상속
- _spawn_job_reaper: register 호출 시 부모 use_rust 전달
- optuna_common.annotate_optuna_period_daily_avg: 결과 JSON 저장 훅에
  use_rust=bool(BACKTEST_USE_RUST=='1') 자동 삽입
  → 4전략(tail/momentum/breakout/scalp) 모두 커버 (근본)

우선순위: 원본 JSON use_rust > 부모 잡 상속 > None()

실매 스모크: logs/test_live_execution_validation_20260906_194154.log → 최종: 통과
웹 재시작 200. 다음 옵투나 실행부터 반영.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-06 19:42:17 +09:00

611 lines
22 KiB
Python

"""
kis_trader/engine/feed_fallback.py — 실매·옵투나 공통 읽기 나이 헬퍼
======================================================================
메인이 지금보다 LIVE_FEED_FALLBACK_MAX_AGE_SEC 초보다 오래면 그 벤더는 없는 것과 같음.
2차(보조 WS) → 3차(LS) → (실매 매도만) 4차 키움 REST.
적재(공책) 임계와 읽기 나이는 분리한다. OHLC 로 틱을 메우지 않음.
"""
from __future__ import annotations
from datetime import datetime
import time
from typing import Any, Dict, List, Optional, Sequence, Tuple
from kis_trader.utils.env import get_env_bool, get_env_float, get_env_from_db
def packet_lag_seconds(pkt_raw: str, *, now_dt: Optional[datetime] = None) -> Optional[float]:
"""체결시각 문자열 vs now_dt(기본=서버 지금). 증권사끼리 비교 아님.
YYYYMMDDHHMMSS 또는 HHMMSS. 파싱 실패면 None(나이를 모름 → 읽기에서 버리지 않음).
"""
tt = str(pkt_raw or "").strip()
if not tt:
return None
wall = now_dt if now_dt is not None else datetime.now()
pkt_dt = None
try:
if len(tt) >= 14 and tt[:14].isdigit():
pkt_dt = datetime.strptime(tt[:14], "%Y%m%d%H%M%S")
elif len(tt) >= 6 and tt[-6:].isdigit():
pkt_dt = datetime.strptime(
wall.strftime("%Y%m%d") + tt[-6:], "%Y%m%d%H%M%S",
)
except Exception:
return None
if pkt_dt is None:
return None
return float((wall - pkt_dt).total_seconds())
def is_feed_read_stale(
lag_sec: Optional[float], max_age_sec: Optional[float] = None,
) -> bool:
"""읽기 나이 초과. lag/나이를 모르면 False(버리지 않음)."""
age = live_feed_fallback_max_age_sec() if max_age_sec is None else float(max_age_sec)
if age <= 0 or lag_sec is None:
return False
try:
return float(lag_sec) > age
except (TypeError, ValueError):
return False
def is_orderbook_snap_time_stale(
snap_time: Any,
*,
max_age_sec: Optional[float] = None,
) -> bool:
"""호가 거래소 시각(snap_time) vs 지금 — 틱 ``packet_lag`` + ``is_feed_read_stale`` 와 동일.
snap_time 없거나 파싱 실패 → False(버리지 않음). LIVE_FEED_FALLBACK 0 이하면 OFF.
"""
return is_feed_read_stale(
packet_lag_seconds(str(snap_time or "").strip()),
max_age_sec=max_age_sec,
)
def live_feed_fallback_max_age_sec() -> float:
"""읽기 유효 나이(초). 기본 3. 0 이하면 나이로 벤더를 버리지 않음(레거시)."""
try:
v = float(get_env_float("LIVE_FEED_FALLBACK_MAX_AGE_SEC", 3.0) or 3.0)
except (TypeError, ValueError):
v = 3.0
return v
def live_tick_primary() -> str:
p = str(get_env_from_db("LIVE_TICK_PROVIDER", "kiwoom") or "kiwoom").strip().lower()
return p if p in ("kis", "kiwoom") else "kiwoom"
def live_ob_primary() -> str:
p = str(get_env_from_db("LIVE_OB_PROVIDER", "kiwoom") or "kiwoom").strip().lower()
return p if p in ("kis", "kiwoom") else "kiwoom"
def vendor_read_max_age_sec(caller_max_age: Optional[float]) -> Optional[float]:
"""벤더 1회 조회에 쓸 나이.
- caller None = 마지막 RAM(명시, 나이 무시)
- caller 0/음수 또는 생략 대체 = 폴백 나이(기본 3초)
- caller 양수 = min(caller, 폴백나이). 폴백 0이면 caller 그대로.
"""
fb = live_feed_fallback_max_age_sec()
if caller_max_age is None:
return None
try:
c = float(caller_max_age)
except (TypeError, ValueError):
c = 0.0
if fb <= 0:
return None if c <= 0 else c
if c <= 0:
return fb
return min(c, fb)
def tick_feed_tier(vendor: str, primary: str) -> int:
v = str(vendor or "").strip().lower()
p = str(primary or "kiwoom").strip().lower()
if v in ("kiwoom_rest", "rest", "ka10007"):
return 4
if v == "ls":
return 3
if v == p:
return 1
return 2
def format_vendor_label(vendor: str, tier: int = 0, spilled: bool = False) -> str:
v = str(vendor or "").strip().lower() or "?"
names = {
"kiwoom": "kiwoom",
"kis": "kis",
"ls": "ls",
"kiwoom_rest": "kiwoom_rest",
"rest": "kiwoom_rest",
"ka10007": "kiwoom_rest",
}
name = names.get(v, v)
extra: List[str] = []
tmap = {1: "1차", 2: "2차", 3: "3차", 4: "4차"}
tlab = tmap.get(int(tier or 0), "")
if tlab:
extra.append(tlab)
if spilled:
extra.append("spill")
if extra:
return f"{name}({','.join(extra)})"
return name
def format_mm_feed_line(tick_lab: str, ob_lab: str) -> str:
t = str(tick_lab or "").strip() or "?"
o = str(ob_lab or "").strip() or "?"
return f"시세: {t} | 호가: {o}"
def trigger_feed_detail_log_enabled() -> bool:
"""트리거(매수체크) 로그에 틱·호가 상세 추적 줄. 기본 ON."""
try:
return bool(get_env_bool("TRIGGER_FEED_DETAIL_LOG", True))
except Exception:
return True
def extract_tick_trace_fields(data: Optional[Dict[str, Any]]) -> Dict[str, Any]:
"""get_price dict → 로그용 현재가·틱타임·나이."""
out: Dict[str, Any] = {}
if not isinstance(data, dict):
return out
px = (
data.get("stck_prpr")
or data.get("price")
or data.get("cur_prc")
or data.get("close")
or 0
)
try:
out["price"] = float(str(px).replace(",", "") or 0)
except (TypeError, ValueError):
out["price"] = 0.0
tt = (
data.get("tick_time")
or data.get("tick_time_raw")
or data.get("chetime")
or data.get("cntr_tm")
or data.get("FID20")
or data.get("stck_cntg_hour")
or data.get("cntg_hour")
or ""
)
out["tick_time"] = str(tt or "").strip()
if not out["tick_time"]:
bs = str(data.get("kis_bsop_date_raw") or "").strip()
hr = str(
data.get("kis_cntg_hour_raw")
or data.get("chetime")
or data.get("kiwoom_fid20")
or ""
).strip()
if bs and hr:
out["tick_time"] = (bs + hr)[:14]
elif hr:
out["tick_time"] = hr
# 키움 캐시에 tick_time 키가 없어도 chetime만으로 lag 계산
if data.get("chetime") and not out.get("tick_time"):
out["tick_time"] = str(data.get("chetime") or "").strip()
if data.get("_age_ms") is not None:
try:
out["age_ms"] = int(data.get("_age_ms") or 0)
except (TypeError, ValueError):
pass
lag = packet_lag_seconds(out.get("tick_time") or "")
if lag is not None:
out["lag_sec"] = round(float(lag), 2)
if data.get("_feed_vendor"):
out["vendor"] = str(data.get("_feed_vendor") or "").strip().lower()
return out
def extract_ob_trace_fields(snap: Any) -> Dict[str, Any]:
"""호가 스냅샷 → bid/ask/스프레드/잔량비·스냅시각."""
out: Dict[str, Any] = {}
if snap is None:
return out
try:
if hasattr(snap, "best_bid"):
out["best_bid"] = int(snap.best_bid() or 0)
out["best_ask"] = int(snap.best_ask() or 0)
try:
out["spread_pct"] = round(float(snap.spread_pct() or 0), 3)
except Exception:
pass
try:
bq = int(getattr(snap, "total_bid_qty", 0) or 0)
aq = int(getattr(snap, "total_ask_qty", 0) or 0)
out["bid_qty"] = bq
out["ask_qty"] = aq
if aq > 0:
out["bid_ask_ratio"] = round(bq / float(aq), 3)
except Exception:
pass
out["snap_time"] = str(getattr(snap, "snap_time", "") or "").strip()
src = str(getattr(snap, "source", "") or "").strip()
if src:
out["ob_source"] = src
if getattr(snap, "ts", None):
try:
out["age_ms"] = int(max(0.0, (time.time() - float(snap.ts)) * 1000))
except Exception:
pass
return out
except Exception:
pass
if isinstance(snap, dict):
try:
out["best_bid"] = int(float(snap.get("best_bid") or snap.get("bidp1") or 0))
out["best_ask"] = int(float(snap.get("best_ask") or snap.get("askp1") or 0))
except (TypeError, ValueError):
pass
out["snap_time"] = str(snap.get("snap_time") or "").strip()
return out
def format_trigger_feed_trace(
tick_rec: Optional[Dict[str, Any]],
ob_rec: Optional[Dict[str, Any]],
*,
tick_primary: str = "",
ob_primary: str = "",
) -> str:
"""매수체크 로그 꼬리 — 벤더(1·2·3차) + 현재가 + 틱타임 + 호가."""
tp = str(tick_primary or live_tick_primary()).strip().lower() or "kiwoom"
op = str(ob_primary or live_ob_primary()).strip().lower() or "kiwoom"
alt_t = "kiwoom" if tp == "kis" else "kis"
alt_o = "kiwoom" if op == "kis" else "kis"
parts: List[str] = [
f"틱1차설정={tp}(체인 {tp}{alt_t}→ls)",
f"호가1차설정={op}(체인 {op}{alt_o}→ls)",
]
tr = tick_rec or {}
if tr.get("label") or tr.get("vendor"):
lab = str(tr.get("label") or format_vendor_label(
str(tr.get("vendor") or ""),
int(tr.get("tier") or tick_feed_tier(str(tr.get("vendor") or ""), tp)),
bool(tr.get("spilled")),
))
bit = [f"틱실제={lab}"]
if tr.get("price"):
try:
bit.append("px=%s" % int(float(tr.get("price") or 0)))
except (TypeError, ValueError):
bit.append("px=%s" % tr.get("price"))
if tr.get("tick_time"):
bit.append("t=%s" % tr.get("tick_time"))
if tr.get("lag_sec") is not None:
bit.append("lag=%ss" % tr.get("lag_sec"))
elif tr.get("age_ms") is not None:
bit.append("age=%sms" % tr.get("age_ms"))
parts.append(" ".join(bit))
else:
parts.append("틱실제=없음")
obr = ob_rec or {}
if obr.get("label") or obr.get("vendor"):
lab = str(obr.get("label") or format_vendor_label(
str(obr.get("vendor") or ""),
int(obr.get("tier") or tick_feed_tier(str(obr.get("vendor") or ""), op)),
bool(obr.get("spilled")),
))
bit = [f"호가실제={lab}"]
if obr.get("best_bid") or obr.get("best_ask"):
bit.append("bid=%s ask=%s" % (obr.get("best_bid") or 0, obr.get("best_ask") or 0))
if obr.get("spread_pct") is not None:
bit.append("spr=%s%%" % obr.get("spread_pct"))
if obr.get("bid_ask_ratio") is not None:
bit.append("or=%s" % obr.get("bid_ask_ratio"))
if obr.get("snap_time"):
bit.append("snap=%s" % obr.get("snap_time"))
if obr.get("ob_source"):
bit.append("src=%s" % obr.get("ob_source"))
if obr.get("age_ms") is not None:
bit.append("age=%sms" % obr.get("age_ms"))
parts.append(" ".join(bit))
else:
parts.append("호가실제=없음")
return " | ".join(parts)
def _tick_second_key(tick: Dict[str, Any]) -> str:
tt = str(tick.get("tick_time") or "")[:14]
if len(tt) >= 14:
return tt[:14]
if len(tt) >= 12:
return tt[:12] + "00"
return tt
def merge_ticks_time_axis_fallback(
ticks: Sequence[Dict[str, Any]],
*,
main_src: str = "kiwoom",
max_lag_sec: Optional[float] = None,
) -> List[Dict[str, Any]]:
"""같은 초에는 메인만. 메인·보조 모두 lag>나이면 그 초는 비움(OHLC 메우지 않음).
분봉에 메인 1건 있다고 보조를 통째로 버리지 않음 (실매 2초 실패→2차와 동일).
"""
if not ticks:
return list(ticks)
main = str(main_src or "kiwoom").strip().lower()
if main not in ("kis", "kiwoom"):
main = "kiwoom"
age = live_feed_fallback_max_age_sec() if max_lag_sec is None else float(max_lag_sec)
grouped: Dict[str, List[Dict[str, Any]]] = {}
order: List[str] = []
for t in ticks:
if not isinstance(t, dict):
continue
sk = _tick_second_key(t)
if not sk:
continue
if sk not in grouped:
grouped[sk] = []
order.append(sk)
grouped[sk].append(t)
out: List[Dict[str, Any]] = []
for sk in order:
bucket = grouped[sk]
mains: List[Dict[str, Any]] = []
aux: List[Dict[str, Any]] = []
for t in bucket:
src = str(t.get("source") or "").strip().lower()
lag = t.get("_lag_sec")
if is_feed_read_stale(lag, age):
continue
if src == main:
mains.append(t)
else:
aux.append(t)
if mains:
out.extend(mains)
else:
out.extend(aux)
return out
def merge_ls_ticks_third_fallback(
ticks_by_code: Dict[str, Dict[str, List[Dict[str, Any]]]],
ls_ticks_by_code: Dict[str, Dict[str, List[Dict[str, Any]]]],
*,
max_lag_sec: Optional[float] = None,
) -> int:
"""실매 틱 3차: 같은 초에 kis/kiwoom 이 이미 있으면 LS 안 넣음.
메인·2차가 비었거나(또는 나이컷으로 비움) 그 초만 ``ls_ws_ticks`` 로 채움.
``LIVE_FEED_FALLBACK_MAX_AGE_SEC``(기본 3) — lag 모르면 버리지 않음(실매와 동일).
반환=추가한 틱 수.
"""
if not ls_ticks_by_code:
return 0
age = live_feed_fallback_max_age_sec() if max_lag_sec is None else float(max_lag_sec)
added = 0
for code, ls_minutes in ls_ticks_by_code.items():
code_s = str(code or "").strip()
if not code_s or not isinstance(ls_minutes, dict):
continue
main_minutes = ticks_by_code.setdefault(code_s, {})
occupied: set = set()
for _mk, ticks in main_minutes.items():
for t in ticks or []:
if not isinstance(t, dict):
continue
sk = _tick_second_key(t)
if sk:
occupied.add(sk)
for mk, ls_ticks in ls_minutes.items():
mk_s = str(mk or "").strip()
if not mk_s:
continue
bucket = main_minutes.setdefault(mk_s, [])
for t in ls_ticks or []:
if not isinstance(t, dict):
continue
sk = _tick_second_key(t)
if not sk or sk in occupied:
continue
lag = t.get("_lag_sec")
if is_feed_read_stale(lag, age):
continue
row = dict(t)
row["source"] = "ls"
bucket.append(row)
occupied.add(sk)
added += 1
if bucket:
bucket.sort(key=lambda x: str(x.get("tick_time") or ""))
return added
def orderbook_row_lag_seconds(row: Dict[str, Any]) -> Optional[int]:
"""recv_ts vs snap_time 초 차이. 파싱 실패면 None(유효로 봄)."""
recv_ts = str(row.get("recv_ts") or "").strip()
st = str(row.get("snap_time") or "").strip()
if len(recv_ts) < 19 or len(st) < 6:
return None
try:
recv_dt = datetime.strptime(recv_ts[:19], "%Y-%m-%d %H:%M:%S")
except Exception:
return None
try:
if len(st) >= 14 and st[:14].isdigit():
pkt = datetime.strptime(st[:14], "%Y%m%d%H%M%S")
elif len(st) >= 6 and st[-6:].isdigit():
pkt = datetime.strptime(recv_ts[:10].replace("-", "") + st[-6:], "%Y%m%d%H%M%S")
else:
return None
return int((recv_dt - pkt).total_seconds())
except Exception:
return None
def candle_garbage_fallback_enabled() -> bool:
"""쓰레기 봉 스킵(2차·3차 소스 폴백) 활성 여부.
2026-09-06: 기본 True 복원. bar_is_garbage 를 wall-clock(recv_ts) 기준으로 정정 후
실매 RAM 3초컷과 동일 논리로 동작 → 실매↔백테 정합. (docs/정합성.md §9)
BACKTEST_CANDLE_GARBAGE_OFF=1 은 실험용 강제 OFF (일반 운영은 사용 금지).
"""
import os
if os.environ.get("BACKTEST_CANDLE_GARBAGE_OFF") == "1":
return False
try:
return bool(get_env_bool("CANDLE_GARBAGE_FALLBACK", True))
except Exception:
return True
def bar_end_datetime(candle_time: str, tf_min: int):
"""봉 시작 candle_time(YYYYMMDDHHMM) + tf → 봉 끝 datetime."""
from kis_trader.engine.candle_rollup import add_candle_minutes
end = add_candle_minutes(str(candle_time or "")[:12], int(tf_min or 1))
if not end or len(end) < 12:
return None
try:
return datetime.strptime(end[:12], "%Y%m%d%H%M")
except Exception:
return None
def tick_in_bar_bucket(tick_time: str, candle_time: str, tf_min: int) -> bool:
from kis_trader.engine.candle_rollup import add_candle_minutes
ct = str(candle_time or "").strip()[:12]
raw = str(tick_time or "").strip().replace(":", "").replace("-", "").replace(" ", "")
if len(ct) < 12:
return False
tmin = raw[:12] if len(raw) >= 12 else ""
if len(tmin) < 12 or not tmin.isdigit():
return False
end = add_candle_minutes(ct, int(tf_min or 1))
if not end:
return False
return ct <= tmin < end[:12]
def bar_is_garbage(
ticks: Sequence[Dict[str, Any]],
*,
candle_time: str,
tf_min: int,
source: str,
missing_policy: str = "hole",
) -> bool:
"""그 분·그 source 틱이 없거나 전부 수신시각(recv_ts) 대비 읽기나이 초과면 True(구멍).
2026-09-06 정정 (docs/정합성.md): 기준 시각을 **봉끝** → **각 틱의 recv_ts** 로 변경.
이유: 실매 RAM 3초컷은 `wall-clock(수신순간) - tick_time` 기준.
이전 봉끝 기준은 유동성 낮은 종목·봉 시작 근처 틱이 부당 스킵되어 실매와 갭 발생.
이제 실매·백테 동일 wall-clock lag 판정 → 정합.
fallback: recv_ts 없으면 옛 봉끝 기준(레거시 데이터 호환).
missing_policy:
hole — 0건이면 쓰레기 (옵투나 공책)
keep — 0건이면 유지 (실매 링이 그 분을 커버 못 할 때 호출측에서 keep)
증권사 시계끼리 비교하지 않음.
"""
if not candle_garbage_fallback_enabled():
return False
src = str(source or "").strip().lower()
bucket: List[Dict[str, Any]] = []
for t in ticks or []:
if not isinstance(t, dict):
continue
tsrc = str(t.get("source") or "").strip().lower()
if tsrc and src and tsrc != src:
continue
raw = t.get("tick_time_raw") or t.get("tick_time") or ""
if tick_in_bar_bucket(str(raw), candle_time, tf_min):
bucket.append(t)
if not bucket:
return str(missing_policy or "hole").strip().lower() == "hole"
bar_end = bar_end_datetime(candle_time, tf_min)
for t in bucket:
# 1) 사전 계산된 _lag_sec 가 있으면 최우선 (로더가 recv_ts - tick_time 을 이미 계산)
# None 이면 판정 불가 → 이 틱은 안전(False), 다음 폴백으로.
pre_lag = t.get("_lag_sec")
if pre_lag is not None:
try:
_pl = float(pre_lag)
if not is_feed_read_stale(_pl):
return False
# _pl 이 stale 이면 다음 틱으로 계속. 이 틱은 컷 대상.
continue
except (TypeError, ValueError):
pass # 아래 폴백으로
raw = str(t.get("tick_time_raw") or t.get("tick_time") or "")
# 2) 각 틱의 wall-clock 기준(recv_ts) 대비 lag = 실매 RAM 판정과 동일 논리
recv_dt: Optional[datetime] = None
rts = t.get("recv_ts")
if rts is not None and str(rts).strip():
try:
if isinstance(rts, datetime):
recv_dt = rts
else:
s = str(rts).strip().replace("T", " ")
recv_dt = datetime.strptime(s[:19], "%Y-%m-%d %H:%M:%S")
except Exception:
recv_dt = None
# 3) recv_ts 없거나 파싱 실패 → 봉끝 기준(레거시 폴백, 실매보다 보수적)
base_dt = recv_dt if recv_dt is not None else bar_end
if base_dt is None:
# 봉끝도 파싱 실패면 판정 포기 → False(버리지 않음)
return False
lag = packet_lag_seconds(raw, now_dt=base_dt)
if not is_feed_read_stale(lag):
return False
return True
def live_bar_is_garbage(
ticks: Sequence[Dict[str, Any]],
*,
candle_time: str,
tf_min: int,
source: str,
) -> bool:
"""실매 링버퍼: 그 소스 틱이 그 분까지 없으면 판정 불가 → False(유지).
LS 틱은 TickRecorder 링에 안 들어갈 수 있음. 다른 증권사 틱만으로 cover 하면
LS/2차 봉을 구멍으로 잘못 버린다. 소스별 cover 만 본다.
"""
if not candle_garbage_fallback_enabled():
return False
ct = str(candle_time or "")[:12]
if len(ct) < 12:
return False
src = str(source or "").strip().lower()
cover = False
for t in ticks or []:
if not isinstance(t, dict):
continue
tsrc = str(t.get("source") or "").strip().lower()
if tsrc and src and tsrc != src:
continue
raw = str(t.get("tick_time_raw") or t.get("tick_time") or "").strip()
digits = raw.replace(":", "").replace("-", "").replace(" ", "")
tmin = digits[:12] if len(digits) >= 12 else ""
if tmin and tmin.isdigit() and tmin <= ct:
cover = True
break
if not cover:
return False
return bar_is_garbage(
ticks, candle_time=candle_time, tf_min=tf_min, source=source,
missing_policy="hole",
)