feat(ws): 키움 WS 시세 마이그레이션 검증 인프라
KIS WS(41 한도) → 키움 WS(100 한도) 본격 전환 전, 두 소스를 동시 운영해 가격 일치성을 데이터로 검증하기 위한 인프라. 신규 - kiwoom_ws.py: 키움 WS 클라이언트 (KIS WS 와 동일 get_price 인터페이스, 메모리 dict 캐시) - kis_trader/network/ws_validator.py: 5초마다 KIS↔키움 가격 비교, ws_price_validation 테이블에 1행 INSERT, |diff|≥WARN_PCT 시 WARN 로그 - test_kiwoom_ws.py: 키움 WS 단독 동작 확인 스크립트 (토큰/LOGIN/REG/시세) 수정 - database.py: ws_price_validation 테이블 + ENV 키 (WS_PROVIDER, WS_VALIDATION_INTERVAL_SEC, WS_VALIDATION_DIFF_WARN_PCT) + insert_ws_price_validation/get_ws_validation_stats 헬퍼 - kis_trader/main.py: WS_PROVIDER=kis_with_validation 시 키움 WS + Validator 백그라운드 기동, 종료 시 정리 운영 영향: 0. 매매·시세 의사결정은 항상 KIS WS만 사용. 키움 WS는 백그라운드 비교 기록만. 적용: DB env_config 에 WS_PROVIDER=kis_with_validation INSERT 후 봇 재시작. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
213
kis_trader/network/ws_validator.py
Normal file
213
kis_trader/network/ws_validator.py
Normal file
@@ -0,0 +1,213 @@
|
||||
"""
|
||||
kis_trader/network/ws_validator.py — KIS↔키움 WS 가격 검증
|
||||
============================================================
|
||||
주기적으로 같은 종목의 KIS WS 캐시와 키움 WS 캐시를 비교해
|
||||
``ws_price_validation`` 테이블에 1행 INSERT.
|
||||
|
||||
운영(매매)에는 영향 없음 — **읽기만** 함.
|
||||
검증 모드(``WS_PROVIDER=kis_with_validation``)에서만 기동.
|
||||
|
||||
데이터 흐름::
|
||||
|
||||
[KIS WS] ─┐
|
||||
│ 각각 메모리 dict 캐시
|
||||
[키움 WS] ─┘
|
||||
│
|
||||
▼
|
||||
[Validator] ── 5s 주기 ──► ws_price_validation 테이블
|
||||
│ │
|
||||
│ └─→ |diff| ≥ WARN_PCT 면 WARN 로그
|
||||
│
|
||||
└─→ 24h/1주 통계 분석 → 마이그레이션 전환 결정 근거
|
||||
|
||||
비교 대상 종목
|
||||
--------------
|
||||
KIS WS 가 구독 중인 종목 = 봇이 매매에 쓰는 가격이 실제로 들어오고 있는 종목.
|
||||
키움 WS 도 같은 종목을 구독하도록 동기화 (subscribe).
|
||||
|
||||
env_config 토글
|
||||
---------------
|
||||
``WS_VALIDATION_INTERVAL_SEC`` (기본 5)
|
||||
``WS_VALIDATION_DIFF_WARN_PCT`` (기본 0.10 — 0.1%p)
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import threading
|
||||
import time
|
||||
from typing import Optional, Set
|
||||
|
||||
from ..utils.env import get_env_float, get_env_int
|
||||
from ..utils.logger import get_logger
|
||||
|
||||
logger = get_logger("kis_trader.ws_validator")
|
||||
|
||||
|
||||
class WSPriceValidator:
|
||||
"""KIS↔키움 WS 가격 비교 백그라운드 워커.
|
||||
|
||||
매 N초마다:
|
||||
1) KIS WS 의 구독 중인 종목 목록을 읽음
|
||||
2) 키움 WS 가 같은 종목을 구독하도록 동기화
|
||||
3) 두 캐시에서 가격 조회 → DB INSERT
|
||||
4) |diff_pct| ≥ warn_pct 시 WARN 로그
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
kis_ws, # kis_ws.KISWebSocketPriceCache 인스턴스
|
||||
kiwoom_ws, # kiwoom_ws.KiwoomWebSocketPriceCache 인스턴스
|
||||
db, # database.TradeDB
|
||||
):
|
||||
self.kis_ws = kis_ws
|
||||
self.kiwoom_ws = kiwoom_ws
|
||||
self.db = db
|
||||
|
||||
self._thread: Optional[threading.Thread] = None
|
||||
self._running = False
|
||||
self._last_warn_ts: dict = {} # code → 최근 WARN 로그 시각 (스팸 방지)
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
def start(self) -> bool:
|
||||
if self._thread and self._thread.is_alive():
|
||||
return True
|
||||
self._running = True
|
||||
self._thread = threading.Thread(
|
||||
target=self._loop, daemon=True, name="WSPriceValidator",
|
||||
)
|
||||
self._thread.start()
|
||||
logger.info(
|
||||
"✅ WS 가격 검증기 시작 — KIS↔키움 비교 (interval=%ds, warn≥%.2f%%)",
|
||||
self._interval_sec(), self._warn_pct(),
|
||||
)
|
||||
return True
|
||||
|
||||
def stop(self) -> None:
|
||||
self._running = False
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
def _interval_sec(self) -> int:
|
||||
return max(1, get_env_int("WS_VALIDATION_INTERVAL_SEC", 5))
|
||||
|
||||
def _warn_pct(self) -> float:
|
||||
return max(0.0, get_env_float("WS_VALIDATION_DIFF_WARN_PCT", 0.10))
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
def _kis_subscribed(self) -> Set[str]:
|
||||
"""KIS WS 가 현재 구독 중인 종목 set."""
|
||||
# KISWebSocketPriceCache 의 _subscribed 직접 참조 (kis_ws.py 정의)
|
||||
try:
|
||||
with self.kis_ws._sub_lock: # type: ignore[attr-defined]
|
||||
return set(self.kis_ws._subscribed) # type: ignore[attr-defined]
|
||||
except AttributeError:
|
||||
return set()
|
||||
|
||||
def _sync_kiwoom_subscriptions(self, target: Set[str]) -> None:
|
||||
"""키움 WS 구독 = KIS WS 구독 으로 맞춤."""
|
||||
try:
|
||||
current = set()
|
||||
with self.kiwoom_ws._sub_lock: # type: ignore[attr-defined]
|
||||
current = set(self.kiwoom_ws._subscribed) # type: ignore[attr-defined]
|
||||
|
||||
to_add = target - current
|
||||
to_remove = current - target
|
||||
for code in to_add:
|
||||
self.kiwoom_ws.subscribe(code)
|
||||
for code in to_remove:
|
||||
self.kiwoom_ws.unsubscribe(code)
|
||||
except Exception as e:
|
||||
logger.debug("키움 구독 동기화 실패: %s", e)
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
def _loop(self) -> None:
|
||||
# 시작 직후 KIS/키움 둘 다 캐시 채워질 시간 약간 줌
|
||||
time.sleep(15)
|
||||
while self._running:
|
||||
try:
|
||||
self._tick()
|
||||
except Exception as e:
|
||||
logger.warning("WS 검증기 tick 예외: %s", e)
|
||||
time.sleep(self._interval_sec())
|
||||
|
||||
def _tick(self) -> None:
|
||||
"""1회 비교."""
|
||||
codes = self._kis_subscribed()
|
||||
if not codes:
|
||||
return
|
||||
|
||||
# 키움 WS 구독 동기화 (구독 안 된 종목은 데이터 없으니 의미 없음)
|
||||
self._sync_kiwoom_subscriptions(codes)
|
||||
|
||||
# 키움 WS 연결되어 있어야 의미 있음. 미연결이면 한쪽만 기록.
|
||||
kiwoom_ready = self.kiwoom_ws.is_connected()
|
||||
|
||||
warn_pct = self._warn_pct()
|
||||
warn_count = 0
|
||||
sample_count = 0
|
||||
|
||||
for code in codes:
|
||||
kis_data = self.kis_ws.get_price(code, max_age_sec=10.0)
|
||||
kw_data = self.kiwoom_ws.get_price(code, max_age_sec=10.0) if kiwoom_ready else None
|
||||
|
||||
kis_price = self._parse_price(kis_data, "stck_prpr")
|
||||
kw_price = self._parse_price(kw_data, "stck_prpr")
|
||||
|
||||
kis_age = self._parse_age_ms(kis_data)
|
||||
kw_age = self._parse_age_ms(kw_data)
|
||||
|
||||
if kis_price is None and kw_price is None:
|
||||
continue
|
||||
|
||||
# DB 1행 INSERT
|
||||
self.db.insert_ws_price_validation(
|
||||
code=code,
|
||||
kis_price=kis_price,
|
||||
kiwoom_price=kw_price,
|
||||
kis_age_ms=kis_age,
|
||||
kiwoom_age_ms=kw_age,
|
||||
)
|
||||
sample_count += 1
|
||||
|
||||
# 차이 경고
|
||||
if kis_price not in (None, 0) and kw_price is not None:
|
||||
diff_pct = (kw_price - kis_price) / kis_price * 100.0
|
||||
if abs(diff_pct) >= warn_pct:
|
||||
# 종목당 60초 1회만 로그 (스팸 방지)
|
||||
now = time.time()
|
||||
last = self._last_warn_ts.get(code, 0)
|
||||
if now - last >= 60:
|
||||
self._last_warn_ts[code] = now
|
||||
logger.warning(
|
||||
"⚠️ [WS 검증] %s 가격 차이 %.3f%% (KIS=%.0f, 키움=%.0f)",
|
||||
code, diff_pct, kis_price, kw_price,
|
||||
)
|
||||
warn_count += 1
|
||||
|
||||
if sample_count > 0:
|
||||
logger.debug(
|
||||
"📊 [WS 검증] tick: %d종목 비교 (warn=%d, 키움연결=%s)",
|
||||
sample_count, warn_count, kiwoom_ready,
|
||||
)
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
@staticmethod
|
||||
def _parse_price(data: Optional[dict], key: str) -> Optional[float]:
|
||||
if not data:
|
||||
return None
|
||||
try:
|
||||
v = float(str(data.get(key, "0")).replace(",", ""))
|
||||
return abs(v) if v else None
|
||||
except (ValueError, TypeError):
|
||||
return None
|
||||
|
||||
@staticmethod
|
||||
def _parse_age_ms(data: Optional[dict]) -> Optional[int]:
|
||||
if not data:
|
||||
return None
|
||||
v = data.get("_age_ms")
|
||||
if v is None:
|
||||
return None
|
||||
try:
|
||||
return int(v)
|
||||
except (ValueError, TypeError):
|
||||
return None
|
||||
Reference in New Issue
Block a user