Files
kis_bot/kis_trader/backtest/param_search_scalping.py
Hwang 61c72a8a4c feat(tests): 신규 키움 웹소켓 조건검색 및 실시간 조건검색 테스트 추가
변경 사항
----
- _test_kiwoom_condition_list.py: 키움 웹소켓 조건검색 '목록조회' 기능을 단독으로 테스트하는 스크립트 추가
- _test_kiwoom_condition_realtime.py: 'momentum' 조건식을 실시간으로 등록하고 초기 매칭 종목 리스트 및 실시간 편입/이탈을 수신하는 테스트 스크립트 추가
- _verify_columnar_bitid.py, _verify_shared_e2e_breakout.py, _verify_shared_e2e.py: 공유 메모리 및 dict 간의 데이터 일관성을 검증하는 테스트 추가

영향
----
- 신규 테스트 스크립트 추가로 키움 웹소켓 API의 기능 검증 및 안정성을 높임
- 기존 기능에 대한 영향 없음

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-07-06 01:27:00 +09:00

1008 lines
48 KiB
Python
Raw 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.
#!/usr/bin/env python3
"""
kis_trader/backtest/param_search_scalping.py — 스캘핑(Reversal) 백테스트 파라미터 자동 탐색 (Grid Search)
======================================================================================================
[전략 = SCALP / Reversal]
- RSI(3) 과매도(<rsi_oversold) → 양봉 전환 V자 반등 진입 ("바닥잡기" 단타).
- 실매매 ScalpingStrategy.check_buy → ``scalping_engine.check_buy_signal_live`` 호출.
- 본 CLI 는 ``scalping_engine.run_scalping_backtest`` 위에서 그리드 서치.
[관련 CLI (전략별 파일 분리, 2026-05 정리)]
SCALP(Reversal): kis_trader/backtest/param_search_scalping.py ← 이 파일
MOMENTUM : kis_trader/backtest/param_search_momentum.py
BREAKOUT : kis_trader/backtest/param_search_breakout.py
SHORT(꼬리잡기) : kis_trader/backtest/tail_param_search.py
실행:
cd /home/hoon/kis_bot
python3 kis_trader/backtest/param_search_scalping.py
# 또는 패키지 모듈로
python3 -m kis_trader.backtest.param_search_scalping --mode fast --apply 1
옵션:
--start 시작일 (기본: 오늘-7일)
--end 종료일 (기본: 오늘)
--mode 탐색 모드: fast(기본, ~192조합·~1215분) / coarse / fine / full / wide
fast = reversal·어깨·tp_max·손익 축만 (MACD/방어는 DB 고정)
--top 상위 N개 출력 (기본: 5000)
--min_trades 최소 거래 건수 필터 (기본: 1)
--apply 결과 N위 조합을 DB에 자동 적용 (기본: False, 지정하면 N 생략 시 1)
--apply-ai Gemini 로 수익·승률 기준 조합 하나 선택해서 DB 자동 적용
설명:
웹 API를 거치지 않고 scalping_engine 을 직접 호출하여 속도를 극대화했습니다.
피뢰침 방지(high_chase_thr)·급등주 필터(max_daily_chg)·방어필터 ON/OFF 도 그리드에 포함됩니다.
위치 이관 이력:
1) backtest_scalping/param_search.py
2) → kis_trader/backtest/param_search.py (2026-04)
3) → kis_trader/backtest/param_search_scalping.py (2026-05, 전략별 파일 분리)
- ROOT = kis_bot 프로젝트 루트 (__file__ 기준 3단계 위)
- 결과 저장: kis_trader/backtest/results/ (구 backtest_scalping/results 는 그대로 유지)
- scalping_engine / database.TradeDB 는 ROOT 에서 그대로 임포트
"""
import sys, os, json, time, argparse, signal
import heapq
from datetime import datetime, timedelta
from itertools import product
from concurrent.futures import as_completed
from typing import Optional, List, Dict, Any, Tuple
# kis_bot 루트를 경로에 추가 (스크립트·패키지 실행 모두 대응)
HERE = os.path.dirname(os.path.abspath(__file__))
ROOT = os.path.dirname(os.path.dirname(HERE))
if ROOT not in sys.path:
sys.path.insert(0, ROOT)
if HERE not in sys.path:
sys.path.insert(0, HERE)
import logging
logging.getLogger("TradeDB").setLevel(logging.WARNING) # 반복 초기화 로그 억제
from database import TradeDB
from kis_trader.engine import scalping_engine as se
from kis_trader.backtest import scalping_backtest_common as sbc
from kis_trader.backtest.backtest_portfolio_common import (
merge_param_search_apply_source,
portfolio_env_patch,
session_env_patch,
)
from kis_trader.backtest.param_search_cli_common import (
add_portfolio_cli_args,
add_search_filter_cli_args,
apply_session_to_fixed,
combo_passes_search_filters,
format_session_hm,
search_json_meta,
)
from kis_trader.backtest.param_search_pool import (
ParamSearchSharedPayload,
ParamSearchProgressETA,
assert_parent_alive,
cap_combos_uniform,
cap_combos_uniform_lazy,
iter_pool_chunk_results,
managed_process_pool,
param_search_chunk_plan,
param_search_worker_budget_line,
try_acquire_run_lock,
worker_shared_get,
)
from kis_trader.utils.env import get_env_float
# ──────────────────────────────────────────────────────────────────────────────
# 스캘핑 기본값 = 엔진에서 DB 로드 (백테스트 API와 동일 단일 소스, 실매매와 동기화)
# ──────────────────────────────────────────────────────────────────────────────
def _fixed_defaults():
"""엔진 get_scalping_defaults_from_db() 사용. UI/파람서치용으로 %/비율 변환."""
_d = se.get_scalping_defaults_from_db()
return {
"rsi_period": _d["rsi_period"],
"rsi_overbought": _d.get("rsi_overbought", 75.0), # 과열 차단 (그리드 탐색 제외, 고정값)
"slot_money": _d["slot_money"],
"vol_mult": _d["vol_mult"],
"trail_trigger": _d["trail_trigger"] * 100, # % 단위 (레거시, 청산 미사용)
"trail_stop": _d["trail_stop"] * 100, # % 단위 (레거시, 청산 미사용)
"shoulder_min_high": _d.get("shoulder_min_high", 0.005) * 100,
"shoulder_cut_pct": _d.get("shoulder_cut_pct", 0.003) * 100,
"tp_max_pct": _d.get("tp_max_pct", 0.02) * 100,
"cooldown_min": _d["cooldown_min"],
"time_start_hm": _d["time_start_hm"],
"time_end_hm": _d["time_end_hm"],
"max_daily": _d["max_daily"],
"fee_rate": _d["fee_rate"] * 100, # % 단위
"sell_tax": _d["sell_tax"] * 100, # % 단위
# 방어로직
"high_chase_thr": _d["high_chase_thr"], # 비율 (0.96)
"max_daily_chg": _d["max_daily_chg"], # % 단위
"min_price": _d["min_price"],
"max_loss_krw": _d["max_loss_krw"],
"min_margin": _d["min_margin"] * 100, # % 단위
"use_defense_filters": _d.get("use_defense_filters", True),
"use_macd_cross": _d.get("use_macd_cross", False),
}
# RSI_OVERSOLD별 JSON 저장 개수 (한 RSI에 치중되지 않도록 균등 분배)
PER_RSI_JSON = 20
# ──────────────────────────────────────────────────────────────────────────────
# 결과 디렉터리 (신규 위치 우선, 구 경로도 계속 조회 가능)
# ──────────────────────────────────────────────────────────────────────────────
def _results_dir_for_write() -> str:
"""새로 저장할 결과는 kis_trader/backtest/results/ 에 둔다."""
d = os.path.join(HERE, "results")
os.makedirs(d, exist_ok=True)
return d
def _results_dirs_for_read() -> List[str]:
"""읽기용 디렉터리: 신규 → 구 순서."""
return [
os.path.join(HERE, "results"),
os.path.join(ROOT, "backtest_scalping", "results"),
]
def _latest_json(prefix: str) -> Optional[str]:
"""읽기용 디렉터리에서 prefix 로 시작하는 가장 최근 JSON 경로 반환."""
best_path, best_mtime = None, -1.0
for d in _results_dirs_for_read():
if not os.path.isdir(d):
continue
for f in os.listdir(d):
if not (f.startswith(prefix) and f.endswith(".json")):
continue
p = os.path.join(d, f)
try:
m = os.path.getmtime(p)
except OSError:
continue
if m > best_mtime:
best_mtime = m
best_path = p
return best_path
# ──────────────────────────────────────────────────────────────────────────────
# 파라미터 그리드 (min_price 등은 DB 앵커 기반 — 고정 1값만 쓰지 않음)
# ──────────────────────────────────────────────────────────────────────────────
def _min_price_grid():
"""DB get_scalping_defaults 의 min_price 를 포함한 최소가격 스윕 (저·중·고가 후보)."""
mp = int(_fixed_defaults()["min_price"])
return sorted(set([500, 1000, max(500, mp), min(mp + 5000, 200000), 50000, 100000]))
def _scalp_grids():
"""탐색 모드별 그리드. 호출 시점 DB 기본값으로 min_price 티어가 잡힘."""
mp_list = _min_price_grid()
# fine: 조합 폭주 완화 — min_price 는 앵커·고가 두 티어 + 전 구간 5만원
fd_mp = int(_fixed_defaults()["min_price"])
min_price_fine = sorted(set([1000, fd_mp, 50000]))
return {
# ─────────────────────────────────────────────────────────────────────
# [FAST] ~1215분 — 어깨 1순위 + tp_max 상한 + reversal 진입 (192조합)
# MACD/방어/최소가격 등은 _fixed_defaults(DB) 고정. use_macd=False 권장 실매 정렬.
# 2×2×3×2×2×2×2 = 192
# ─────────────────────────────────────────────────────────────────────
"fast": {
"rsi_oversold": [21, 24],
"sl_pct": [1.2, 1.5],
"tp_pct": [2.0, 2.5, 3.0],
"tp_max_pct": [1.8, 2.0],
"shoulder_min_high": [0.3, 0.5],
"shoulder_cut_pct": [0.2, 0.3],
"drop_rate": [1.0, 1.5],
},
# coarse — 1차 스크리닝 (~576조합, history 유니버스 기준 약 1~1.5시간)
# MACD/reversal × 방어 × 손익 × RSI × 낙폭 × 트레일 (구 5,184 대비 축소)
# 조합: 4×3×3×2×2×2×2 = 576
"coarse": {
"rsi_oversold": [18, 21, 24, 28], # reversal 전용 (MACD ON 일 땐 무시)
"sl_pct": [0.8, 1.2, 1.5],
"tp_pct": [2.0, 2.5, 3.5],
"drop_rate": [1.0, 1.5],
"shoulder_min_high": [0.3, 0.5],
"shoulder_cut_pct": [0.2, 0.3],
"tp_max_pct": [1.8, 2.0],
"high_chase_thr": [0.96],
"max_daily_chg": [28.0],
"min_price": [1000],
"use_defense_filters": [True, False],
"use_macd_cross": [False, True],
"max_loss_krw": [200000],
"min_drop_pct_for_loss_cut": [0.015],
},
# 세밀 2차 (트레일·쿨다운 포함 + 최소가격·방어)
"fine": {
"rsi_oversold": [15, 17, 19, 21, 25],
"sl_pct": [0.8, 1.0, 1.2, 1.5],
"tp_pct": [1.5, 2.0, 2.5, 3.0, 3.5],
"drop_rate": [1.0, 1.5, 2.0],
"shoulder_min_high": [0.2, 0.3, 0.5],
"shoulder_cut_pct": [0.2, 0.3, 0.4],
"cooldown_min": [5, 10],
"high_chase_thr": [0.96, 0.98],
"max_daily_chg": [15.0, 20.0, 25.0],
"min_price": min_price_fine,
"use_defense_filters": [True, False],
"max_loss_krw": [150000, 200000, 300000],
"min_drop_pct_for_loss_cut": [0.01, 0.015, 0.02],
},
# 전체 탐색 (매우 오래 걸림)
"full": {
"rsi_oversold": [15, 17, 20, 25],
"sl_pct": [0.8, 1.0, 1.2, 1.5, 2.0],
"tp_pct": [1.5, 2.0, 2.5, 3.0, 4.0],
"drop_rate": [0.8, 1.0, 1.5, 2.0, 2.5],
"shoulder_min_high": [0.2, 0.3, 0.5, 0.7],
"shoulder_cut_pct": [0.2, 0.3, 0.4, 0.5],
"cooldown_min": [5, 10],
"high_chase_thr": [0.96, 0.98, 1.0],
"max_daily_chg": [15.0, 20.0, 30.0],
"min_price": mp_list,
"use_defense_filters": [True, False],
"max_loss_krw": [100000, 200000, 300000],
"min_drop_pct_for_loss_cut": [0.01, 0.015, 0.02, 0.025],
},
# 방어축 광범위 (RSI·슬립은 coarse 보다 약간 줄이고 급등·min_price·방어를 넓게)
"wide": {
"rsi_oversold": [15, 17, 19, 21, 23, 25, 28],
"sl_pct": [0.8, 1.0, 1.2, 1.5],
"tp_pct": [1.5, 2.0, 2.5, 3.0],
"drop_rate": [1.0, 1.5, 2.0, 2.5],
"high_chase_thr": [0.96, 0.98, 1.0],
"max_daily_chg": [8.0, 12.0, 16.0, 20.0, 25.0, 30.0, 40.0],
"min_price": mp_list,
"use_defense_filters": [True, False],
"max_loss_krw": [100000, 200000, 300000],
"min_drop_pct_for_loss_cut": [0.01, 0.015, 0.02],
},
}
def _get_scalp_field_map():
"""스캘핑 파라미터 → env_config 컬럼 매핑 (apply / db_snapshot 공용). 웹·봇과 동일 키 저장."""
return {
"rsi_oversold": ("SCALP_RSI_OVERSOLD", lambda v: str(int(v))),
"rsi_overbought": ("SCALP_RSI_OVERBOUGHT", lambda v: str(int(v))),
"sl_pct": ("SCALP_STOP_LOSS_PCT", lambda v: str(float(v) / 100)),
"tp_pct": ("SCALP_TAKE_PROFIT_PCT", lambda v: str(float(v) / 100)),
"drop_rate": ("SCALP_MIN_DROP_RATE", lambda v: str(float(v) / 100)),
"trail_trigger": ("SCALP_ATR_UP_MULT", lambda v: str(abs(float(v)) / 100.0)),
"trail_stop": ("SCALP_ATR_DOWN_MULT", lambda v: str(abs(float(v)) / 100.0)),
"shoulder_min_high": ("SCALP_SHOULDER_MIN_HIGH_PCT", lambda v: str(float(v) / 100)),
"shoulder_cut_pct": ("SCALP_SHOULDER_CUT_PCT", lambda v: str(float(v) / 100)),
"tp_max_pct": ("SCALP_TP_MAX_PCT", lambda v: str(float(v) / 100)),
"cooldown_min": ("SCALP_COOLDOWN_SEC", lambda v: str(int(float(v)) * 60)),
# 스캘핑 전용 방어로직 (꼬리잡기와 값 분리)
"high_chase_thr": ("SCALP_HIGH_PRICE_CHASE_THRESHOLD", lambda v: str(float(v))),
"max_daily_chg": ("SCALP_MAX_DAILY_CHANGE_PCT", lambda v: str(float(v))),
"min_price": ("SCALP_MIN_PRICE", lambda v: str(int(float(v)))),
"max_loss_krw": ("SCALP_MAX_LOSS_PER_TRADE_KRW", lambda v: str(int(float(v)))),
"min_drop_pct_for_loss_cut": ("SCALP_MIN_DROP_PCT_FOR_LOSS_CUT", lambda v: str(round(float(v)*100, 2)) if float(v) < 1 else str(round(float(v), 2))), # % (1.5)
"min_margin": ("SCALP_MIN_PROFIT_PCT", lambda v: str(float(v))), # % 단위 (0.2 등)
"use_defense_filters": ("SCALP_USE_DEFENSE_FILTERS", lambda v: "true" if v else "false"),
"use_macd_cross": ("SCALP_USE_MACD_CROSS", lambda v: "true" if v else "false"),
}
def _params_to_db_snapshot(params: dict) -> dict:
"""그리드 params(표시 단위) + 포트폴리오 → env_config 컬럼명:값 문자열 dict."""
field_map = _get_scalp_field_map()
snap = {
db_col: fmt(params[param_k])
for param_k, (db_col, fmt) in field_map.items()
if param_k in params
}
snap.update(portfolio_env_patch("SCALP", params))
snap.update(session_env_patch("SCALP", params))
return snap
def _apply_from_latest_json(rank: int):
"""최근 search_*.json에서 rank번째(1-based) 항목의 merged_params를 DB에 적용."""
latest_path = _latest_json("search_")
if not latest_path:
print("⚠️ search_*.json 파일이 없습니다 (신규/구 경로 모두).")
return
with open(latest_path, "r", encoding="utf-8") as f:
data = json.load(f)
top = data.get("top") or []
if rank < 1 or rank > len(top):
print(f"⚠️ 순번 {rank}이(가) 유효하지 않습니다. (1~{len(top)})")
return
item = top[rank - 1]
merged = merge_param_search_apply_source(item, data)
if not merged:
print("⚠️ 해당 항목에 merged_params/params가 없습니다.")
return
if item.get("total_pnl", 0) <= 0:
print(f"⚠️ {rank}번째 결과는 총손익 ≤ 0 (조건 미충족). DB 미적용. 기존 설정 유지.")
return
print(f"📂 {latest_path} 에서 {rank}번째 적용합니다.")
_apply_to_db(merged)
def _ui_to_engine_params(ui_params: dict) -> dict:
"""UI 표시용(% 등) → 엔진용 비율 단위. 워커에서 공통 사용."""
engine_params = dict(ui_params)
engine_params["sl_pct"] = ui_params["sl_pct"] / 100
engine_params["tp_pct"] = ui_params["tp_pct"] / 100
if "tp_max_pct" in ui_params:
engine_params["tp_max_pct"] = ui_params["tp_max_pct"] / 100
elif "tp_max_pct" not in engine_params:
engine_params["tp_max_pct"] = 0.02
engine_params["drop_rate"] = ui_params["drop_rate"] / 100
if "shoulder_min_high" in ui_params:
engine_params["shoulder_min_high"] = ui_params["shoulder_min_high"] / 100
if "shoulder_cut_pct" in ui_params:
engine_params["shoulder_cut_pct"] = ui_params["shoulder_cut_pct"] / 100
if "trail_trigger" in ui_params:
engine_params["trail_trigger"] = ui_params["trail_trigger"] / 100
if "trail_stop" in ui_params:
engine_params["trail_stop"] = ui_params["trail_stop"] / 100
engine_params["fee_rate"] = ui_params["fee_rate"] / 100
engine_params["sell_tax"] = ui_params["sell_tax"] / 100
engine_params["min_margin"] = ui_params.get("min_margin", 0.2) / 100
if "use_defense_filters" in ui_params:
engine_params["use_defense_filters"] = bool(ui_params["use_defense_filters"])
if "use_macd_cross" in ui_params:
engine_params["use_macd_cross"] = bool(ui_params["use_macd_cross"])
# 👇 [핵심] 손절 퍼센트에 맞춰 1회 투자금(slot_money) 자동 계산 (봇과 동일 공식)
# 5억 고정이면 수수료만으로 손절컷 걸려 좋은 조합(RSI 17 등)이 버려짐 → max_loss_krw/sl_pct 로 보정
max_loss = engine_params.get("max_loss_krw", 200000)
if engine_params["sl_pct"] > 0 and max_loss > 0:
engine_params["slot_money"] = max_loss / engine_params["sl_pct"]
return engine_params
def _evaluate_scalp_chunk(
param_chunk: List[Dict[str, Any]],
base_fixed: Dict[str, Any],
keys: List[str],
codes_candles: Optional[Dict[str, List[Dict]]],
min_trades: int,
min_win_rate: float,
min_pf: float,
top_n: int,
universe_by_slot: Optional[Dict[str, List[str]]] = None,
sort_by: str = "pnl",
slot_money: float = 3_000_000.0,
max_stocks: int = 3,
total_budget_krw: float = 9_000_000.0,
fee_rate: float = 0.00015,
sell_tax: float = 0.0018,
period_days: int = 1,
) -> List[Tuple[float, float, int, Dict]]:
"""워커: 청크 내 조합 평가 — 시각순 포트폴리오·총한도 (scalping_backtest_common)."""
shared = worker_shared_get()
if shared:
if codes_candles is None:
codes_candles = shared.get("codes_candles") or {}
if universe_by_slot is None:
universe_by_slot = shared.get("universe_by_slot")
if codes_candles is None:
codes_candles = {}
local_heap: List[Tuple[float, float, int, Dict]] = []
for combo in param_chunk:
assert_parent_alive()
ui_params = dict(base_fixed)
ui_params.update(combo)
engine_params = _ui_to_engine_params(ui_params)
engine_params["slot_money"] = float(slot_money)
engine_params["max_stocks"] = int(max_stocks)
engine_params["total_budget_krw"] = float(total_budget_krw)
engine_params["portfolio_mode"] = True
meta: Dict[str, Any] = {}
trades = sbc.run_scalping_backtest_web_aligned(
codes_candles, engine_params, universe_by_slot,
slot_money=slot_money, fee_rate=fee_rate, sell_tax=sell_tax,
max_stocks=max_stocks, total_budget_krw=total_budget_krw,
meta_out=meta, mode="reversal",
)
stats = sbc.summarize_scalp_trades(
trades, total_budget_krw=total_budget_krw, period_days=period_days,
)
total_trades = stats["total_trades"]
if total_trades < min_trades:
continue
total_pnl = stats["total_pnl"]
win_rate = stats["win_rate"]
pf = float(stats.get("pf") or 0)
if not combo_passes_search_filters(
win_rate=win_rate, pf=pf,
min_win_rate=min_win_rate, min_pf=min_pf,
):
continue
avg_hold = stats["avg_hold_min"]
peak, mdd, cum = 0.0, 0.0, 0.0
for t in trades:
cum += t["pnl"]
if cum > peak:
peak = cum
dd = peak - cum
if dd > mdd:
mdd = dd
merged = dict(ui_params)
merged["slot_money"] = float(slot_money)
merged["max_stocks"] = int(max_stocks)
merged["total_budget_krw"] = float(total_budget_krw)
result_pkg = {
"params": {k: ui_params[k] for k in keys},
"total_pnl": int(total_pnl),
"win_rate": round(win_rate, 2),
"total_trades": total_trades,
"pf": round(pf, 2),
"avg_hold": round(avg_hold, 1),
"mdd": round(mdd),
"bot_pct": stats["bot_pct"],
"daily_avg_pct": stats["daily_avg_pct"],
"skipped_micro_buys": int(
(meta.get("skip_stats") or {}).get("skipped_micro_buys") or 0
),
"merged_params": merged,
}
# 청크 내 상위 top_n: sort_by 에 따라 heapreplace (과거 min-heap+pushpop 로 최악 조합만 남는 버그 수정)
if str(sort_by).strip().lower() == "pnl":
item_t = (total_pnl, win_rate, id(result_pkg), result_pkg)
if len(local_heap) < top_n:
heapq.heappush(local_heap, item_t)
elif total_pnl > local_heap[0][0]:
heapq.heapreplace(local_heap, item_t)
else:
item_t = (win_rate, total_pnl, id(result_pkg), result_pkg)
if len(local_heap) < top_n:
heapq.heappush(local_heap, item_t)
elif win_rate > local_heap[0][0]:
heapq.heapreplace(local_heap, item_t)
return local_heap
def _load_candles_for_search(start: str, end: str, rsi_period: int) -> dict:
"""엔진 직접 타격을 위해 DB에서 캔들을 한 번만 메모리에 로드합니다."""
db = TradeDB()
codes_candles = {}
try:
start_key = (start.replace("-", "") + "0000") if start else "20260101"
end_key = (end.replace("-", "") + "2359") if end else "99991231"
codes_raw = db.conn.execute(
"SELECT DISTINCT code FROM ws_candles WHERE timeframe=1 "
"AND candle_time >= %s AND candle_time <= %s ORDER BY code",
[start_key, end_key]
).fetchall()
codes = [r["code"] for r in codes_raw]
for code in codes:
rows = db.conn.execute(
"SELECT candle_time, open, high, low, close, volume "
"FROM ws_candles "
"WHERE timeframe=1 AND code=%s "
"AND candle_time >= %s AND candle_time <= %s "
"AND is_confirmed=1 "
"ORDER BY candle_time ASC",
[code, start_key, end_key]
).fetchall()
if len(rows) < rsi_period + 5:
continue
codes_candles[code] = [dict(r) for r in rows]
finally:
db.close()
return codes_candles
def run_search(start: str, end: str, mode: str, top_n: int,
min_trades: int, min_win_rate: float, min_pf: float,
apply_rank: Optional[int], from_file_only: bool,
sort_by: str = "pnl", use_fallback_universe: bool = False,
slot_money: Optional[float] = None,
max_stocks: Optional[int] = None,
total_budget_krw: Optional[float] = None,
time_start_hm: Optional[int] = None,
time_end_hm: Optional[int] = None,
max_combos: Optional[int] = None) -> bool:
"""탐색 실행. 결과가 있어서 JSON 저장까지 했으면 True, 조건 만족 조합 없이 조기 return 시 False."""
if from_file_only and apply_rank is not None and apply_rank >= 1:
_apply_from_latest_json(apply_rank)
return True
grid = _scalp_grids()[mode]
keys = list(grid.keys())
axes = [grid[k] for k in keys]
# 데카르트곱 전체를 RAM 에 펼치지 않는다 (대형 그리드 OOM 방지).
dict_combos, total_grid, max_combos_cap, dropped_by_cap = cap_combos_uniform_lazy(
keys,
axes,
mode,
strategy_env_prefix="SCALP",
default_fast=192,
max_combos_override=max_combos,
)
total = len(dict_combos)
print(f"\n[{mode.upper()} 모드] 그리드: {total_grid:,} → 백테: {total:,} | 기간: {start} ~ {end}")
if mode == "fast":
print(f"📌 [fast] {total_grid:,}{max_combos_cap}균등샘플 · reversal·어깨·손익 축")
if dropped_by_cap:
print(f" (max-combos={max_combos_cap} 균등 샘플, 제외 {dropped_by_cap:,}개)")
print(f"📌 1위 정렬 기준: {'총손익 최대 (수익 나는 조합 우선)' if sort_by == 'pnl' else '승률 최대'}")
print(f"📌 필터: 승률≥{min_win_rate}% · PF≥{min_pf} · 거래≥{min_trades}")
print("=" * 70)
results = []
t0 = time.time()
FIXED_DEFAULTS = _fixed_defaults()
apply_session_to_fixed(
FIXED_DEFAULTS,
time_start_hm=time_start_hm,
time_end_hm=time_end_hm,
)
db = TradeDB()
try:
row = db.conn.execute("SELECT * FROM env_config ORDER BY id DESC LIMIT 1").fetchone()
env_row = dict(row) if row else {}
finally:
db.close()
fee_rate, sell_tax, slot_from_env = sbc.fee_and_slot_from_env(env_row, strategy="SCALP")
portfolio = sbc.resolve_scalp_portfolio_params(
env_row,
None,
strategy="SCALP",
slot_money=slot_money if slot_money is not None else slot_from_env,
max_stocks=max_stocks,
total_budget_krw=total_budget_krw,
)
slot_money_v = float(portfolio["slot_money"])
max_stocks_v = int(portfolio["max_stocks"])
total_budget_v = float(portfolio["total_budget_krw"])
period_days = max(
1,
(datetime.strptime(end, "%Y-%m-%d") - datetime.strptime(start, "%Y-%m-%d")).days + 1,
)
print(
f"💼 포트폴리오: 1회 {slot_money_v:,.0f}원 | 동시 {max_stocks_v}종 | "
f"총한도 {total_budget_v:,.0f}원 | 매매 {format_session_hm(FIXED_DEFAULTS)}"
)
if portfolio.get("budget_warning"):
print(f"💰 {portfolio['budget_warning']}")
# 캔들 데이터를 메모리에 1회 로드 (워커에 전달)
print("⏳ DB에서 캔들 데이터를 메모리로 불러오는 중...")
codes_candles = _load_candles_for_search(start, end, FIXED_DEFAULTS.get("rsi_period", 3))
print(f"✅ 데이터 로드 완료: {len(codes_candles)}종목")
# 유니버스: --fallback-universe 이면 이력 무시하고 시뮬레이션만 사용 (조합별 거래 수 확대)
# 신봇 기준:
# * 실매매는 10초 REST 폴링 + 변동 tick 마다 초단위 event_time 으로 저장.
# * 백테스트는 TradeDBExt.get_universe_by_candle_time("SCALP", ...) 로
# 1분 캔들 시각 키를 가진 dict 로 받아 엔진에 그대로 주입.
# * fallback 시뮬레이션은 1분봉 근사 점수 기반 5분 버킷팅 (과거 호환).
universe_by_slot = None
fallback_sim_interval = 5 # --fallback-universe 전용 시뮬레이션 버킷 (분)
start_ymd = start.replace("-", "") if start else ""
end_ymd = end.replace("-", "") if end else ""
if use_fallback_universe:
print("📌 [유니버스] --fallback-universe: 저장 이력 무시 → 시뮬레이션 유니버스 (조합별 거래 수 확대)")
elif start_ymd and end_ymd:
try:
# 신봇 이력 (초단위 event_time) → 1분 캔들 시각으로 리샘플링
from kis_trader.database.db_manager import get_db as _get_ext_db # type: ignore
_ext = _get_ext_db()
history = _ext.get_universe_by_candle_time(
strategy_id="SCALP",
start_ymd=start_ymd,
end_ymd=end_ymd,
)
if history:
universe_by_slot = history
n_bins = len(history)
avg = sum(len(v) for v in history.values()) / max(1, n_bins)
print(
f"✅ 유니버스: 신봇 이력 사용 (event_time → 1분 캔들 리샘플링) | "
f"{n_bins:,}분봉 · 평균 {avg:.1f}종목 (SCALP)"
)
print(
"📌 [유니버스] 이력만 쓰면 매수 기회 적어 거래 0~1건 나올 수 있음. "
"조합 많을 때는 --fallback-universe 권장."
)
else:
universe_by_slot = None
except Exception as _e:
logging.getLogger("param_search").debug(
"신봇 유니버스 이력 조회 스킵: %s", _e,
)
universe_by_slot = None
# 엔진이 쓸 슬롯 단위 결정:
# * 신봇 이력 경로 → 1분봉 키(=passthrough)
# * 시뮬 fallback → fallback_sim_interval 분 버킷
engine_scan_interval_min = 1
if universe_by_slot is None:
universe_top_n = int(os.environ.get("UPDATE_UNIVERSE_TOP_N", "20"))
universe_min_score = float(os.environ.get("UPDATE_UNIVERSE_MIN_SCORE", "4.0"))
universe_by_slot = se.build_universe_simulation(
codes_candles,
top_n=universe_top_n,
min_score=universe_min_score,
scan_interval_min=fallback_sim_interval,
)
engine_scan_interval_min = fallback_sim_interval
n_slots = len(universe_by_slot)
avg_per_slot = sum(len(c) for c in universe_by_slot.values()) / max(1, n_slots)
print(
f"✅ 유니버스: 시뮬레이션 사용 (이력 없음) | "
f"{fallback_sim_interval}분 슬롯 {n_slots}개 · 슬롯당 평균 {avg_per_slot:.1f}종목"
)
print("📌 [유니버스] 서치는 '시뮬레이션 유니버스' 기준입니다. (신봇 이력 없음)")
# 엔진에 슬롯 단위 주입 (engine._slot_key 가 이 값으로 캔들 시각 정규화)
FIXED_DEFAULTS["scan_interval_min"] = engine_scan_interval_min
# 조합을 딕셔너리 리스트로 변환 후 청크 분할 (dict_combos 는 fast 균등샘플 적용 완료)
shared = ParamSearchSharedPayload({
"codes_candles": codes_candles,
"universe_by_slot": universe_by_slot,
})
payload_bytes = shared.estimate_bytes()
n_cpu = os.cpu_count() or 4
_cpu_frac = get_env_float("PARAM_SEARCH_CPU_FRAC", 0.8)
max_workers, chunk_size, _ = param_search_chunk_plan(total, payload_bytes)
chunks = [dict_combos[i:i + chunk_size] for i in range(0, len(dict_combos), chunk_size)]
print(param_search_worker_budget_line(payload_bytes))
print(f"⚙️ 멀티프로세싱 시작 (코어: {n_cpu}, 워커: {max_workers}, CPU {_cpu_frac*100:.0f}%) | 청크: {len(chunks):,}개 (청크당 ~{chunk_size}조합)")
start_time = time.time()
global_heap: List[Tuple[float, float, Tuple[int, int], Dict]] = []
progress_eta = ParamSearchProgressETA(len(chunks), max_workers)
with managed_process_pool(max_workers, shared_payload=shared) as executor:
def _submit(chunk: List[Dict[str, Any]]):
return executor.submit(
_evaluate_scalp_chunk, chunk, FIXED_DEFAULTS, keys, None,
min_trades, min_win_rate, min_pf, top_n, None, sort_by,
slot_money_v, max_stocks_v, total_budget_v, fee_rate, sell_tax, period_days,
)
combos_per_chunk = chunk_size
print(f"⏳ 청크 처리 중… (청크당 최대 {combos_per_chunk:,}개 조합, 완료되는 대로 진행률·ETA 출력)")
processed = 0
use_carriage_return = sys.stdout.isatty()
for local_results in iter_pool_chunk_results(
executor, chunks, _submit, max_workers=max_workers,
):
processed += 1
for idx, item in enumerate(local_results):
if sort_by == "pnl":
pnl, wr, _, result_pkg = item
entry = (pnl, wr, (processed, idx), result_pkg)
if len(global_heap) < top_n:
heapq.heappush(global_heap, entry)
elif pnl > global_heap[0][0]:
heapq.heapreplace(global_heap, entry)
else:
wr, pnl, _, result_pkg = item
entry = (wr, pnl, (processed, idx), result_pkg)
if len(global_heap) < top_n:
heapq.heappush(global_heap, entry)
elif wr > global_heap[0][0]:
heapq.heapreplace(global_heap, entry)
# ETA — 워밍업 1파도 제외 누적 평균 (param_search_pool.ParamSearchProgressETA)
progress = (processed / len(chunks)) * 100
elapsed_so_far = time.time() - start_time
eta_str = ParamSearchProgressETA.format_sec(
progress_eta.remaining_sec(processed, elapsed_so_far),
)
elapsed_str = ParamSearchProgressETA.format_elapsed(elapsed_so_far)
eta_msg = f" | 경과: {elapsed_str} | 남은시간: {eta_str}"
line = f"⏳ 진행률: {progress:.1f}% ({processed:,}/{len(chunks):,} 청크 완료){eta_msg}"
if use_carriage_return:
print(f"\r{line}", end="", flush=True)
else:
print(line, flush=True)
if use_carriage_return:
print(flush=True) # 줄바꿈으로 진행률 줄 마무리
elapsed = time.time() - start_time
if not global_heap:
print("⚠️ 조건을 만족하는 조합이 없습니다. (min_trades를 낮추거나 기간을 늘려보세요.)")
print("📌 DB 미적용. 기존 설정 유지.")
return False
# 힙: heapreplace 로 상위 유지 → heappop 순은 정렬 보조용일 뿐, 최종은 아래 sort 로 확정
results = [heapq.heappop(global_heap)[3] for _ in range(len(global_heap))]
if sort_by == "win_rate":
results.sort(key=lambda r: (-r["win_rate"], -r["total_pnl"]))
else:
results.sort(key=lambda r: (-r["total_pnl"], -r["win_rate"]))
# RSI_OVERSOLD 기준으로 균등 분배 → JSON에 RSI별 PER_RSI_JSON개씩 (한 RSI에 치중 방지)
if "rsi_oversold" in keys and results:
rsi_vals = sorted(set(r["params"]["rsi_oversold"] for r in results))
by_rsi: Dict[Any, List[Dict]] = {}
for r in results:
v = r["params"]["rsi_oversold"]
if v not in by_rsi:
by_rsi[v] = []
if len(by_rsi[v]) < PER_RSI_JSON:
by_rsi[v].append(r)
results = []
for v in rsi_vals:
results.extend(by_rsi.get(v, []))
# 손익 순 유지, 동점이면 RSI 낮은 쪽(과매도 강함) 우선
results.sort(key=lambda r: (-r["total_pnl"], r["params"]["rsi_oversold"], -r["win_rate"]))
print(f"✅ RSI별 상위 {PER_RSI_JSON}개씩 보장 → {len(results)}건 (총 {len(rsi_vals)}가지 RSI)")
# 승률 min_win_rate 이상만 우선; 없으면 차악
filtered = [
r for r in results
if combo_passes_search_filters(
win_rate=float(r.get("win_rate") or 0),
pf=float(r.get("pf") or 0),
min_win_rate=min_win_rate,
min_pf=min_pf,
)
]
if filtered:
results = filtered
order_msg = "수익→승률 순" if sort_by == "pnl" else "승률→수익 순"
print(f"✅ 승률≥{min_win_rate}% · PF≥{min_pf} {len(results)}건 중 {order_msg} 상위")
else:
print(f"⚠️ 승률≥{min_win_rate}% · PF≥{min_pf} 없음 → 차악(상위) 적용")
# 손익 마이너스인 조합 제외 (수익 나는 것만 표시·저장)
profitable = [r for r in results if r["total_pnl"] > 0]
if profitable:
results = profitable
print(f"✅ 총손익 플러스만 사용: {len(results)}건 (손실 조합 제외)")
else:
print(f"⚠️ 수익 나는 조합 없음 → 손실 최소 순으로 표시")
print(f"\n완료: {elapsed:.1f}초 | 유효 결과: {len(results)}")
# ── 결과 출력 ──
hdr_keys = [k for k in keys]
col_w = max(len(k) for k in hdr_keys) + 2
order_label = "수익" if sort_by == "pnl" else "승률"
print(f"\n{'='*90}")
print(f" 🏆 {order_label} TOP {min(top_n, len(results))} (투자금 1,000,000원 기준)")
print(f"{'='*90}")
# 헤더
hdr = " ".join(f"{k:>{col_w}}" for k in hdr_keys)
print(f"{hdr} | {'손익(원)':>12} {'승률':>6} {'거래':>5} {'PF':>5} {'보유':>6}")
print("-" * (len(hdr) + 55))
for r in results[:top_n]:
p = r["params"]
row = " ".join(
(f"{str(p[k]):>{col_w}}" if isinstance(p[k], bool) else f"{p[k]:>{col_w}.4g}")
for k in hdr_keys
)
print(f"{row} | {r['total_pnl']:>+12,.0f} {r['win_rate']:>5.1f}% "
f"{r['total_trades']:>5} {r['pf']:>5.2f} {r['avg_hold']:>5.1f}")
best = results[0]
bp = best["params"]
print(f"""
╔══════════════════════════════════════════╗
║ 🏆 1위 최적 파라미터 ║
╠══════════════════════════════════════════╣""")
_disp_map = {
"rsi_oversold": "SCALP_RSI_OVERSOLD",
"rsi_overbought": "SCALP_RSI_OVERBOUGHT",
"sl_pct": "SCALP_STOP_LOSS_PCT (÷100)",
"tp_pct": "SCALP_TAKE_PROFIT_PCT(÷100)",
"drop_rate": "SCALP_MIN_DROP_RATE (÷100)",
"trail_trigger": "SCALP_ATR_UP_MULT (레거시)",
"trail_stop": "SCALP_ATR_DOWN_MULT(레거시)",
"shoulder_min_high": "SCALP_SHOULDER_MIN_HIGH_PCT (÷100)",
"shoulder_cut_pct": "SCALP_SHOULDER_CUT_PCT (÷100)",
"tp_max_pct": "SCALP_TP_MAX_PCT (÷100)",
"cooldown_min": "SCALP_COOLDOWN_SEC(÷60=분)",
"high_chase_thr": "HIGH_PRICE_CHASE_THRESHOLD",
"max_daily_chg": "MAX_DAILY_CHANGE_PCT",
"max_loss_krw": "MAX_LOSS_PER_TRADE_KRW(원)",
"min_price": "SCALP_MIN_PRICE(원)",
"use_defense_filters": "SCALP_USE_DEFENSE_FILTERS",
}
for k, v in bp.items():
label = _disp_map.get(k, k)
print(f"{label:<30s} : {v!s:>6}")
print(f"""╠══════════════════════════════════════════╣
║ 총 손익 : {best['total_pnl']:>+12,.0f} 원 ║
║ 승률 : {best['win_rate']:>6.1f}% ║
║ 총 거래 : {best['total_trades']:>5} 건 ║
║ Profit Factor : {best['pf']:>5.2f}
║ 평균 보유 : {best['avg_hold']:>5.1f} 분 ║
║ 최대 낙폭(MDD): {best['mdd']:>12,.0f} 원 ║
╚══════════════════════════════════════════╝""")
# JSON 저장용: RSI별 200개씩 모은 전체 리스트 (콘솔 TOP은 여전히 top_n만 출력)
top_list = []
for idx, r in enumerate(results):
p = r["params"]
merged = r.get("merged_params", p)
top_list.append({
"rank": idx + 1,
"params": p,
"merged_params": merged,
"db_snapshot": _params_to_db_snapshot(merged),
"total_pnl": r["total_pnl"],
"win_rate": r["win_rate"],
"total_trades": r["total_trades"],
"pf": r["pf"],
"avg_hold": r["avg_hold"],
"mdd": r["mdd"],
"bot_pct": r.get("bot_pct"),
"daily_avg_pct": r.get("daily_avg_pct"),
"skipped_micro_buys": r.get("skipped_micro_buys", 0),
})
# ── JSON 결과 저장 (kis_trader/backtest/results/) ──
# 권한 문제 대비: 쓰기 실패 시 사용자 홈(~/.kis_bot_search_results/)으로 폴백.
out_dir = _results_dir_for_write()
ts = datetime.now().strftime("%Y%m%d_%H%M%S")
out_path = os.path.join(out_dir, f"search_{mode}_{ts}.json")
payload = {
"mode": mode,
"start": start,
"end": end,
"min_win_rate": min_win_rate,
"min_pf": min_pf,
"top": top_list,
}
payload.update(search_json_meta(portfolio, FIXED_DEFAULTS))
try:
with open(out_path, "w", encoding="utf-8") as f:
json.dump(payload, f, ensure_ascii=False, indent=2)
print(f"\n💾 결과 저장: {out_path}")
except (PermissionError, OSError) as _e:
fallback_dir = os.path.join(os.path.expanduser("~"), ".kis_bot_search_results")
os.makedirs(fallback_dir, exist_ok=True)
out_path = os.path.join(fallback_dir, f"search_{mode}_{ts}.json")
with open(out_path, "w", encoding="utf-8") as f:
json.dump(payload, f, ensure_ascii=False, indent=2)
print(f"\n⚠️ 기본 경로({out_dir}) 쓰기 실패({type(_e).__name__}). 폴백 저장: {out_path}")
print(f" 권한 복구: sudo chown -R $USER:$USER {out_dir}")
# ── DB 자동 적용 (만족 조건: 총손익 > 0. 미충족 시 기존 설정 유지) ──
if apply_rank is not None and apply_rank >= 1 and apply_rank <= len(top_list):
cand = top_list[apply_rank - 1]
if cand.get("total_pnl", 0) <= 0:
print(f"⚠️ {apply_rank}번째 결과는 총손익 ≤ 0 (조건 미충족). DB 미적용. 기존 설정 유지.")
else:
merged_apply = merge_param_search_apply_source(cand, payload)
_apply_to_db(merged_apply)
print(f"{apply_rank}번째 결과 적용 완료")
return True
def _apply_to_db(best_params: dict):
"""1위 파라미터 → insert_env_snapshot (config_scalp + env_config)."""
from kis_trader.backtest.param_search_apply_snapshot import apply_env_patch
patch = _params_to_db_snapshot(best_params)
if not patch:
print("DB 적용할 파라미터가 없습니다.")
return
eid = apply_env_patch(patch)
if eid is None:
print("❌ insert_env_snapshot 실패")
return
print(f"\n✅ config_scalp + env_config INSERT id={eid}:")
for k, v in sorted(patch.items()):
print(f" {k:<30s} = {v}")
# ──────────────────────────────────────────────────────────────────────────────
# CLI 진입점
# ──────────────────────────────────────────────────────────────────────────────
def main():
today = datetime.now().strftime("%Y-%m-%d")
week_ago = (datetime.now() - timedelta(days=7)).strftime("%Y-%m-%d")
parser = argparse.ArgumentParser(description="스캘핑 백테스트 파라미터 Grid Search")
parser.add_argument("--start", default=week_ago, help="시작일 (YYYY-MM-DD)")
parser.add_argument("--end", default=today, help="종료일 (YYYY-MM-DD)")
parser.add_argument("--mode", default="fast", choices=["fast", "coarse", "fine", "full", "wide"],
help="탐색 모드: fast(그리드→균등192·~1215분) / coarse / fine / full / wide")
parser.add_argument(
"--max-combos", type=int, default=None, dest="max_combos",
help="백테 조합 상한 (fast 기본 env SCALP_FAST_MAX_COMBOS 또는 PARAM_SEARCH_FAST_MAX_COMBOS=192, 0=무제한)",
)
parser.add_argument("--top", default=5000, type=int, help="상위 N개 출력·JSON 저장 (기본 5000)")
parser.add_argument("--min_trades", default=1, type=int, help="최소 거래 건수")
add_search_filter_cli_args(parser)
parser.add_argument("--apply", nargs="?", const=1, type=int, default=None, metavar="N",
help="N번째 결과를 DB에 적용 (기본 1). --from-file 시 최근 JSON에서 적용")
parser.add_argument("--from-file", action="store_true", help="--apply N 과 함께 사용 시, 최근 결과 JSON에서만 적용 (탐색 생략)")
parser.add_argument("--sort-by", default="pnl", choices=["pnl", "win_rate"],
help="1위 기준: pnl=총손익 최대(기본), win_rate=승률 최대")
parser.add_argument("--fallback-universe", action="store_true", dest="fallback_universe",
help="저장 이력 무시, 시뮬레이션 유니버스만 사용. 조합 많을 때 거래 수 확대용")
add_portfolio_cli_args(parser)
parser.add_argument("--apply-ai", action="store_true", dest="apply_ai",
help="Gemini가 수익·승률 기준으로 하나 골라 DB 적용. --from-file 과 함께 쓰면 탐색 없이 최근 JSON만 사용; 그 외에는 탐색 완료 후 방금 생성된 JSON으로 적용")
args = parser.parse_args()
# 탐색 없이 최근 JSON으로만 AI 적용 (--from-file --apply-ai)
if args.apply_ai and args.from_file:
import param_apply_ai
param_apply_ai.apply_ai_scalp()
return
# SIGTERM 을 KeyboardInterrupt 와 동일하게 처리 (systemd·운영자 kill 대응)
def _sigterm_to_kbd(_sig, _frm):
raise KeyboardInterrupt("SIGTERM 수신 → 워커 정리 후 종료")
try:
signal.signal(signal.SIGTERM, _sigterm_to_kbd)
except Exception:
pass
run_lock = None
if not (args.from_file and args.apply is not None):
run_lock = try_acquire_run_lock("param_search_scalping")
if run_lock is None:
print(
"⛔ 이미 실행 중인 param_search_scalping 이 있습니다.\n"
" ps -ef | grep param_search_scalping\n"
" pkill -f 'param_search_scalping.py' 후 재실행하세요.",
flush=True,
)
sys.exit(2)
try:
had_results = run_search(
start = args.start,
end = args.end,
mode = args.mode,
top_n = args.top,
min_trades = args.min_trades,
min_win_rate = args.min_win_rate,
min_pf = args.min_pf,
apply_rank = args.apply,
from_file_only = args.from_file,
sort_by = args.sort_by,
use_fallback_universe = args.fallback_universe,
slot_money = args.slot_money,
max_stocks = args.max_stocks,
total_budget_krw = args.total_budget,
time_start_hm = args.time_start,
time_end_hm = args.time_end,
max_combos = args.max_combos,
)
except KeyboardInterrupt as e:
print(f"\n{e} — 미완료 결과 없이 종료합니다.", flush=True)
sys.exit(130)
finally:
if run_lock is not None:
run_lock.release()
# 탐색에서 조건 만족 조합이 있었을 때만 방금 저장된 JSON으로 AI 적용 (없으면 기존 설정 유지)
if args.apply_ai and had_results:
import param_apply_ai
param_apply_ai.apply_ai_scalp()
elif args.apply_ai and not had_results:
print("📌 이번 탐색에서 조건 만족 조합 없음 → apply_ai 스킵. DB 미적용. 기존 설정 유지.")
if __name__ == "__main__":
main()