From ab693d332c46267cbc3d5b4808e603a19f9f657f Mon Sep 17 00:00:00 2001 From: I Luk Kim Date: Sat, 30 May 2026 12:30:03 -0700 Subject: [PATCH] =?UTF-8?q?refactor:=20=EB=B0=B1=ED=95=84=20=EC=8A=A4?= =?UTF-8?q?=ED=81=AC=EB=A6=BD=ED=8A=B8=20=E2=86=92=20REST=20admin=20?= =?UTF-8?q?=EC=97=94=EB=93=9C=ED=8F=AC=EC=9D=B8=ED=8A=B8=EB=A1=9C=20?= =?UTF-8?q?=EC=A0=84=ED=99=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 직접 DB/서비스에 접근하는 스크립트 패턴을 REST API 패턴으로 교체: - 삭제: scripts/backfill_finra_2020_gap.py (기존 POST /finra/admin/ingest으로 충분) - 삭제: scripts/backfill_alpaca_daily_pit.py - 삭제: scripts/export_pit_panel.py (외부 DB 연결으로 대체) - 추가: POST /api/v1/alpaca/admin/backfill-pit BackgroundTasks 패턴, FINRA 전체 PIT 심볼 백필, adjustment=all Co-Authored-By: Claude Opus 4.8 --- app/api/v1/endpoints/alpaca.py | 189 +++++++++++++++- docs/DATA_COVERAGE.md | 13 +- scripts/backfill_alpaca_daily_pit.py | 313 --------------------------- scripts/backfill_finra_2020_gap.py | 83 ------- scripts/export_pit_panel.py | 257 ---------------------- 5 files changed, 198 insertions(+), 657 deletions(-) delete mode 100644 scripts/backfill_alpaca_daily_pit.py delete mode 100644 scripts/backfill_finra_2020_gap.py delete mode 100644 scripts/export_pit_panel.py diff --git a/app/api/v1/endpoints/alpaca.py b/app/api/v1/endpoints/alpaca.py index d193d0c..6e99e76 100644 --- a/app/api/v1/endpoints/alpaca.py +++ b/app/api/v1/endpoints/alpaca.py @@ -4,11 +4,17 @@ Alpaca Market Data endpoints — standalone price data via Alpaca API import asyncio import gc +import logging from datetime import date, datetime, timezone, timedelta from typing import Optional from zoneinfo import ZoneInfo -from fastapi import APIRouter, HTTPException, Query +from fastapi import APIRouter, BackgroundTasks, HTTPException, Query +from sqlalchemy import text + +from app.core.database import AsyncSessionLocal + +logger = logging.getLogger(__name__) _ET = ZoneInfo("America/New_York") _MARKET_CLOSE_HOUR = 16 # 4:00 PM ET @@ -325,6 +331,187 @@ def _parse_snapshot(ticker: str, raw: dict) -> AlpacaSnapshotResponse: ) +# ------------------------------------------------------------------ +# Admin: PIT price backfill +# ------------------------------------------------------------------ + +def _finra_to_alpaca(symbol: str) -> str: + """FINRA "/" → Alpaca "." (BRK/B → BRK.B, AAC/U → AAC.U).""" + return symbol.replace("/", ".").replace("-", ".") + + +async def _run_pit_backfill(start_str: str, end_str: str, force: bool) -> None: + """Background task: backfill Alpaca 1d bars for all FINRA PIT symbols.""" + from sqlalchemy.dialects.postgresql import insert as pg_insert + from app.models.alpaca_price import AlpacaPriceData + + BATCH_SIZE = 100 + CHUNK_SIZE = 2300 # asyncpg 32767 bind-param limit + ADJUSTMENT = "all" + TOLERANCE_DAYS = 5 + + client = AlpacaClient() + if not client.is_configured(): + logger.error("PIT backfill: Alpaca keys not configured") + return + + try: + # Step 1: FINRA PIT 심볼 + 활동기간 + async with AsyncSessionLocal() as db: + rows = (await db.execute(text(""" + SELECT symbol, + MIN(date)::date AS finra_first, + MAX(date)::date AS finra_last, + COUNT(DISTINCT date::date) AS finra_days + FROM finra_short_volume + GROUP BY symbol + """))).fetchall() + finra_info = {r.symbol: (r.finra_first, r.finra_last, r.finra_days) for r in rows} + total_obs = sum(v[2] for v in finra_info.values()) + logger.info(f"PIT backfill: {len(finra_info):,} FINRA symbols") + + # Step 2: 기존 Alpaca 1d 커버리지 (정규화 키로 저장) + async with AsyncSessionLocal() as db: + cov_rows = (await db.execute(text(""" + SELECT ticker, MAX(date)::date AS alpaca_max + FROM alpaca_price_data WHERE interval = '1d' + GROUP BY ticker + """))).fetchall() + alpaca_cov: dict = {} + for r in cov_rows: + key = _finra_to_alpaca(r.ticker).upper() + if key not in alpaca_cov or r.alpaca_max > alpaca_cov[key]: + alpaca_cov[key] = r.alpaca_max + + # Step 3: 수집 필요 심볼 결정 + if force: + need_fetch = list(finra_info.keys()) + else: + need_fetch = [ + sym for sym, (_, finra_last, _) in finra_info.items() + if (am := alpaca_cov.get(_finra_to_alpaca(sym).upper())) is None + or (finra_last - am).days > TOLERANCE_DAYS + ] + logger.info(f"PIT backfill: fetching {len(need_fetch):,} / {len(finra_info):,} symbols") + + # Step 4: 배치 수집 + upsert + total_inserted = 0 + n_batches = (len(need_fetch) + BATCH_SIZE - 1) // BATCH_SIZE + + for batch_idx in range(0, len(need_fetch), BATCH_SIZE): + batch = need_fetch[batch_idx: batch_idx + BATCH_SIZE] + batch_num = batch_idx // BATCH_SIZE + 1 + if batch_num == 1 or batch_num % 50 == 0: + logger.info(f"PIT backfill batch {batch_num}/{n_batches} | inserted={total_inserted:,}") + + reverse_map = {_finra_to_alpaca(s).upper(): s for s in batch} + try: + raw = await client.get_multi_bars( + symbols=batch, timeframe="1d", + start=start_str, end=end_str, + adjustment=ADJUSTMENT, + ) + except Exception as exc: + logger.warning(f"PIT backfill batch {batch_num} error: {exc}") + continue + + rows_to_insert = [] + for alpaca_sym, bar_list in raw.items(): + if not bar_list: + continue + original = reverse_map.get(alpaca_sym.upper(), alpaca_sym) + for bar in bar_list: + dt = datetime.fromisoformat(bar["t"].replace("Z", "+00:00")) + if dt.tzinfo is None: + dt = dt.replace(tzinfo=timezone.utc) + rows_to_insert.append({ + "ticker": original, + "date": dt, + "interval": "1d", + "open": float(bar.get("o") or 0), + "high": float(bar.get("h") or 0), + "low": float(bar.get("l") or 0), + "close": float(bar.get("c") or 0), + "volume": float(bar.get("v") or 0), + "vwap": float(bar["vw"]) if bar.get("vw") else None, + "trade_count": int(bar["n"]) if bar.get("n") else None, + "data_source": "ALPACA", + }) + + if rows_to_insert: + async with AsyncSessionLocal() as db: + for i in range(0, len(rows_to_insert), CHUNK_SIZE): + stmt = pg_insert(AlpacaPriceData).values( + rows_to_insert[i: i + CHUNK_SIZE] + ) + stmt = stmt.on_conflict_do_nothing(constraint="uq_alpaca_price_data") + result = await db.execute(stmt) + total_inserted += result.rowcount + await asyncio.sleep(0) + await db.commit() + + # Step 5: 완료 리포트 + async with AsyncSessionLocal() as db: + new_cov = { + _finra_to_alpaca(r.ticker).upper() + for r in (await db.execute(text( + "SELECT DISTINCT ticker FROM alpaca_price_data WHERE interval='1d'" + ))).fetchall() + } + covered_obs = sum(fd for sym, (_, _, fd) in finra_info.items() if _finra_to_alpaca(sym).upper() in new_cov) + missing_obs = total_obs - covered_obs + bias_pct = 100.0 * missing_obs / total_obs if total_obs else 0 + logger.info( + f"PIT backfill complete — inserted={total_inserted:,} | " + f"covered={len(new_cov):,} symbols | " + f"residual bias={bias_pct:.2f}% (row-weighted)" + ) + finally: + await client.close() + + +@router.post( + "/admin/backfill-pit", + summary="PIT 가격 백필 — 상폐 종목 포함 전체 FINRA 심볼", + description=( + "FINRA short-volume DB에 등장한 모든 심볼(현 활성 유니버스 + 상폐/합병 과거 심볼)의 " + "Alpaca SIP 일봉(1d)을 백필합니다. 생존편향-0 수익 계산에 필요.\n\n" + "**특성**:\n" + "- `adjustment=all` (분할+배당 조정) — 상폐 종목은 future-proof\n" + "- DB-first, idempotent (`on_conflict_do_nothing`) — 재실행 안전\n" + "- Alpaca SIP 일봉은 무료 플랜에서 2016-01-04부터 제공\n" + "- 기본 시작일: 2018-08-01 (FINRA DB 시작일)\n\n" + "**백그라운드 실행**: 즉시 `started` 응답, 1-3시간 소요.\n" + "진행 상황: `GET /alpaca/status` 또는 DB `SELECT COUNT(DISTINCT ticker) FROM alpaca_price_data WHERE interval='1d';`\n\n" + "**PIT 뷰** (백필 후): `SELECT DISTINCT symbol FROM pit_universe_membership WHERE d='2023-03-09';`" + ), +) +async def backfill_pit_prices( + background_tasks: BackgroundTasks, + start_date: date = Query(date(2018, 8, 1), description="백필 시작일 (기본: 2018-08-01)"), + force: bool = Query(False, description="이미 커버된 심볼도 재수집"), +): + svc = _require_alpaca() # API 키 확인 + _ = svc # 키 확인용 + + end_date = datetime.now(timezone.utc).date() - timedelta(days=1) + start_str = start_date.isoformat() + end_str = end_date.isoformat() + + background_tasks.add_task(_run_pit_backfill, start_str, end_str, force) + + return { + "status": "started", + "start_date": start_str, + "end_date": end_str, + "adjustment": "all", + "note": ( + "Backfilling Alpaca 1d bars for all FINRA PIT symbols in background. " + "Typically 1-3 hours for ~22k symbols. Check container logs for progress." + ), + } + + @router.get( "/snapshot", response_model=AlpacaMultiSnapshotResponse, diff --git a/docs/DATA_COVERAGE.md b/docs/DATA_COVERAGE.md index d79c125..e0fbd4b 100644 --- a/docs/DATA_COVERAGE.md +++ b/docs/DATA_COVERAGE.md @@ -685,9 +685,16 @@ docker exec stock_oracle_api python scripts/backfill_finra_2020_gap.py # adjustment='all' (분할+배당 조정), 2018-08-01부터 docker exec stock_oracle_api python scripts/backfill_alpaca_daily_pit.py -# Step 3: 리서치 레이어용 parquet 익스포트 (pandas + pyarrow 필요) -docker exec stock_oracle_api python scripts/export_pit_panel.py -# 출력: ./data/pit_panel.parquet +# Step 3 (옵션): 리서치 레이어용 parquet export +# 외부 환경에서 직접 DB 연결 (port 15433) 후 pandas로 export: +# python -c " +# import pandas as pd +# from sqlalchemy import create_engine +# eng = create_engine('postgresql+psycopg2://stockoracle:stockoracle2024@localhost:15433/stock_oracle') +# df = pd.read_sql('SELECT p.d, p.symbol, p.short_ratio, a.close FROM pit_universe_membership p LEFT JOIN alpaca_price_data a ON a.ticker=p.symbol AND a.date::date=p.d AND a.interval=\'1d\' WHERE p.d BETWEEN \'2018-08-01\' AND NOW()', eng) +# df.to_parquet('pit_panel.parquet', index=False) +# print(df.shape) +# " ``` ### PIT 뷰 (DB 직접 쿼리 시) diff --git a/scripts/backfill_alpaca_daily_pit.py b/scripts/backfill_alpaca_daily_pit.py deleted file mode 100644 index dcf44ab..0000000 --- a/scripts/backfill_alpaca_daily_pit.py +++ /dev/null @@ -1,313 +0,0 @@ -""" -PIT(Point-in-Time) Alpaca 일봉 가격 백필 — 생존편향-0 가격 데이터셋. - -FINRA short-volume DB에 등장한 모든 심볼(현 활성 유니버스 9,635 + 상폐/합병 -과거 심볼 ~16,143 포함, 계 22,722)의 일봉(1d)을 Alpaca SIP에서 수집한다. - -이 스크립트 없이는 가격(returns) 측이 현 생존자(~3k) 종목으로만 계산 가능해 -숏볼륨 신호 검정 자체가 생존편향으로 무효화된다. - -전략: - 1. FINRA 심볼 리스트 + 활동기간(min/max date) 로드. - 2. alpaca_price_data 1d 현재 커버리지 로드. - 3. 커버리지 불충분 심볼만 Alpaca 요청(DB-first, idempotent). - 4. adjustment='all'(분할+배당 조정) 저장 — 상폐 종목은 future-proof. - 5. 완료 후: covered/missing 심볼 수 + row-weighted 잔존편향 % 리포트. - -심볼 정규화: - - FINRA "/" (BRK/B, BF/B, AAC/U 등) → Alpaca "." (BRK.B, BF.B, AAC.U) - - 워런트/유닛(/U, /WS 등)은 Alpaca에 데이터 없음 → 빈 응답 (무해) - - 일반 심볼(대다수): 그대로 전달 - -환경 요건: - - ALPACA_API_KEY, ALPACA_SECRET_KEY 환경변수 설정 필요 - - Alpaca 일봉 SIP 데이터는 무료 플랜에서 2016-01-04까지 제공 - - 2018-08-01 기점 전 구간 PIT 달성 가능 (실증 검증 완료) - -소요 시간: - - 신규 심볼 ~19k × 8년치 / 200req/min ≈ 1-3시간 - - 이미 수집된 심볼은 스킵 (on_conflict_do_nothing, idempotent) - -Usage (컨테이너 내부): - python scripts/backfill_alpaca_daily_pit.py - -Usage (호스트): - docker exec stock_oracle_api python scripts/backfill_alpaca_daily_pit.py -""" - -import asyncio -import fcntl -import logging -import os -import sys -from datetime import datetime, timezone, timedelta, date - -sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) - -logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s") -logger = logging.getLogger("pit_price_backfill") - -_LOCK_FILE = "/tmp/backfill_alpaca_pit.lock" - -# 2018-08-01: earliest date where all current FINRA data starts -# (Alpaca 일봉은 2016-01-04부터 가능, 2018-08 기점 전 구간 커버) -_BACKFILL_START_STR = "2018-08-01" - -_BATCH_SIZE = 100 # 심볼/요청 (Alpaca multi-bar 한도 내) -_CHUNK_SIZE = 2300 # 행/upsert (asyncpg 32767 bind-param 한도: 2300×14=32,200) -_ADJUSTMENT = "all" # split+dividend 조정 (상폐 종목은 future-proof, 생존종목도 정확) - -# Coverage tolerance: alpaca_max vs finra_last 이 N일 이상 차이나면 re-fetch -_COVERAGE_TOLERANCE_DAYS = 5 - - -def finra_to_alpaca(symbol: str) -> str: - """FINRA 심볼을 Alpaca 심볼 형식으로 변환. - - FINRA는 주식 클래스와 워런트/유닛에 '/'를 사용 (BRK/B → BRK.B). - Yahoo Finance 호환 하이픈도 처리 (BRK-B → BRK.B). - """ - return symbol.replace("/", ".").replace("-", ".") - - -async def run_backfill(): - from sqlalchemy import text - from sqlalchemy.dialects.postgresql import insert as pg_insert - - from app.core.database import AsyncSessionLocal - from app.models.alpaca_price import AlpacaPriceData - from app.services.alpaca_client import AlpacaClient - - client = AlpacaClient() - if not client.is_configured(): - logger.error("ALPACA_API_KEY / ALPACA_SECRET_KEY 미설정. 중단.") - return - - now_utc = datetime.now(timezone.utc) - # 당일 봉은 아직 확정 전일 수 있으므로 1일 버퍼 - end_str = (now_utc - timedelta(days=1)).strftime("%Y-%m-%d") - start_str = _BACKFILL_START_STR - - logger.info(f"PIT 가격 백필 시작: {start_str} → {end_str}") - - # ------------------------------------------------------------------ # - # Step 1: FINRA 심볼 + 활동기간 로드 # - # ------------------------------------------------------------------ # - logger.info("FINRA PIT 심볼 리스트 로드 중 ...") - async with AsyncSessionLocal() as db: - rows = await db.execute( - text(""" - SELECT symbol, - MIN(date)::date AS finra_first, - MAX(date)::date AS finra_last, - COUNT(DISTINCT date::date) AS finra_days - FROM finra_short_volume - GROUP BY symbol - ORDER BY symbol - """) - ) - finra_info = {r.symbol: (r.finra_first, r.finra_last, r.finra_days) for r in rows} - - total_syms = len(finra_info) - total_obs = sum(v[2] for v in finra_info.values()) - logger.info(f"FINRA 심볼 수: {total_syms:,} | 총 관측치(심볼-일): {total_obs:,}") - - # ------------------------------------------------------------------ # - # Step 2: Alpaca 1d 현재 커버리지 로드 # - # ------------------------------------------------------------------ # - logger.info("기존 Alpaca 1d 커버리지 로드 중 ...") - async with AsyncSessionLocal() as db: - rows = await db.execute( - text(""" - SELECT ticker, MAX(date)::date AS alpaca_max - FROM alpaca_price_data - WHERE interval = '1d' - GROUP BY ticker - """) - ) - # 정규화 키로 저장 (BF-B, BF.B 등 혼재 → 점 형식으로 통일) - alpaca_coverage: dict[str, date] = {} - for r in rows: - key = finra_to_alpaca(r.ticker).upper() - # 동일 정규화 키가 여러 ticker로 존재할 수 있음 (BF-B, BF.B) → 최신값 사용 - if key not in alpaca_coverage or r.alpaca_max > alpaca_coverage[key]: - alpaca_coverage[key] = r.alpaca_max - - logger.info(f"Alpaca 1d 기존 커버리지 심볼 수: {len(alpaca_coverage):,}") - - # ------------------------------------------------------------------ # - # Step 3: 수집 필요 심볼 결정 # - # ------------------------------------------------------------------ # - # 조건: alpaca 데이터 없음, 또는 alpaca_max < finra_last - tolerance - need_fetch: list[str] = [] - for sym, (finra_first, finra_last, _) in finra_info.items(): - alpaca_key = finra_to_alpaca(sym).upper() - alpaca_max = alpaca_coverage.get(alpaca_key) - if alpaca_max is None: - need_fetch.append(sym) - elif (finra_last - alpaca_max).days > _COVERAGE_TOLERANCE_DAYS: - need_fetch.append(sym) - - logger.info( - f"수집 필요 심볼: {len(need_fetch):,} / {total_syms:,} " - f"(기존 충분: {total_syms - len(need_fetch):,})" - ) - - if not need_fetch: - logger.info("모든 심볼 커버리지 충분. 수집 불필요.") - else: - # ------------------------------------------------------------------ # - # Step 4: 배치 수집 + upsert # - # ------------------------------------------------------------------ # - total_inserted = 0 - total_bars_fetched = 0 - n_batches = (len(need_fetch) + _BATCH_SIZE - 1) // _BATCH_SIZE - - for batch_idx in range(0, len(need_fetch), _BATCH_SIZE): - batch = need_fetch[batch_idx: batch_idx + _BATCH_SIZE] - batch_num = batch_idx // _BATCH_SIZE + 1 - - if batch_num == 1 or batch_num % 20 == 0: - logger.info( - f"배치 {batch_num}/{n_batches}: {len(batch)}개 심볼 " - f"({batch[0]}..{batch[-1]}) | 누적 삽입: {total_inserted:,}" - ) - - # FINRA 심볼 → Alpaca 정규화 역방향 맵 (Alpaca 응답 키 → FINRA 원본 심볼) - reverse_map: dict[str, str] = { - finra_to_alpaca(s).upper(): s for s in batch - } - - try: - raw = await client.get_multi_bars( - symbols=batch, - timeframe="1d", - start=start_str, - end=end_str, - adjustment=_ADJUSTMENT, - ) - except Exception as exc: - logger.warning(f"배치 {batch_num} 요청 실패: {exc} — 스킵") - continue - - rows_to_insert = [] - for alpaca_sym, bar_list in raw.items(): - if not bar_list: - continue - # Alpaca 응답 키는 정규화된 형식 (예: BRK.B) — 역맵으로 FINRA 원본 복원 - original = reverse_map.get(alpaca_sym.upper(), alpaca_sym) - for bar in bar_list: - ts = bar["t"] - dt = datetime.fromisoformat(ts.replace("Z", "+00:00")) - if dt.tzinfo is None: - dt = dt.replace(tzinfo=timezone.utc) - rows_to_insert.append({ - "ticker": original, - "date": dt, - "interval": "1d", - "open": float(bar.get("o") or 0), - "high": float(bar.get("h") or 0), - "low": float(bar.get("l") or 0), - "close": float(bar.get("c") or 0), - "volume": float(bar.get("v") or 0), - "vwap": float(bar["vw"]) if bar.get("vw") else None, - "trade_count": int(bar["n"]) if bar.get("n") else None, - "data_source": "ALPACA", - }) - - total_bars_fetched += len(rows_to_insert) - - if rows_to_insert: - async with AsyncSessionLocal() as db: - for i in range(0, len(rows_to_insert), _CHUNK_SIZE): - stmt = pg_insert(AlpacaPriceData).values( - rows_to_insert[i: i + _CHUNK_SIZE] - ) - stmt = stmt.on_conflict_do_nothing( - constraint="uq_alpaca_price_data" - ) - result = await db.execute(stmt) - total_inserted += result.rowcount - await asyncio.sleep(0) # event loop yield - await db.commit() - - logger.info( - f"수집 완료 — 총 봉: {total_bars_fetched:,} | DB 삽입: {total_inserted:,} " - f"(중복 스킵: {total_bars_fetched - total_inserted:,})" - ) - - # ------------------------------------------------------------------ # - # Step 5: 완료 후 커버리지 리포트 # - # ------------------------------------------------------------------ # - logger.info("최종 커버리지 리포트 계산 중 ...") - async with AsyncSessionLocal() as db: - rows = await db.execute( - text(""" - SELECT ticker, - COUNT(DISTINCT date::date) AS alpaca_days - FROM alpaca_price_data - WHERE interval = '1d' - GROUP BY ticker - """) - ) - new_coverage_keys: set[str] = { - finra_to_alpaca(r.ticker).upper() for r in rows - } - - covered_syms = 0 - covered_obs = 0 - missing_syms_list: list[str] = [] - missing_obs = 0 - - for sym, (finra_first, finra_last, finra_days) in finra_info.items(): - key = finra_to_alpaca(sym).upper() - if key in new_coverage_keys: - covered_syms += 1 - covered_obs += finra_days - else: - missing_syms_list.append(sym) - missing_obs += finra_days - - bias_pct = 100.0 * missing_obs / total_obs if total_obs else 0.0 - - logger.info("=" * 70) - logger.info("PIT 가격 백필 완료 — 커버리지 리포트") - logger.info(f" FINRA 전체 심볼 : {total_syms:,}") - logger.info(f" Alpaca 가격 있는 심볼 : {covered_syms:,}") - logger.info(f" 가격 없는 심볼 (워런트/OTC 등): {len(missing_syms_list):,}") - logger.info(f" FINRA 전체 관측치(심볼-일) : {total_obs:,}") - logger.info(f" 가격 매칭 관측치 : {covered_obs:,}") - logger.info(f" 가격 누락 관측치 : {missing_obs:,}") - logger.info(f" 잔존 편향 (row-weighted) : {bias_pct:.2f}%") - logger.info("=" * 70) - - if missing_syms_list: - sample = sorted(missing_syms_list)[:30] - logger.info(f"가격 누락 심볼 샘플 (상위 30): {sample}") - logger.info("(대부분은 워런트/유닛 (/U, /WS) 또는 초저유동 OTC — 분석단에서 필터)") - - logger.info("") - logger.info("검증 쿼리:") - logger.info(" SELECT ticker, MIN(date)::date, MAX(date)::date, COUNT(*)") - logger.info(" FROM alpaca_price_data WHERE interval='1d'") - logger.info(" AND ticker IN ('SIVB','SBNY','FRC','TWTR','ATVI','DISCA')") - logger.info(" GROUP BY 1;") - - -if __name__ == "__main__": - lock_fh = open(_LOCK_FILE, "w") - try: - fcntl.flock(lock_fh, fcntl.LOCK_EX | fcntl.LOCK_NB) - except OSError: - logger.error("백필이 이미 실행 중 (lock 파일 보유). 종료.") - sys.exit(1) - - try: - asyncio.run(run_backfill()) - finally: - fcntl.flock(lock_fh, fcntl.LOCK_UN) - lock_fh.close() - try: - os.unlink(_LOCK_FILE) - except OSError: - pass diff --git a/scripts/backfill_finra_2020_gap.py b/scripts/backfill_finra_2020_gap.py deleted file mode 100644 index 9582183..0000000 --- a/scripts/backfill_finra_2020_gap.py +++ /dev/null @@ -1,83 +0,0 @@ -""" -FINRA short-volume 2020 갭 메우기 (2020-04-01 ~ 2020-10-31). - -DB 조사 결과: 전 구간(2018-08 ~ 2026-05)에서 유일한 데이터 결손이 -2020-04-01 ~ 2020-10-31 (약 138 평일). COVID 폭락장 + 회복장을 포함하는 -이 구간이 없으면 멀티레짐 검정력이 크게 훼손된다. - -특성: - - FINRA CDN 공개 무료 (API 키 불필요) - - idempotent: 이미 수집된 날짜는 자동 스킵 (_count_for_date 검사) - - 주말/공휴일(파일 없는 날)은 404 수신 후 자동 스킵 - - 1년치 약 20-40분 소요 → 이 스크립트는 7개월, 약 10-20분 예상 - -Usage (컨테이너 내부): - python scripts/backfill_finra_2020_gap.py - -Usage (호스트): - docker exec stock_oracle_api python scripts/backfill_finra_2020_gap.py -""" - -import asyncio -import fcntl -import logging -import os -import sys -from datetime import date - -sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) - -logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s") -logger = logging.getLogger("finra_gap_fill") - -_LOCK_FILE = "/tmp/backfill_finra_2020.lock" - -# 2020-10 is partially ingested (15 days) → start from 2020-04-01, end 2020-11-01 -_GAP_START = date(2020, 4, 1) -_GAP_END = date(2020, 11, 1) - - -async def run_gap_fill(): - from app.core.database import AsyncSessionLocal - from app.services.finra_short_volume_service import FinraShortVolumeService - - logger.info(f"Starting FINRA gap fill: {_GAP_START} → {_GAP_END}") - - svc = FinraShortVolumeService() - async with AsyncSessionLocal() as db: - inserted_total = await svc.ingest_date_range( - db=db, - start_date=_GAP_START, - end_date=_GAP_END, - force_refresh=False, # skip already-ingested dates - ) - - logger.info("=" * 60) - logger.info("FINRA 2020 GAP FILL COMPLETE") - logger.info(f" Period : {_GAP_START} → {_GAP_END}") - logger.info(f" Inserted : {inserted_total:,} rows (0 = all dates already present or holiday)") - logger.info("=" * 60) - logger.info("Verify with:") - logger.info(" SELECT to_char(date,'YYYY-MM'), COUNT(DISTINCT date::date)") - logger.info(" FROM finra_short_volume") - logger.info(" WHERE date >= '2020-03-01' AND date < '2021-02-01'") - logger.info(" GROUP BY 1 ORDER BY 1;") - - -if __name__ == "__main__": - lock_fh = open(_LOCK_FILE, "w") - try: - fcntl.flock(lock_fh, fcntl.LOCK_EX | fcntl.LOCK_NB) - except OSError: - logger.error("Gap fill is already running (lock file held). Exiting.") - sys.exit(1) - - try: - asyncio.run(run_gap_fill()) - finally: - fcntl.flock(lock_fh, fcntl.LOCK_UN) - lock_fh.close() - try: - os.unlink(_LOCK_FILE) - except OSError: - pass diff --git a/scripts/export_pit_panel.py b/scripts/export_pit_panel.py deleted file mode 100644 index 4804f65..0000000 --- a/scripts/export_pit_panel.py +++ /dev/null @@ -1,257 +0,0 @@ -""" -PIT(Point-in-Time) 패널 데이터 parquet 익스포트. - -FINRA short-volume ⋈ Alpaca 일봉 가격을 (date, symbol)로 조인하여 -생존편향-0 리서치 패널을 data/pit_panel.parquet에 저장한다. - -리서치 레이어에 Postgres 직접 접속이 없을 때 이 파일을 핸드오프로 사용. - -패널 컬럼: - date, symbol, short_volume, short_exempt_volume, total_volume, short_ratio, - open, high, low, close, volume, vwap, - (옵션) sector, exchange ← 현재 활성 종목 한정 - -심볼 정규화: - FINRA symbol과 alpaca_price_data ticker 간 컨벤션 차이 (BRK/B vs BRK.B)를 - 동일 정규화 함수로 매핑. 조인 후 매칭 0행인 심볼 수를 리포트한다. - -사전 조건: - - backfill_alpaca_daily_pit.py 완료 후 실행 - - pandas + pyarrow 설치 필요: pip install pandas pyarrow - -Usage (컨테이너 내부): - python scripts/export_pit_panel.py - -Usage (호스트): - docker exec stock_oracle_api python scripts/export_pit_panel.py - # 출력: ./data/pit_panel.parquet (볼륨 마운트로 호스트에서 접근 가능) - -옵션: - --start YYYY-MM-DD 데이터 시작일 (기본: 2018-08-01) - --end YYYY-MM-DD 데이터 종료일 (기본: 오늘) - --out PATH 출력 경로 (기본: data/pit_panel.parquet) -""" - -import argparse -import asyncio -import logging -import os -import sys -from datetime import datetime, timezone - -sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) - -logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s") -logger = logging.getLogger("pit_export") - -_DEFAULT_START = "2018-08-01" -_DEFAULT_OUT = os.path.join( - os.path.dirname(os.path.dirname(os.path.abspath(__file__))), - "data", "pit_panel.parquet" -) - - -def finra_to_alpaca(symbol: str) -> str: - """FINRA 심볼 → Alpaca 정규화 (/ → ., - → .).""" - return symbol.replace("/", ".").replace("-", ".") - - -async def run_export(start_str: str, end_str: str, out_path: str): - try: - import pandas as pd - except ImportError: - logger.error("pandas not installed. Run: pip install pandas pyarrow") - return - - try: - import pyarrow # noqa: F401 - except ImportError: - logger.error("pyarrow not installed. Run: pip install pyarrow") - return - - from sqlalchemy import text - from app.core.database import AsyncSessionLocal - - logger.info(f"PIT 패널 익스포트: {start_str} → {end_str}") - - # ------------------------------------------------------------------ # - # Step 1: FINRA short-volume 전체 로드 # - # ------------------------------------------------------------------ # - logger.info("FINRA short-volume 로드 중 ...") - async with AsyncSessionLocal() as db: - result = await db.execute( - text(""" - SELECT - date::date AS date, - symbol, - short_volume, - short_exempt_volume, - total_volume, - short_ratio - FROM finra_short_volume - WHERE date >= :start AND date <= :end - ORDER BY date, symbol - """), - {"start": start_str, "end": end_str}, - ) - finra_rows = result.fetchall() - - logger.info(f"FINRA rows: {len(finra_rows):,}") - - if not finra_rows: - logger.error("FINRA 데이터 없음. 백필 먼저 실행하세요.") - return - - df_finra = pd.DataFrame(finra_rows, columns=[ - "date", "symbol", "short_volume", "short_exempt_volume", - "total_volume", "short_ratio", - ]) - # 조인 키: 정규화된 심볼 - df_finra["symbol_norm"] = df_finra["symbol"].apply(finra_to_alpaca).str.upper() - - # ------------------------------------------------------------------ # - # Step 2: Alpaca 1d 가격 로드 # - # ------------------------------------------------------------------ # - logger.info("Alpaca 1d 가격 로드 중 ...") - async with AsyncSessionLocal() as db: - result = await db.execute( - text(""" - SELECT - date::date AS date, - ticker, - open, high, low, close, volume, vwap - FROM alpaca_price_data - WHERE interval = '1d' - AND date >= :start AND date <= :end - ORDER BY date, ticker - """), - {"start": start_str, "end": end_str}, - ) - price_rows = result.fetchall() - - logger.info(f"Alpaca price rows: {len(price_rows):,}") - - df_price = pd.DataFrame(price_rows, columns=[ - "date", "ticker", "open", "high", "low", "close", "volume", "vwap", - ]) - df_price["symbol_norm"] = df_price["ticker"].apply(finra_to_alpaca).str.upper() - - # ------------------------------------------------------------------ # - # Step 3: 조인 # - # ------------------------------------------------------------------ # - logger.info("조인 중 ...") - # FINRA에 market(B,Q,N) 컬럼이 있어 (date, symbol)이 여러 행일 수 있음. - # 리서치용으로 시장 통합(aggregated) 값을 사용하는 것이 일반적. - # 먼저 (date, symbol) 기준으로 집계 (이미 집계된 B,Q,N 합산 행이 보통이지만 혹시 분리돼 있을 경우) - df_finra_agg = ( - df_finra - .groupby(["date", "symbol", "symbol_norm"], as_index=False) - .agg( - short_volume = ("short_volume", "sum"), - short_exempt_volume= ("short_exempt_volume","sum"), - total_volume = ("total_volume", "sum"), - ) - ) - df_finra_agg["short_ratio"] = ( - df_finra_agg["short_volume"] / df_finra_agg["total_volume"].replace(0, float("nan")) - ) - - # Alpaca에도 (date, ticker) 중복이 있을 수 있음 (adjustment 차이 등) → dedup - df_price_dedup = df_price.drop_duplicates(subset=["date", "symbol_norm"], keep="last") - - df_panel = df_finra_agg.merge( - df_price_dedup[["date", "symbol_norm", "open", "high", "low", "close", "volume", "vwap"]], - on=["date", "symbol_norm"], - how="left", - ) - - # ------------------------------------------------------------------ # - # Step 4: 커버리지 체크 # - # ------------------------------------------------------------------ # - matched = df_panel["close"].notna().sum() - total_rows = len(df_panel) - missing_pct = 100.0 * (total_rows - matched) / total_rows if total_rows else 0 - - no_price_syms = set( - df_panel.loc[df_panel["close"].isna(), "symbol"].unique() - ) - - logger.info(f"패널 총 행: {total_rows:,}") - logger.info(f"가격 매칭 행: {matched:,} ({100-missing_pct:.1f}%)") - logger.info(f"가격 누락 행: {total_rows-matched:,} ({missing_pct:.1f}%)") - logger.info(f"가격 없는 고유 심볼 수: {len(no_price_syms):,}") - - if no_price_syms: - sample = sorted(no_price_syms)[:30] - logger.info(f"가격 없는 심볼 샘플(30): {sample}") - - # ------------------------------------------------------------------ # - # Step 5: 옵션 — sector/exchange 추가 (활성 종목 한정) # - # ------------------------------------------------------------------ # - logger.info("sector/exchange 메타데이터 조인 중 (활성 종목 한정) ...") - async with AsyncSessionLocal() as db: - result = await db.execute( - text(""" - SELECT ticker, sector, exchange - FROM universe_ticker_registry - WHERE is_active = true - """) - ) - meta_rows = result.fetchall() - - df_meta = pd.DataFrame(meta_rows, columns=["ticker", "sector", "exchange"]) - df_meta["symbol_norm"] = df_meta["ticker"].apply(finra_to_alpaca).str.upper() - - df_panel = df_panel.merge( - df_meta[["symbol_norm", "sector", "exchange"]], - on="symbol_norm", - how="left", - ) - - # ------------------------------------------------------------------ # - # Step 6: 컬럼 정리 + 저장 # - # ------------------------------------------------------------------ # - # 최종 컬럼 순서 - keep_cols = [ - "date", "symbol", - "short_volume", "short_exempt_volume", "total_volume", "short_ratio", - "open", "high", "low", "close", "volume", "vwap", - "sector", "exchange", - ] - df_out = df_panel[keep_cols].copy() - df_out["date"] = pd.to_datetime(df_out["date"]) - - os.makedirs(os.path.dirname(out_path), exist_ok=True) - df_out.to_parquet(out_path, index=False, engine="pyarrow") - - file_mb = os.path.getsize(out_path) / 1_048_576 - logger.info("=" * 70) - logger.info("PIT 패널 익스포트 완료") - logger.info(f" 출력 파일 : {out_path}") - logger.info(f" 파일 크기 : {file_mb:.1f} MB") - logger.info(f" 패널 shape : {df_out.shape}") - logger.info(f" 날짜 범위 : {df_out['date'].min().date()} ~ {df_out['date'].max().date()}") - logger.info(f" 고유 심볼 : {df_out['symbol'].nunique():,}") - logger.info(f" 가격 커버 : {df_out['close'].notna().mean()*100:.1f}%") - logger.info(f" 잔존편향 : {missing_pct:.2f}% (row-weighted)") - logger.info("=" * 70) - logger.info("") - logger.info("사용 예 (Python):") - logger.info(" import pandas as pd") - logger.info(f" df = pd.read_parquet('{out_path}')") - logger.info(" # PIT 유니버스 (특정 날짜 활동 종목):") - logger.info(" pit_2023_03_09 = df[df.date == '2023-03-09']['symbol'].unique()") - logger.info(" # 그날 공매도비율:") - logger.info(" df[df.date == '2023-03-09'][['symbol','short_ratio','close']]") - - -if __name__ == "__main__": - parser = argparse.ArgumentParser(description="PIT 패널 parquet 익스포트") - parser.add_argument("--start", default=_DEFAULT_START, help="시작일 (YYYY-MM-DD)") - parser.add_argument("--end", - default=datetime.now(timezone.utc).strftime("%Y-%m-%d"), - help="종료일 (YYYY-MM-DD)") - parser.add_argument("--out", default=_DEFAULT_OUT, help="출력 parquet 경로") - args = parser.parse_args() - - asyncio.run(run_export(args.start, args.end, args.out))