"""일자별 틱·호가 수집 통계 (백테웹 통계 탭). 실매 벤더별 DB 적재량 + 틱↔호가 공백 비율. - 틱/호가 DB는 벤더마다 SAVE 플래그로 **각각** 적재한다 (2·3차 spill일 때만 아님). - MODE=tick: 체결 1건당 그 벤더 RAM 호가가 ``WS_ORDERBOOK_TICK_MAX_AGE_SEC``(기본 3초) 안일 때만 호가행. 스냅 없으면 그 틱의 호가행은 안 생김. - ``LIVE_FEED_FALLBACK``(기본 3초) 초과 틱은 매매 RAM·호가틱동기 skip (틱 DB는 기본 전부 적재). """ from __future__ import annotations import time as _time from datetime import datetime from typing import Any, Dict, List, Optional, Set # 통계 결과 메모리 캐시 (key: day8, value: (ts, result)) # 웹 탭 재조회 시 매번 100초대 쉷 나는 일일 직접 스캔을 피하기 위해 TTL 캠시 적용. _STATS_CACHE: Dict[str, tuple] = {} # {day8: (expire_ts, result)} _STATS_CACHE_TTL = float(300) # 5분 (ENV으로 오버라이드 가능) def _ymd8(day: str) -> str: s = (day or "").strip().replace("-", "")[:8] if len(s) == 8 and s.isdigit(): return s return datetime.now().strftime("%Y%m%d") def _like(day8: str) -> str: return f"{day8}%" def _day_bounds(day8: str) -> tuple: """일자 범위 tick_time/snap_time (VARCHAR14 숫자문자열 — LIKE 대신 range).""" return f"{day8}000000", f"{day8}235959" class _FeedStatsDaySlice: """하루치 ws_ticks/ws_orderbook/ls_ws_ticks/ls_ws_orderbook TEMP. lag·pick 집계는 이 TEMP만 스캔 → LS 풀스캔 제거. ls_ws_ticks TEMP는 생성 시 lag_sec(TIMESTAMPDIFF)까지 미리 계산해 저장. """ TICK_TMP = "tmp_fs_ticks" OB_TMP = "tmp_fs_ob" LS_TICK_TMP = "tmp_fs_ls_ticks" LS_OB_TMP = "tmp_fs_ls_ob" def __init__(self, db, day8: str) -> None: self.db = db self.day8 = day8 self.t0, self.t1 = _day_bounds(day8) self.like = _like(day8) self.ready = False self.ls_ready = False # LS TEMP 별도 플래그 (LS WS 꺼진 날도 KIS/키움 TEMP는 동작) self.tick_n = 0 self.ob_n = 0 self.ls_tick_n = 0 self.ls_ob_n = 0 # ls_ws_ticks 당일 ts 범위 (DATETIME 형) from datetime import datetime as _dt self._ls_d0 = _dt.strptime(day8, "%Y%m%d") self._ls_d1 = self._ls_d0.replace(hour=23, minute=59, second=59) def ensure(self) -> bool: """KIS/키움 TEMP + LS TEMP 모두 생성. 실패해도 KIS/키움 TEMP는 유지.""" self._ensure_kis_kiwoom() self._ensure_ls() return self.ready def _ensure_kis_kiwoom(self) -> bool: if self.ready: return True conn = self.db.conn try: conn.execute(f"DROP TEMPORARY TABLE IF EXISTS {self.TICK_TMP}") conn.execute(f"DROP TEMPORARY TABLE IF EXISTS {self.OB_TMP}") conn.execute( f"CREATE TEMPORARY TABLE {self.TICK_TMP} (" "id BIGINT, code VARCHAR(32), tick_time VARCHAR(14), " "tick_time_raw VARCHAR(64), source VARCHAR(16), recv_ts VARCHAR(30), " "KEY idx_src (source), KEY idx_recv (recv_ts)" ") AS " "SELECT id, code, tick_time, tick_time_raw, source, recv_ts " "FROM ws_ticks WHERE tick_time >= %s AND tick_time <= %s", (self.t0, self.t1), ) conn.execute( f"CREATE TEMPORARY TABLE {self.OB_TMP} (" "id BIGINT, code VARCHAR(32), snap_time VARCHAR(14), " "source VARCHAR(16), recv_ts VARCHAR(30), " "reject_code VARCHAR(64), strategy VARCHAR(64), " "KEY idx_src (source)" ") AS " "SELECT id, code, snap_time, source, recv_ts, reject_code, strategy " "FROM ws_orderbook WHERE snap_time >= %s AND snap_time <= %s", (self.t0, self.t1), ) r = conn.execute(f"SELECT COUNT(*) AS n FROM {self.TICK_TMP}").fetchone() self.tick_n = int((r.get("n") if isinstance(r, dict) else r[0]) or 0) r = conn.execute(f"SELECT COUNT(*) AS n FROM {self.OB_TMP}").fetchone() self.ob_n = int((r.get("n") if isinstance(r, dict) else r[0]) or 0) self.ready = True return True except Exception: self.ready = False return False def _ensure_ls(self) -> bool: """ls_ws_ticks/ls_ws_orderbook 당일치 TEMP 생성 (lag_sec 미리 계산 포함).""" if self.ls_ready: return True conn = self.db.conn d0, d1 = self._ls_d0, self._ls_d1 try: conn.execute(f"DROP TEMPORARY TABLE IF EXISTS {self.LS_TICK_TMP}") # ls_ws_ticks: tk, lag_sec 가 이미 계산되어 있다면 바로 사용 (폴백 지원) conn.execute( f"CREATE TEMPORARY TABLE {self.LS_TICK_TMP} (" "code VARCHAR(20), tk CHAR(14), lag_sec DOUBLE, " "KEY idx_code_tk (code, tk)" ") AS " "SELECT code, " " IFNULL(tk, CONCAT(DATE_FORMAT(ts,'%%Y%%m%%d'), RIGHT(IFNULL(chetime,'000000'),6))) AS tk, " " IFNULL(lag_sec, TIMESTAMPDIFF(SECOND, " " STR_TO_DATE(" " CONCAT(DATE_FORMAT(ts,'%%Y%%m%%d'), RIGHT(IFNULL(chetime,'000000'),6))," " '%%Y%%m%%d%%H%%i%%s'), " " ts)) AS lag_sec " "FROM ls_ws_ticks WHERE ts >= %s AND ts <= %s", (d0, d1), ) r = conn.execute(f"SELECT COUNT(*) AS n FROM {self.LS_TICK_TMP}").fetchone() self.ls_tick_n = int((r.get("n") if isinstance(r, dict) else r[0]) or 0) except Exception: self.ls_tick_n = 0 try: conn.execute(f"DROP TEMPORARY TABLE IF EXISTS {self.LS_OB_TMP}") # ls_ws_orderbook: tk, lag_sec 활용 conn.execute( f"CREATE TEMPORARY TABLE {self.LS_OB_TMP} (" "code VARCHAR(32), snap_time CHAR(14), lag_sec DOUBLE, " "KEY idx_code_snap (code, snap_time)" ") AS " "SELECT code, IFNULL(tk, LEFT(snap_time,14)) AS snap_time, " " IFNULL(lag_sec, TIMESTAMPDIFF(SECOND, " " STR_TO_DATE(LEFT(snap_time,14), '%%Y%%m%%d%%H%%i%%s'), " " recv_ts)) AS lag_sec " "FROM ls_ws_orderbook WHERE snap_time >= %s AND snap_time <= %s", (self.t0, self.t1), ) r = conn.execute(f"SELECT COUNT(*) AS n FROM {self.LS_OB_TMP}").fetchone() self.ls_ob_n = int((r.get("n") if isinstance(r, dict) else r[0]) or 0) except Exception: self.ls_ob_n = 0 self.ls_ready = True return True def tick_from(self) -> str: return self.TICK_TMP if self.ready else "ws_ticks" def ob_from(self) -> str: return self.OB_TMP if self.ready else "ws_orderbook" def ls_tick_from(self) -> str: return self.LS_TICK_TMP if self.ls_ready else None def ls_ob_from(self) -> str: return self.LS_OB_TMP if self.ls_ready else None def tick_filter_sql(self) -> tuple: if self.ready: return "", () return " WHERE tick_time >= %s AND tick_time <= %s", (self.t0, self.t1) def ob_filter_sql(self) -> tuple: if self.ready: return "", () return " WHERE snap_time >= %s AND snap_time <= %s", (self.t0, self.t1) def _row_nc(r: Any) -> Dict[str, Any]: if r is None: return {"n": 0, "codes": 0} if isinstance(r, dict): return {"n": int(r.get("n") or 0), "codes": int(r.get("c") or r.get("codes") or 0)} return {"n": int(r[0] or 0), "codes": int(r[1] or 0)} def _pct(num: float, den: float) -> Optional[float]: if den is None or float(den) <= 0: return None return round(100.0 * float(num) / float(den), 2) def _age_cut_row( *, channel: str, source: str, total: int, pass_n: int, fail_n: int, unknown_n: int, empty_minutes: int = 0, usable_share_pct: Optional[float] = None, pick_n: Optional[int] = None, pick_pct: Optional[float] = None, ) -> Dict[str, Any]: """옵투나/실매 읽기나이(LIVE_FEED_FALLBACK) 합격 요약 1행. usable% = 그 source 총행 대비(자사). 벤더끼리 합이 100이 아님. usable_share% = 같은 채널 usable 합 대비 분배. pick% = LIVE_*_PROVIDER 체인 시뮬레이션 간택(종목×초). """ usable = int(pass_n) + int(unknown_n) # lag 미상=실매와 같이 유지 tot = int(total or 0) return { "channel": channel, "source": source, "total": tot, "pass_n": int(pass_n), "fail_n": int(fail_n), "unknown_n": int(unknown_n), "usable_n": usable, "empty_minutes": int(empty_minutes or 0), "pass_pct": _pct(pass_n, tot), "fail_pct": _pct(fail_n, tot), "usable_pct": _pct(usable, tot), "usable_share_pct": usable_share_pct, "pick_n": pick_n, "pick_pct": pick_pct, } def _vendor_read_chain(primary: str) -> List[str]: """실매 get_price/get_orderbook 체인: 1차 → 나머지 → ls.""" p = str(primary or "kiwoom").strip().lower() if p not in ("kis", "kiwoom"): p = "kiwoom" alt = "kis" if p == "kiwoom" else "kiwoom" return [p, alt, "ls"] def _ob_src_to_vendor(src: str) -> Optional[str]: s = str(src or "").strip().lower() if s in ("kis_h0stasp0", "kis"): return "kis" if s in ("kiwoom_0d", "kiwoom"): return "kiwoom" if s in ("ls_uh1", "ls"): return "ls" return None def _pick_winner(flags: Dict[str, bool], chain: List[str]) -> Optional[str]: for v in chain: if flags.get(v): return v return None def _count_tick_picks( db, day8: str, age: float, chain: List[str], day_slice: Optional[_FeedStatsDaySlice] = None, ) -> Dict[str, int]: """종목×초(YYYYMMDDHHMMSS)마다 usable 벤더 중 체인 1순위 간택 횟수.""" wins = {"kis": 0, "kiwoom": 0, "ls": 0} ds = day_slice or _FeedStatsDaySlice(db, day8) tick_tbl = ds.tick_from() tick_where, tick_params = ds.tick_filter_sql() exch_dt = _sql_ws_tick_exchange_dt() slots: Dict[tuple, Dict[str, bool]] = {} try: sql = ( "SELECT code, LEFT(tick_time, 14) AS tk, source, " "MAX(CASE WHEN lag_sec IS NULL OR lag_sec <= %s THEN 1 ELSE 0 END) AS ok " "FROM (" " SELECT code, tick_time, source, " f" TIMESTAMPDIFF(SECOND, {exch_dt}, recv_ts) AS lag_sec " f" FROM {tick_tbl}{tick_where}" ") x GROUP BY code, LEFT(tick_time, 14), source" ) for r in db.conn.execute(sql, (age, *tick_params)).fetchall() or []: d = dict(r) if not isinstance(r, dict) else r src = str(d.get("source") or "").strip().lower() if src not in wins: continue if int(d.get("ok") or 0) <= 0: continue key = (str(d.get("code") or ""), str(d.get("tk") or "")) if not key[0] or len(key[1]) < 14: continue slots.setdefault(key, {})[src] = True except Exception: pass try: ls_tbl = ds.ls_tick_from() if ls_tbl: # TEMP: tk 콜럼 이미 계산(YYYYMMDDHHMMSS), lag_sec 이미 계산 ls_sql = ( "SELECT code, tk, " "MAX(CASE WHEN lag_sec IS NULL OR lag_sec <= %s THEN 1 ELSE 0 END) AS ok " f"FROM {ls_tbl} " "GROUP BY code, tk" ) for r in db.conn.execute(ls_sql, (age,)).fetchall() or []: d = dict(r) if not isinstance(r, dict) else r if int(d.get("ok") or 0) <= 0: continue key = (str(d.get("code") or ""), str(d.get("tk") or "")) if not key[0] or len(key[1]) < 14: continue slots.setdefault(key, {})["ls"] = True else: # TEMP 없으면 본 테이블 폴백 (REGEXP 제거한 단순 CONCAT 버전) d0 = datetime.strptime(day8, "%Y%m%d") d1 = d0.replace(hour=23, minute=59, second=59) ls_sql = ( "SELECT code, " " CONCAT(DATE_FORMAT(ts,'%%Y%%m%%d'), RIGHT(IFNULL(chetime,'000000'),6)) AS tk, " "MAX(CASE WHEN lag_sec IS NULL OR lag_sec <= %s THEN 1 ELSE 0 END) AS ok " "FROM (" " SELECT code, " " IFNULL(tk, CONCAT(DATE_FORMAT(ts,'%%Y%%m%%d'), RIGHT(IFNULL(chetime,'000000'),6))) AS tk, " " IFNULL(lag_sec, TIMESTAMPDIFF(SECOND, " " STR_TO_DATE(" " CONCAT(DATE_FORMAT(ts,'%%Y%%m%%d'), RIGHT(IFNULL(chetime,'000000'),6))," " '%%Y%%m%%d%%H%%i%%s'), " " ts)) AS lag_sec " " FROM ls_ws_ticks WHERE ts>=%s AND ts<=%s" ") t GROUP BY code, tk" ) for r in db.conn.execute(ls_sql, (age, d0, d1)).fetchall() or []: d = dict(r) if not isinstance(r, dict) else r if int(d.get("ok") or 0) <= 0: continue key = (str(d.get("code") or ""), str(d.get("tk") or "")) if not key[0] or len(key[1]) < 14: continue slots.setdefault(key, {})["ls"] = True except Exception: pass for flags in slots.values(): w = _pick_winner(flags, chain) if w and w in wins: wins[w] += 1 return wins def _count_ob_picks( db, day8: str, age: float, chain: List[str], day_slice: Optional[_FeedStatsDaySlice] = None, ) -> Dict[str, int]: """종목×초마다 호가 usable 벤더 체인 간택.""" wins = {"kis": 0, "kiwoom": 0, "ls": 0} ds = day_slice or _FeedStatsDaySlice(db, day8) ob_tbl = ds.ob_from() ob_where, ob_params = ds.ob_filter_sql() t0, t1 = ds.t0, ds.t1 slots: Dict[tuple, Dict[str, bool]] = {} try: if ds.ready: ob_src_where = f" FROM {ob_tbl} WHERE source<>%s" ob_src_params: tuple = ("filter_eval",) else: ob_src_where = f" FROM {ob_tbl}{ob_where} AND source<>%s" ob_src_params = (*ob_params, "filter_eval") sql = ( "SELECT code, tk, source, " "MAX(CASE WHEN lag_sec IS NULL OR lag_sec <= %s THEN 1 ELSE 0 END) AS ok " "FROM (" " SELECT code, source, " " CASE WHEN CHAR_LENGTH(snap_time)>=14 THEN LEFT(snap_time,14) " " WHEN CHAR_LENGTH(snap_time)>=6 THEN CONCAT(%s, RIGHT(snap_time,6)) " " ELSE NULL END AS tk, " " TIMESTAMPDIFF(SECOND, " " STR_TO_DATE(" " CASE WHEN CHAR_LENGTH(snap_time)>=14 THEN LEFT(snap_time,14) " " WHEN CHAR_LENGTH(snap_time)>=6 THEN CONCAT(%s, RIGHT(snap_time,6)) " " ELSE NULL END, " " '%%Y%%m%%d%%H%%i%%s'), " " recv_ts) AS lag_sec " f"{ob_src_where}" ") x WHERE tk IS NOT NULL " "GROUP BY code, tk, source" ) for r in db.conn.execute( sql, (age, day8, day8, *ob_src_params) ).fetchall() or []: d = dict(r) if not isinstance(r, dict) else r vend = _ob_src_to_vendor(str(d.get("source") or "")) if not vend or int(d.get("ok") or 0) <= 0: continue key = (str(d.get("code") or ""), str(d.get("tk") or "")) if not key[0] or len(key[1]) < 14: continue slots.setdefault(key, {})[vend] = True except Exception: pass try: ls_ob_tbl = ds.ls_ob_from() if ls_ob_tbl: # TEMP 사용: lag_sec 이미 계산됨 ls_sql = ( "SELECT code, snap_time AS tk, " "MAX(CASE WHEN lag_sec IS NULL OR lag_sec <= %s THEN 1 ELSE 0 END) AS ok " f"FROM {ls_ob_tbl} " "GROUP BY code, snap_time" ) for r in db.conn.execute(ls_sql, (age,)).fetchall() or []: d = dict(r) if not isinstance(r, dict) else r if int(d.get("ok") or 0) <= 0: continue key = (str(d.get("code") or ""), str(d.get("tk") or "")) if not key[0] or len(key[1]) < 14: continue slots.setdefault(key, {})["ls"] = True else: ls_sql = ( "SELECT code, LEFT(snap_time,14) AS tk, " "MAX(CASE WHEN lag_sec IS NULL OR lag_sec <= %s THEN 1 ELSE 0 END) AS ok " "FROM (" " SELECT code, IFNULL(tk, LEFT(snap_time,14)) AS tk, " " IFNULL(lag_sec, TIMESTAMPDIFF(SECOND, " " STR_TO_DATE(LEFT(snap_time,14), '%%Y%%m%%d%%H%%i%%s'), " " recv_ts)) AS lag_sec " " FROM ls_ws_orderbook WHERE snap_time >= %s AND snap_time <= %s" ") t GROUP BY code, tk" ) for r in db.conn.execute(ls_sql, (age, t0, t1)).fetchall() or []: d = dict(r) if not isinstance(r, dict) else r if int(d.get("ok") or 0) <= 0: continue key = (str(d.get("code") or ""), str(d.get("tk") or "")) if not key[0] or len(key[1]) < 14: continue slots.setdefault(key, {})["ls"] = True except Exception: pass for flags in slots.values(): w = _pick_winner(flags, chain) if w and w in wins: wins[w] += 1 return wins def _annotate_share_and_pick( rows: List[Dict[str, Any]], *, tick_wins: Dict[str, int], ob_wins: Dict[str, int], tick_chain: List[str], ob_chain: List[str], ) -> None: """행에 usable 분배% · 간택(n/%) 주입 (Σ 제외 벤더끼리 합=100).""" tick_vendors = [r for r in rows if r.get("channel") == "tick" and r.get("source") != "Σ(틱)"] ob_vendors = [r for r in rows if r.get("channel") == "orderbook"] tick_usable_sum = sum(int(r.get("usable_n") or 0) for r in tick_vendors) or 0 ob_usable_sum = sum(int(r.get("usable_n") or 0) for r in ob_vendors) or 0 tick_pick_sum = sum(int(tick_wins.get(v) or 0) for v in ("kis", "kiwoom", "ls")) ob_pick_sum = sum(int(ob_wins.get(v) or 0) for v in ("kis", "kiwoom", "ls")) for r in rows: ch = r.get("channel") src = str(r.get("source") or "") if ch == "tick" and src == "Σ(틱)": r["usable_share_pct"] = 100.0 if tick_usable_sum else None r["pick_n"] = tick_pick_sum r["pick_pct"] = 100.0 if tick_pick_sum else None continue if ch == "tick": r["usable_share_pct"] = _pct(int(r.get("usable_n") or 0), tick_usable_sum) pn = int(tick_wins.get(src) or 0) if src in tick_wins else None r["pick_n"] = pn r["pick_pct"] = _pct(pn or 0, tick_pick_sum) if pn is not None else None r["pick_chain"] = "→".join(tick_chain) elif ch == "orderbook": r["usable_share_pct"] = _pct(int(r.get("usable_n") or 0), ob_usable_sum) vend = _ob_src_to_vendor(src) if vend: pn = int(ob_wins.get(vend) or 0) r["pick_n"] = pn r["pick_pct"] = _pct(pn, ob_pick_sum) r["pick_chain"] = "→".join(ob_chain) else: r["pick_n"] = None r["pick_pct"] = None def _sql_ws_tick_exchange_dt() -> str: """ws_ticks 거래소 체결시각 → DATETIME 식 (packet_lag_seconds 와 동일 축). - ``tick_time`` 14자리(YYYYMMDDHHMMSS) 우선 — 키움은 DB에 일자+FID20 로 저장 - ``tick_time_raw`` 만 14자리면 그대로 - raw 가 HHMMSS(키움 FID20)면 일자(tick_time 앞 8 또는 recv_ts)와 결합 주의: raw 를 LEFT14 로만 파싱하면 키움 6자리가 STR_TO_DATE 실패 → 전량 미상. """ return ( "STR_TO_DATE(" " CASE" " WHEN CHAR_LENGTH(IFNULL(tick_time,'')) >= 14 THEN LEFT(tick_time, 14)" " WHEN CHAR_LENGTH(IFNULL(tick_time_raw,'')) >= 14 THEN LEFT(tick_time_raw, 14)" " WHEN CHAR_LENGTH(IFNULL(tick_time_raw,'')) >= 6 THEN" " CONCAT(" " COALESCE(" " NULLIF(LEFT(IFNULL(tick_time,''), 8), '')," " DATE_FORMAT(recv_ts, '%%Y%%m%%d')" " )," " RIGHT(tick_time_raw, 6)" " )" " ELSE NULL" " END," " '%%Y%%m%%d%%H%%i%%s')" ) def _feed_age_cut_stats( db, day8: str, age_sec: float, day_slice: Optional[_FeedStatsDaySlice] = None, ) -> Dict[str, Any]: """recv_ts vs 체결/스냅 시각 지연 → 3초(폴백나이) 합격·탈락·미상. 옵투나 ``merge_ticks_time_axis_fallback`` / 호가 lag 컷과 동일 축. lag 미상(NULL)은 실매와 같이 **버리지 않음** → usable 에 포함. empty_minutes = 그 분 틱이 전부 fail 인 (code,분) 수 (비어 재현되는 분). """ age = max(0.0, float(age_sec or 3.0)) ds = day_slice or _FeedStatsDaySlice(db, day8) tick_tbl = ds.tick_from() ob_tbl = ds.ob_from() tick_where, tick_params = ds.tick_filter_sql() ob_where, ob_params = ds.ob_filter_sql() t0, t1 = ds.t0, ds.t1 rows_out: List[Dict[str, Any]] = [] exch_dt = _sql_ws_tick_exchange_dt() # ── ws_ticks ── try: tick_sql = ( "SELECT source AS src, " "COUNT(*) AS total, " "SUM(CASE WHEN lag_sec IS NOT NULL AND lag_sec <= %s THEN 1 ELSE 0 END) AS pass_n, " "SUM(CASE WHEN lag_sec IS NOT NULL AND lag_sec > %s THEN 1 ELSE 0 END) AS fail_n, " "SUM(CASE WHEN lag_sec IS NULL THEN 1 ELSE 0 END) AS unknown_n " "FROM (" " SELECT source, " f" TIMESTAMPDIFF(SECOND, {exch_dt}, recv_ts) AS lag_sec " f" FROM {tick_tbl}{tick_where}" ") t GROUP BY source" ) for r in db.conn.execute(tick_sql, (age, age, *tick_params)).fetchall() or []: d = dict(r) if not isinstance(r, dict) else r rows_out.append( _age_cut_row( channel="tick", source=str(d.get("src") or ""), total=int(d.get("total") or 0), pass_n=int(d.get("pass_n") or 0), fail_n=int(d.get("fail_n") or 0), unknown_n=int(d.get("unknown_n") or 0), ) ) except Exception: pass # 전량 컷된 분 (틱 채널) — 옵투나에서 그 분 시세 공백 empty_min = 0 try: empty_sql = ( "SELECT COUNT(*) AS n FROM (" " SELECT code, LEFT(tick_time,12) AS mk " " FROM (" " SELECT code, tick_time, " f" TIMESTAMPDIFF(SECOND, {exch_dt}, recv_ts) AS lag_sec " f" FROM {tick_tbl}{tick_where}" " ) x " " GROUP BY code, LEFT(tick_time,12) " " HAVING SUM(CASE WHEN lag_sec IS NULL OR lag_sec <= %s THEN 1 ELSE 0 END)=0 " " AND COUNT(*)>0" ") z" ) r = db.conn.execute(empty_sql, (*tick_params, age)).fetchone() empty_min = int((r["n"] if isinstance(r, dict) else r[0]) or 0) except Exception: empty_min = 0 # ── ls_ws_ticks ── try: ls_tbl = ds.ls_tick_from() if ls_tbl: # TEMP 사용: lag_sec 이미 계산됨 → REGEXP 연산 없음 ls_sql = ( "SELECT 'ls' AS src, " "COUNT(*) AS total, " "SUM(CASE WHEN lag_sec IS NOT NULL AND lag_sec <= %s THEN 1 ELSE 0 END) AS pass_n, " "SUM(CASE WHEN lag_sec IS NOT NULL AND lag_sec > %s THEN 1 ELSE 0 END) AS fail_n, " "SUM(CASE WHEN lag_sec IS NULL THEN 1 ELSE 0 END) AS unknown_n " f"FROM {ls_tbl}" ) r = db.conn.execute(ls_sql, (age, age)).fetchone() else: # TEMP 없으면 본 테이블 폴백 d0 = datetime.strptime(day8, "%Y%m%d") d1 = d0.replace(hour=23, minute=59, second=59) ls_sql = ( "SELECT 'ls' AS src, " "COUNT(*) AS total, " "SUM(CASE WHEN lag_sec IS NOT NULL AND lag_sec <= %s THEN 1 ELSE 0 END) AS pass_n, " "SUM(CASE WHEN lag_sec IS NOT NULL AND lag_sec > %s THEN 1 ELSE 0 END) AS fail_n, " "SUM(CASE WHEN lag_sec IS NULL THEN 1 ELSE 0 END) AS unknown_n " "FROM (" " SELECT IFNULL(lag_sec, TIMESTAMPDIFF(SECOND, " " STR_TO_DATE(" " CASE WHEN CHAR_LENGTH(REGEXP_REPLACE(IFNULL(chetime,''), '[^0-9]', ''))>=14 " " THEN LEFT(REGEXP_REPLACE(chetime, '[^0-9]', ''), 14) " " WHEN CHAR_LENGTH(REGEXP_REPLACE(IFNULL(chetime,''), '[^0-9]', ''))>=6 " " THEN CONCAT(DATE_FORMAT(ts,'%%Y%%m%%d'), " " RIGHT(REGEXP_REPLACE(chetime, '[^0-9]', ''), 6)) " " ELSE DATE_FORMAT(ts,'%%Y%%m%%d%%H%%i%%s') END, " " '%%Y%%m%%d%%H%%i%%s'), " " ts)) AS lag_sec " " FROM ls_ws_ticks WHERE ts>=%s AND ts<=%s" ") t" ) r = db.conn.execute(ls_sql, (age, age, d0, d1)).fetchone() if r: d = dict(r) if not isinstance(r, dict) else r if int(d.get("total") or 0) > 0: rows_out.append( _age_cut_row( channel="tick", source="ls", total=int(d.get("total") or 0), pass_n=int(d.get("pass_n") or 0), fail_n=int(d.get("fail_n") or 0), unknown_n=int(d.get("unknown_n") or 0), ) ) else: rows_out.append( _age_cut_row( channel="tick", source="ls", total=0, pass_n=0, fail_n=0, unknown_n=0, ) ) except Exception: rows_out.append( _age_cut_row( channel="tick", source="ls", total=0, pass_n=0, fail_n=0, unknown_n=0, ) ) # ── ws_orderbook 스트림 ── try: if ds.ready: ob_stream_from = f" FROM {ob_tbl} WHERE source<>%s" ob_stream_params: tuple = ("filter_eval",) else: ob_stream_from = f" FROM {ob_tbl}{ob_where} AND source<>%s" ob_stream_params = (*ob_params, "filter_eval") ob_sql = ( "SELECT source AS src, " "COUNT(*) AS total, " "SUM(CASE WHEN lag_sec IS NOT NULL AND lag_sec <= %s THEN 1 ELSE 0 END) AS pass_n, " "SUM(CASE WHEN lag_sec IS NOT NULL AND lag_sec > %s THEN 1 ELSE 0 END) AS fail_n, " "SUM(CASE WHEN lag_sec IS NULL THEN 1 ELSE 0 END) AS unknown_n " "FROM (" " SELECT source, " " TIMESTAMPDIFF(SECOND, " " STR_TO_DATE(" " CASE WHEN CHAR_LENGTH(snap_time)>=14 THEN LEFT(snap_time,14) " " WHEN CHAR_LENGTH(snap_time)>=6 THEN CONCAT(LEFT(%s,8), RIGHT(snap_time,6)) " " ELSE NULL END, " " '%%Y%%m%%d%%H%%i%%s'), " " recv_ts" " ) AS lag_sec " f"{ob_stream_from}" ") t GROUP BY source" ) for r in db.conn.execute( ob_sql, (age, age, day8, *ob_stream_params) ).fetchall() or []: d = dict(r) if not isinstance(r, dict) else r rows_out.append( _age_cut_row( channel="orderbook", source=str(d.get("src") or ""), total=int(d.get("total") or 0), pass_n=int(d.get("pass_n") or 0), fail_n=int(d.get("fail_n") or 0), unknown_n=int(d.get("unknown_n") or 0), ) ) except Exception: pass # LS 호가 try: ls_ob_tbl = ds.ls_ob_from() if ls_ob_tbl: # TEMP 사용: lag_sec 이미 계산됨 ls_ob_sql = ( "SELECT 'ls_uh1' AS src, " "COUNT(*) AS total, " "SUM(CASE WHEN lag_sec IS NOT NULL AND lag_sec <= %s THEN 1 ELSE 0 END) AS pass_n, " "SUM(CASE WHEN lag_sec IS NOT NULL AND lag_sec > %s THEN 1 ELSE 0 END) AS fail_n, " "SUM(CASE WHEN lag_sec IS NULL THEN 1 ELSE 0 END) AS unknown_n " f"FROM {ls_ob_tbl}" ) r = db.conn.execute(ls_ob_sql, (age, age)).fetchone() else: ls_ob_sql = ( "SELECT 'ls_uh1' AS src, " "COUNT(*) AS total, " "SUM(CASE WHEN lag_sec IS NOT NULL AND lag_sec <= %s THEN 1 ELSE 0 END) AS pass_n, " "SUM(CASE WHEN lag_sec IS NOT NULL AND lag_sec > %s THEN 1 ELSE 0 END) AS fail_n, " "SUM(CASE WHEN lag_sec IS NULL THEN 1 ELSE 0 END) AS unknown_n " "FROM (" " SELECT IFNULL(lag_sec, TIMESTAMPDIFF(SECOND, " " STR_TO_DATE(LEFT(snap_time,14), '%%Y%%m%%d%%H%%i%%s'), " " recv_ts)) AS lag_sec " " FROM ls_ws_orderbook WHERE snap_time >= %s AND snap_time <= %s" ") t" ) r = db.conn.execute(ls_ob_sql, (age, age, t0, t1)).fetchone() if r: d = dict(r) if not isinstance(r, dict) else r rows_out.append( _age_cut_row( channel="orderbook", source="ls_uh1", total=int(d.get("total") or 0), pass_n=int(d.get("pass_n") or 0), fail_n=int(d.get("fail_n") or 0), unknown_n=int(d.get("unknown_n") or 0), ) ) except Exception: rows_out.append( _age_cut_row( channel="orderbook", source="ls_uh1", total=0, pass_n=0, fail_n=0, unknown_n=0, ) ) # 틱 합계 행 tick_rows = [x for x in rows_out if x.get("channel") == "tick"] if tick_rows: rows_out.append( _age_cut_row( channel="tick", source="Σ(틱)", total=sum(x["total"] for x in tick_rows), pass_n=sum(x["pass_n"] for x in tick_rows), fail_n=sum(x["fail_n"] for x in tick_rows), unknown_n=sum(x["unknown_n"] for x in tick_rows), empty_minutes=empty_min, ) ) tick_primary = "kiwoom" ob_primary = "kiwoom" try: from kis_trader.engine.feed_fallback import live_ob_primary, live_tick_primary tick_primary = live_tick_primary() ob_primary = live_ob_primary() except Exception: try: from kis_trader.utils.env import get_env_from_db tick_primary = str( get_env_from_db("LIVE_TICK_PROVIDER", "kiwoom") or "kiwoom" ).strip().lower() ob_primary = str( get_env_from_db("LIVE_OB_PROVIDER", "kiwoom") or "kiwoom" ).strip().lower() except Exception: pass tick_chain = _vendor_read_chain(tick_primary) ob_chain = _vendor_read_chain(ob_primary) tick_wins = _count_tick_picks(db, day8, age, tick_chain, day_slice=ds) ob_wins = _count_ob_picks(db, day8, age, ob_chain, day_slice=ds) _annotate_share_and_pick( rows_out, tick_wins=tick_wins, ob_wins=ob_wins, tick_chain=tick_chain, ob_chain=ob_chain, ) tick_pick_sum = sum(tick_wins.values()) ob_pick_sum = sum(ob_wins.values()) return { "age_sec": age, "empty_minutes_all_fail": empty_min, "tick_chain": tick_chain, "ob_chain": ob_chain, "tick_picks": tick_wins, "ob_picks": ob_wins, "rows": rows_out, "note": ( f"나이={age:g}s (LIVE_FEED_FALLBACK). " "자사usable%=그 source 총행 대비(합격+미상)/총행 — 벤더 합≠100. " "usable분배%=채널 안 usable 건수 비중(합≈100). " f"간택%=종목×초마다 체인({'→'.join(tick_chain)}) 1순위 usable 벤더 " f"(틱슬롯={tick_pick_sum}, 호가슬롯={ob_pick_sum}). " f"전량탈락분={empty_min}. " "키움 raw=HHMMSS → tick_time(14)로 lag." ), } def _disconnect_thresholds() -> tuple: """수집통계 끊김: soft / hard / cap (초).""" soft, hard, cap = 10, 60, 1800 try: from kis_trader.utils.env import get_env_int soft = max(1, int(get_env_int("FEED_STATS_DISCONNECT_SOFT_SEC", 10) or 10)) hard = max(soft, int(get_env_int("FEED_STATS_DISCONNECT_HARD_SEC", 60) or 60)) cap = max(hard + 1, int(get_env_int("FEED_STATS_DISCONNECT_CAP_SEC", 1800) or 1800)) except Exception: pass return soft, hard, cap def _gap_row_from_sql(d: Dict[str, Any], *, soft: int, hard: int) -> Dict[str, Any]: soft_n = int(d.get("soft_n") or 0) hard_n = int(d.get("hard_n") or 0) pair_n = int(d.get("pair_n") or 0) max_gap = d.get("max_gap") avg_soft = d.get("avg_soft_gap") return { "vendor": str(d.get("vendor") or ""), "pair_n": pair_n, "soft_n": soft_n, "hard_n": hard_n, "soft_pct": _pct(soft_n, pair_n), "hard_pct": _pct(hard_n, pair_n), "max_gap_sec": int(max_gap) if max_gap is not None else None, "avg_soft_gap_sec": ( round(float(avg_soft), 1) if avg_soft is not None else None ), "soft_sec": soft, "hard_sec": hard, } def _feed_disconnect_stats( db, day8: str, day_slice: Optional[_FeedStatsDaySlice] = None, ) -> Dict[str, Any]: """종목별 recv_ts 간격으로 끊김(연결 공백) 집계 — usable(체결시각 lag)과 다른 축. 장중(09:00~15:30) · 연속 틱 LAG. cap 초과 공백은 미구독/장외로 제외. """ soft, hard, cap = _disconnect_thresholds() ds = day_slice or _FeedStatsDaySlice(db, day8) tick_tbl = ds.tick_from() tick_where, tick_params = ds.tick_filter_sql() rows: List[Dict[str, Any]] = [] by_vendor: Dict[str, Dict[str, Any]] = {} try: if ds.ready: tick_lag_where = ( f" FROM {tick_tbl} " "WHERE TIME(recv_ts) BETWEEN '09:00:00' AND '15:30:00'" ) tick_lag_params: tuple = () else: tick_lag_where = ( f" FROM {tick_tbl}{tick_where} " "AND TIME(recv_ts) BETWEEN '09:00:00' AND '15:30:00'" ) tick_lag_params = tick_params sql = ( "SELECT source AS vendor, " "COUNT(*) AS pair_n, " "SUM(CASE WHEN gap_sec >= %s THEN 1 ELSE 0 END) AS soft_n, " "SUM(CASE WHEN gap_sec >= %s THEN 1 ELSE 0 END) AS hard_n, " "MAX(gap_sec) AS max_gap, " "AVG(CASE WHEN gap_sec >= %s THEN gap_sec END) AS avg_soft_gap " "FROM (" " SELECT source, " " TIMESTAMPDIFF(SECOND, " " LAG(recv_ts) OVER (PARTITION BY source, code ORDER BY recv_ts, id), " " recv_ts) AS gap_sec " f"{tick_lag_where}" ") t " "WHERE gap_sec IS NOT NULL AND gap_sec > 0 AND gap_sec < %s " "GROUP BY source" ) for r in db.conn.execute( sql, (soft, hard, soft, *tick_lag_params, cap) ).fetchall() or []: d = dict(r) if not isinstance(r, dict) else r row = _gap_row_from_sql(d, soft=soft, hard=hard) v = row["vendor"] if v: by_vendor[v] = row except Exception: pass try: d0 = datetime.strptime(day8, "%Y%m%d") d1 = d0.replace(hour=23, minute=59, second=59) ls_sql = ( "SELECT 'ls' AS vendor, " "COUNT(*) AS pair_n, " "SUM(CASE WHEN gap_sec >= %s THEN 1 ELSE 0 END) AS soft_n, " "SUM(CASE WHEN gap_sec >= %s THEN 1 ELSE 0 END) AS hard_n, " "MAX(gap_sec) AS max_gap, " "AVG(CASE WHEN gap_sec >= %s THEN gap_sec END) AS avg_soft_gap " "FROM (" " SELECT TIMESTAMPDIFF(SECOND, " " LAG(ts) OVER (PARTITION BY code ORDER BY ts, id), " " ts) AS gap_sec " " FROM ls_ws_ticks " " WHERE ts >= %s AND ts <= %s " " AND TIME(ts) BETWEEN '09:00:00' AND '15:30:00'" ") t " "WHERE gap_sec IS NOT NULL AND gap_sec > 0 AND gap_sec < %s" ) r = db.conn.execute(ls_sql, (soft, hard, soft, d0, d1, cap)).fetchone() if r: d = dict(r) if not isinstance(r, dict) else r if int(d.get("pair_n") or 0) > 0: by_vendor["ls"] = _gap_row_from_sql(d, soft=soft, hard=hard) except Exception: pass for name in ("kis", "kiwoom", "ls"): if name in by_vendor: rows.append(by_vendor[name]) else: rows.append( { "vendor": name, "pair_n": 0, "soft_n": 0, "hard_n": 0, "soft_pct": None, "hard_pct": None, "max_gap_sec": None, "avg_soft_gap_sec": None, "soft_sec": soft, "hard_sec": hard, } ) # 상대 비교 메모 hard_map = {r["vendor"]: int(r.get("hard_n") or 0) for r in rows} worst = max(rows, key=lambda x: int(x.get("hard_n") or 0)) if rows else None note = ( f"축=종목별 recv_ts 공백(연결 끊김). usable(체결시각 lag)과 다름. " f"장중 09:00~15:30 · ≥{soft}s=소프트 · ≥{hard}s=하드 · >={cap}s 제외. " "키움 FID20 동결은 틱이 계속 오면 여기선 끊김으로 안 잡힘." ) if worst and int(worst.get("hard_n") or 0) > 0: note += ( f" 오늘 하드끊김 최다(건수)={worst.get('vendor')}" f"({worst.get('hard_n')}회)." ) def _hp(v: str) -> float: for r in rows: if r.get("vendor") == v and r.get("hard_pct") is not None: return float(r["hard_pct"]) return 0.0 ls_hp, kis_hp = _hp("ls"), _hp("kis") if ls_hp > 0 and kis_hp > 0 and ls_hp > kis_hp * 1.5: note += ( f" LS 하드%(={ls_hp:.1f})≫한투(={kis_hp:.1f}) " "→ 연결 끊김 비율은 LS가 높아 3차 둔 판단과 정합." ) elif ls_hp > 0 and kis_hp > 0 and kis_hp >= ls_hp: note += ( f" 하드% 한투({kis_hp:.1f})≥LS({ls_hp:.1f}) — " "건수·구독종목 수와 함께 볼 것." ) return { "soft_sec": soft, "hard_sec": hard, "cap_sec": cap, "rows": rows, "note": note, } def _safe_group(db, sql: str, params: tuple) -> List[Dict[str, Any]]: try: rows = db.conn.execute(sql, params).fetchall() or [] except Exception: return [] out: List[Dict[str, Any]] = [] for r in rows: if isinstance(r, dict): out.append( { "key": str(r.get("k") or r.get("source") or r.get("strategy") or ""), "n": int(r.get("n") or 0), "codes": int(r.get("c") or 0), } ) else: out.append({"key": str(r[0] or ""), "n": int(r[1] or 0), "codes": int(r[2] or 0)}) return out def _safe_nc(db, sql: str, params: tuple) -> Dict[str, Any]: try: r = db.conn.execute(sql, params).fetchone() except Exception: return {"n": 0, "codes": 0} return _row_nc(r) def _distinct_codes(db, sql: str, params: tuple) -> Set[str]: try: rows = db.conn.execute(sql, params).fetchall() or [] except Exception: return set() out: Set[str] = set() for r in rows: if isinstance(r, dict): c = str(r.get("code") or "").strip() else: c = str(r[0] or "").strip() if c: out.add(c) return out def _env_flags(db) -> Dict[str, Any]: """운영설정 스냅샷 (표시용). get_env_* 실패해도 탭은 동작.""" out: Dict[str, Any] = {} try: from kis_trader.utils.env import get_env_bool, get_env_float, get_env_from_db except Exception: return out for k, dflt in ( ("LIVE_TICK_PROVIDER", "kis"), ("LIVE_OB_PROVIDER", "kis"), ("LS_WS_ORDERBOOK_SAVE_MODE", "tick"), ("WS_ORDERBOOK_SAVE_MODE", "tick"), ): try: out[k] = str(get_env_from_db(k, dflt) or dflt) except Exception: out[k] = dflt for k, dflt in ( ("LS_WS_ORDERBOOK_SAVE", True), ("WS_ORDERBOOK_SAVE_KIS", True), ("WS_ORDERBOOK_SAVE_KIWOOM", True), ("WS_TICK_SAVE_KIS", True), ("WS_TICK_SAVE_KIWOOM", True), ("LS_WS_TICK_SAVE", False), ("LS_WS_CANDLE_SAVE", True), ("LS_WS_TICK_MIRROR_WS_TICKS", False), ): try: out[k] = bool(get_env_bool(k, dflt)) except Exception: out[k] = dflt for k, dflt in ( ("LIVE_FEED_FALLBACK_MAX_AGE_SEC", 3.0), ("WS_ORDERBOOK_TICK_MAX_AGE_SEC", 3.0), ("FEED_STATS_DISCONNECT_SOFT_SEC", 10.0), ("FEED_STATS_DISCONNECT_HARD_SEC", 60.0), ("FEED_STATS_DISCONNECT_CAP_SEC", 1800.0), ): try: out[k] = float(get_env_float(k, dflt) or dflt) except Exception: out[k] = float(dflt) return out def _vendor_coverage( *, vendor: str, tick: Dict[str, Any], ob: Dict[str, Any], tick_codes: Optional[Set[str]], ob_codes: Optional[Set[str]], max_age_sec: float, ) -> Dict[str, Any]: """같은 벤더 틱 vs 호가 공백·동기율. tick_codes/ob_codes 가 None 이면(빠른 조회) 종목교집합 지표는 null. """ tick_n = int(tick.get("n") or 0) tick_c = int(tick.get("codes") or (len(tick_codes) if tick_codes is not None else 0) or 0) ob_n = int(ob.get("n") or 0) ob_c = int(ob.get("codes") or (len(ob_codes) if ob_codes is not None else 0) or 0) if tick_codes is None or ob_codes is None: tick_no_ob_n: Optional[int] = None ob_no_tick_n: Optional[int] = None tick_no_ob_pct: Optional[float] = None else: tick_no_ob_n = len(tick_codes - ob_codes) ob_no_tick_n = len(ob_codes - tick_codes) tick_no_ob_pct = _pct(tick_no_ob_n, tick_c) # 호가틱동기 미적재 대리지표: 틱행 대비 호가행이 없는 비율 # (MODE=tick 에서 snap None / skip_ram / 큐드롭 포함) missing_ob_rows = max(0, tick_n - ob_n) return { "vendor": vendor, "tick_n": tick_n, "tick_codes": tick_c, "ob_n": ob_n, "ob_codes": ob_c, "tick_codes_no_ob": tick_no_ob_n, "tick_codes_no_ob_pct": tick_no_ob_pct, "ob_codes_no_tick": ob_no_tick_n, "ob_per_tick_row_pct": _pct(ob_n, tick_n), "missing_ob_row_pct": _pct(missing_ob_rows, tick_n), "avg_tick_per_code": round(tick_n / tick_c, 1) if tick_c else 0.0, "avg_ob_per_code": round(ob_n / ob_c, 1) if ob_c else 0.0, "orderbook_tick_max_age_sec": float(max_age_sec), "sample_tick_no_ob_codes": ( sorted(tick_codes - ob_codes)[:30] if tick_codes is not None and ob_codes is not None else [] ), } def _ob_per_tick_ratio(ob_n: int, tick_n: int) -> Optional[float]: """호가행/틱행 (소수). 틱 0이면 None.""" if tick_n is None or int(tick_n) <= 0: return None return round(float(ob_n) / float(tick_n), 3) def _pack_vendor(label: str, tick_n: int, tick_c: int, ob_n: int, ob_c: int) -> Dict[str, Any]: tn = int(tick_n or 0) on = int(ob_n or 0) ratio = _ob_per_tick_ratio(on, tn) return { "vendor": label, "tick_n": tn, "tick_codes": int(tick_c or 0), "ob_n": on, "ob_codes": int(ob_c or 0), "ob_per_tick": ratio, "ob_per_tick_pct": _pct(on, tn), } def build_feed_collect_stats( db, day: Optional[str] = None, *, heavy: bool = False, invalidate_cache: bool = False, ) -> Dict[str, Any]: """한 거래일(KST YYYYMMDD) 수집 요약 + 증권사 대비. 기본(heavy=False): GROUP BY COUNT + 폴백나이(STR_TO_DATE 전수, 보통 십수 초). heavy=True: 추가로 DISTINCT 종목교집합·filter_eval pass/reject. 종료: 결과는 5분 TTL 메모리 캐시에 저장 (웹탭 재조회 시 즉시 반환). invalidate_cache=True 이면 캐시 무시 및 강제 재계산. """ global _STATS_CACHE, _STATS_CACHE_TTL day8 = _ymd8(day or "") # 캐시 히트 확인 if not invalidate_cache: cached = _STATS_CACHE.get(day8) if cached: expire_ts, result = cached if _time.monotonic() < expire_ts: return result result = _build_feed_collect_stats_inner(db, day8, heavy=heavy) # 캐시 저장 (과거일의 경우 1시간, 당일은 5분) today8 = datetime.now().strftime("%Y%m%d") ttl = _STATS_CACHE_TTL if day8 >= today8 else 3600.0 _STATS_CACHE[day8] = (_time.monotonic() + ttl, result) # 낡은 캐시 항목 정리 _STATS_CACHE = { k: v for k, v in _STATS_CACHE.items() if _time.monotonic() < v[0] } return result def invalidate_feed_collect_stats_cache(day8: Optional[str] = None) -> None: """캐시 무효화. day8 지정 시 해당 일자만, None이면 전체.""" global _STATS_CACHE if day8: _STATS_CACHE.pop(day8, None) else: _STATS_CACHE.clear() def _build_feed_collect_stats_inner( db, day8: str, *, heavy: bool = False, ) -> Dict[str, Any]: """build_feed_collect_stats 실제 연산 바디 (캐시 래퍼 제외).""" t0, t1 = _day_bounds(day8) like = _like(day8) notes: List[str] = [] heavy = bool(heavy) day_slice = _FeedStatsDaySlice(db, day8) if not day_slice.ensure(): notes.append( "TEMP 일별 슬라이스 생성 실패 — tick_time range 직접 스캔으로 폴백." ) elif day_slice.tick_n > 0 or day_slice.ob_n > 0: notes.append( f"일별 TEMP 슬라이스: tick={day_slice.tick_n:,} ob={day_slice.ob_n:,} " "(lag·끊김·pick은 TEMP만 스캔)." ) tick_tbl = day_slice.tick_from() ob_tbl = day_slice.ob_from() tick_where, tick_params = day_slice.tick_filter_sql() ob_where, ob_params = day_slice.ob_filter_sql() ticks_by_src = _safe_group( db, "SELECT source AS k, COUNT(*) AS n, COUNT(DISTINCT code) AS c " f"FROM {tick_tbl}{tick_where} GROUP BY source ORDER BY n DESC", tick_params, ) ob_by_src = _safe_group( db, "SELECT source AS k, COUNT(*) AS n, COUNT(DISTINCT code) AS c " f"FROM {ob_tbl}{ob_where} GROUP BY source ORDER BY n DESC", ob_params, ) if day_slice.ready: fe_where = f" FROM {ob_tbl} WHERE source=%s" fe_params: tuple = ("filter_eval",) else: fe_where = f" FROM {ob_tbl}{ob_where} AND source=%s" fe_params = (*ob_params, "filter_eval") fe_by_strat = _safe_group( db, "SELECT COALESCE(strategy,'') AS k, COUNT(*) AS n, COUNT(DISTINCT code) AS c " f"{fe_where} " "GROUP BY strategy ORDER BY n DESC", fe_params, ) fe_reject = {"n": 0, "codes": 0} fe_pass = {"n": 0, "codes": 0} if heavy: fe_reject = _safe_nc( db, "SELECT COUNT(*) AS n, COUNT(DISTINCT code) AS c " f"{fe_where} AND reject_code IS NOT NULL AND reject_code<>''", fe_params, ) fe_pass = _safe_nc( db, "SELECT COUNT(*) AS n, COUNT(DISTINCT code) AS c " f"{fe_where} AND (reject_code IS NULL OR reject_code='')", fe_params, ) ls_ob = _safe_nc( db, "SELECT COUNT(*) AS n, COUNT(DISTINCT code) AS c FROM ls_ws_orderbook " "WHERE snap_time >= %s AND snap_time <= %s", (t0, t1), ) try: d0 = datetime.strptime(day8, "%Y%m%d") d1 = d0.replace(hour=23, minute=59, second=59) ls_ticks = _safe_nc( db, "SELECT COUNT(*) AS n, COUNT(DISTINCT code) AS c FROM ls_ws_ticks " "WHERE ts >= %s AND ts <= %s", (d0, d1), ) except Exception: ls_ticks = {"n": 0, "codes": 0} ls_candles = {"n": 0, "codes": 0} try: ls_candles = _safe_nc( db, "SELECT COUNT(*) AS n, COUNT(DISTINCT code) AS c FROM ls_ws_candles " "WHERE candle_time LIKE %s", (like,), ) except Exception: pass kis_dedicated = _safe_nc( db, "SELECT COUNT(*) AS n, COUNT(DISTINCT code) AS c FROM kis_ws_orderbook WHERE snap_time LIKE %s", (like,), ) def _pick(rows: List[Dict[str, Any]], *keys: str) -> Dict[str, Any]: want = {k.lower() for k in keys} n, c = 0, 0 for r in rows: if str(r.get("key") or "").lower() in want: n += int(r.get("n") or 0) c = max(c, int(r.get("codes") or 0)) return {"n": n, "codes": c} tick_kis = _pick(ticks_by_src, "kis") tick_kw = _pick(ticks_by_src, "kiwoom") tick_ls = _pick(ticks_by_src, "ls") ob_kis = _pick(ob_by_src, "kis_h0stasp0", "kis") ob_kw = _pick(ob_by_src, "kiwoom_0d", "kiwoom") ob_fe = _pick(ob_by_src, "filter_eval") env = _env_flags(db) max_age = float(env.get("WS_ORDERBOOK_TICK_MAX_AGE_SEC") or 3.0) fb_age = float(env.get("LIVE_FEED_FALLBACK_MAX_AGE_SEC") or 3.0) kis_tick_codes: Optional[Set[str]] = None kis_ob_codes: Optional[Set[str]] = None kw_tick_codes: Optional[Set[str]] = None kw_ob_codes: Optional[Set[str]] = None ls_tick_codes: Optional[Set[str]] = None ls_ob_codes: Optional[Set[str]] = None if heavy: if day_slice.ready: kis_tick_codes = _distinct_codes( db, f"SELECT DISTINCT code FROM {tick_tbl} WHERE source=%s", ("kis",), ) kis_ob_codes = _distinct_codes( db, f"SELECT DISTINCT code FROM {ob_tbl} WHERE source=%s", ("kis_h0stasp0",), ) kw_tick_codes = _distinct_codes( db, f"SELECT DISTINCT code FROM {tick_tbl} WHERE source=%s", ("kiwoom",), ) kw_ob_codes = _distinct_codes( db, f"SELECT DISTINCT code FROM {ob_tbl} WHERE source=%s", ("kiwoom_0d",), ) else: kis_tick_codes = _distinct_codes( db, f"SELECT DISTINCT code FROM {tick_tbl}{tick_where} AND source=%s", (*tick_params, "kis"), ) kis_ob_codes = _distinct_codes( db, f"SELECT DISTINCT code FROM {ob_tbl}{ob_where} AND source=%s", (*ob_params, "kis_h0stasp0"), ) kw_tick_codes = _distinct_codes( db, f"SELECT DISTINCT code FROM {tick_tbl}{tick_where} AND source=%s", (*tick_params, "kiwoom"), ) kw_ob_codes = _distinct_codes( db, f"SELECT DISTINCT code FROM {ob_tbl}{ob_where} AND source=%s", (*ob_params, "kiwoom_0d"), ) try: d0 = datetime.strptime(day8, "%Y%m%d") d1 = d0.replace(hour=23, minute=59, second=59) ls_tick_codes = _distinct_codes( db, "SELECT DISTINCT code FROM ls_ws_ticks WHERE ts >= %s AND ts <= %s", (d0, d1), ) except Exception: ls_tick_codes = set() ls_ob_codes = _distinct_codes( db, "SELECT DISTINCT code FROM ls_ws_orderbook WHERE snap_time >= %s AND snap_time <= %s", (t0, t1), ) coverage = [ _vendor_coverage( vendor="kis", tick=tick_kis, ob=ob_kis, tick_codes=kis_tick_codes, ob_codes=kis_ob_codes, max_age_sec=max_age, ), _vendor_coverage( vendor="kiwoom", tick=tick_kw, ob=ob_kw, tick_codes=kw_tick_codes, ob_codes=kw_ob_codes, max_age_sec=max_age, ), _vendor_coverage( vendor="ls", tick=ls_ticks, ob=ls_ob, tick_codes=ls_tick_codes, ob_codes=ls_ob_codes, max_age_sec=max_age, ), ] try: slip_rows = db.conn.execute( "SELECT avg_buy_price, target_price FROM active_trades WHERE buy_date >= %s", (f"{day8[:4]}-{day8[4:6]}-{day8[6:8]} 00:00:00",) ).fetchall() slips = [] for r in slip_rows: r = dict(r) bp = float(r.get("avg_buy_price") or 0) tp = float(r.get("target_price") or 0) if tp > 0 and bp > 0: slips.append((bp - tp) / tp * 100.0) avg_slippage = round(sum(slips) / len(slips), 3) if slips else None slippage_count = len(slips) except Exception as e: avg_slippage = None slippage_count = 0 vendor_matrix = { "day8": day8, "day": f"{day8[:4]}-{day8[4:6]}-{day8[6:8]}", "avg_slippage_pct": avg_slippage, "slippage_count": slippage_count, "vendors": [ _pack_vendor( "kis", int(tick_kis.get("n") or 0), int(tick_kis.get("codes") or 0), int(ob_kis.get("n") or 0), int(ob_kis.get("codes") or 0), ), _pack_vendor( "kiwoom", int(tick_kw.get("n") or 0), int(tick_kw.get("codes") or 0), int(ob_kw.get("n") or 0), int(ob_kw.get("codes") or 0), ), _pack_vendor( "ls", int(ls_ticks.get("n") or 0), int(ls_ticks.get("codes") or 0), int(ls_ob.get("n") or 0), int(ls_ob.get("codes") or 0), ), ], "filter_eval_n": int(ob_fe.get("n") or 0), "filter_eval_codes": int(ob_fe.get("codes") or 0), "kis_tick": int(tick_kis.get("n") or 0), "kis_ob": int(ob_kis.get("n") or 0), "kis_ob_tick": _ob_per_tick_ratio(int(ob_kis.get("n") or 0), int(tick_kis.get("n") or 0)), "kw_tick": int(tick_kw.get("n") or 0), "kw_ob": int(ob_kw.get("n") or 0), "kw_ob_tick": _ob_per_tick_ratio(int(ob_kw.get("n") or 0), int(tick_kw.get("n") or 0)), "ls_tick": int(ls_ticks.get("n") or 0), "ls_ob": int(ls_ob.get("n") or 0), "ls_ob_tick": _ob_per_tick_ratio(int(ls_ob.get("n") or 0), int(ls_ticks.get("n") or 0)), } if heavy and kis_tick_codes is not None and kw_tick_codes is not None: both_tick = len(kis_tick_codes & kw_tick_codes) overlap = { "both_tick_codes": both_tick, "kis_only_tick_codes": len(kis_tick_codes - kw_tick_codes), "kiwoom_only_tick_codes": len(kw_tick_codes - kis_tick_codes), "both_ob_codes": len((kis_ob_codes or set()) & (kw_ob_codes or set())), "kis_only_ob_codes": len((kis_ob_codes or set()) - (kw_ob_codes or set())), "kiwoom_only_ob_codes": len((kw_ob_codes or set()) - (kis_ob_codes or set())), } else: overlap = { "both_tick_codes": None, "kis_only_tick_codes": None, "kiwoom_only_tick_codes": None, "both_ob_codes": None, "kis_only_ob_codes": None, "kiwoom_only_ob_codes": None, "deferred": True, } notes.append( "틱·호가 DB는 벤더별 SAVE로 각각 적재 (2·3차 spill일 때만 쌓이는 구조 아님). " f"읽기 폴백나이={fb_age:g}s · 호가틱동기 max_age={max_age:g}s." ) if not heavy: notes.append( "기본 조회에 폴백나이(usable) 포함. " "호가0종목 교집합·filter_eval 합격/거부는 「상세 집계」." ) for cov in coverage: if int(cov.get("tick_codes") or 0) <= 0 and int(cov.get("ob_n") or 0) <= 0: continue no_ob = cov.get("tick_codes_no_ob") miss = cov.get("missing_ob_row_pct") if no_ob is None: if int(cov.get("tick_n") or 0) > 0 and miss is not None and float(miss) >= 5.0: notes.append( f"[{cov['vendor']}] 틱행 대비 호가행 공백≈{miss}% " f"(호가행/틱행={cov.get('ob_per_tick_row_pct')}%). " "호가0종목 수는 상세 집계에서 확인." ) continue no_ob_i = int(no_ob or 0) if int(cov.get("tick_n") or 0) > 0 and ( no_ob_i > 0 or (miss is not None and float(miss) >= 5.0) ): notes.append( f"[{cov['vendor']}] 틱종목 중 호가0종목 {no_ob_i}/{cov['tick_codes']}" f"({cov.get('tick_codes_no_ob_pct')}%) · " f"틱행 대비 호가행 공백≈{miss}% " f"(호가행/틱행={cov.get('ob_per_tick_row_pct')}%). " "원인후보: MODE=tick에서 해당벤더 호가RAM이 max_age 밖·미수신, " "또는 체결 lag>폴백나이로 호가틱동기 skip(틱DB는 기본 유지)." ) if cov.get("vendor") == "ls" and int(cov.get("ob_n") or 0) == 0 and int(cov.get("tick_n") or 0) == 0: notes.append("[ls] 당일 틱·호가 DB 0 — hold외·Bye·SAVE OFF·구독 실패 점검.") if int(ls_ob.get("n") or 0) == 0 and env.get("LS_WS_ORDERBOOK_SAVE"): mode = str(env.get("LS_WS_ORDERBOOK_SAVE_MODE") or "tick").lower() if mode in ("tick", "on_tick", "tick_sync", "sync"): notes.append( "LS 호가 SAVE=ON·MODE=tick 인데 당일 0건 → UH1만으로 DB 안 씀. " "체결 콜백 필요. Bye/구독 빈약이면 0." ) else: notes.append("LS 호가 SAVE=ON 인데 당일 0건 → Bye·구독·interval 스로틀 점검.") if not env.get("LS_WS_TICK_SAVE"): notes.append("LS_WS_TICK_SAVE=false (ls_ws_ticks 미적재). ON이면 구독전체 적재.") else: notes.append("LS_WS_TICK_SAVE=true → 구독 종목 US3를 ls_ws_ticks 에 적재 (영구만이 아님).") # 폴백나이: 기본 조회에도 포함 (실측 ~10초대, 요약테이블/크론 불필요) age_cut = _feed_age_cut_stats(db, day8, fb_age, day_slice=day_slice) disconnect = _feed_disconnect_stats(db, day8, day_slice=day_slice) return { "day": f"{day8[:4]}-{day8[4:6]}-{day8[6:8]}", "day8": day8, "heavy": heavy, "env": env, "notes": notes, "vendor_matrix": vendor_matrix, "recent_days": [], "recent_days_n": 0, "age_cut": age_cut, "disconnect": disconnect, "summary": { "ticks": { "kis": tick_kis, "kiwoom": tick_kw, "ls_mirror": tick_ls, "ls_table": ls_ticks, }, "orderbook": { "kis_h0stasp0": ob_kis, "kiwoom_0d": ob_kw, "filter_eval": ob_fe, "ls_ws_orderbook": ls_ob, "kis_ws_orderbook": kis_dedicated, }, "ls_candles": ls_candles, "filter_eval_pass": fe_pass, "filter_eval_reject": fe_reject, }, "coverage": coverage, "overlap": overlap, "ticks_by_source": ticks_by_src, "orderbook_by_source": ob_by_src, "filter_eval_by_strategy": fe_by_strat, }