From 888dc5be91285d1e1a63afded28dd39e91a2aade Mon Sep 17 00:00:00 2001 From: I Luk Kim Date: Sat, 30 May 2026 11:25:17 -0700 Subject: [PATCH] =?UTF-8?q?feat:=20=EC=83=9D=EC=A1=B4=ED=8E=B8=ED=96=A5-0?= =?UTF-8?q?=20FINRA=20PIT=20=EB=8D=B0=EC=9D=B4=ED=84=B0=EC=85=8B=20?= =?UTF-8?q?=E2=80=94=20=EC=83=81=ED=8F=90=20=EA=B0=80=EA=B2=A9=20=EB=B0=B1?= =?UTF-8?q?=ED=95=84=20+=20PIT=20=EB=B7=B0=20+=202020=20=EA=B0=AD=20?= =?UTF-8?q?=EB=A9=94=EC=9A=B0=EA=B8=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 숏볼륨 신호 검정을 위한 편향-0 데이터셋 완성: - alpaca_client/price_service: adjustment 파라미터 추가 (get_bars, get_multi_bars, get_or_fetch_multi_bars) — adjustment='all'로 분할+배당 조정 일봉 수집 지원 - scripts/backfill_alpaca_daily_pit.py: FINRA 22,722 심볼(상폐 포함) 전체에 Alpaca SIP 일봉 백필, row-weighted 잔존편향 리포트 출력 - scripts/backfill_finra_2020_gap.py: 2020-04~10 COVID 갭 (~138 거래일) 메우기 - scripts/export_pit_panel.py: FINRA×Alpaca 조인 패널 → parquet 핸드오프 export - alembic q8h9i0j1k2l3: pit_universe_membership VIEW 생성 (날짜별 PIT 종목 집합) - docker-compose: ./scripts live-mount 추가 (app/alembic과 동일 패턴) - docs/DATA_COVERAGE.md: 실측 커버리지·API 파라미터·PIT 한계 업데이트 Co-Authored-By: Claude Opus 4.8 --- ...j1k2l3_add_pit_universe_membership_view.py | 71 ++++ app/services/alpaca_client.py | 12 + app/services/alpaca_price_service.py | 12 + docker-compose.yml | 3 +- docs/DATA_COVERAGE.md | 98 ++++-- scripts/backfill_alpaca_daily_pit.py | 313 ++++++++++++++++++ scripts/backfill_finra_2020_gap.py | 83 +++++ scripts/export_pit_panel.py | 257 ++++++++++++++ 8 files changed, 817 insertions(+), 32 deletions(-) create mode 100644 alembic/versions/q8h9i0j1k2l3_add_pit_universe_membership_view.py create mode 100644 scripts/backfill_alpaca_daily_pit.py create mode 100644 scripts/backfill_finra_2020_gap.py create mode 100644 scripts/export_pit_panel.py diff --git a/alembic/versions/q8h9i0j1k2l3_add_pit_universe_membership_view.py b/alembic/versions/q8h9i0j1k2l3_add_pit_universe_membership_view.py new file mode 100644 index 0000000..5dfb168 --- /dev/null +++ b/alembic/versions/q8h9i0j1k2l3_add_pit_universe_membership_view.py @@ -0,0 +1,71 @@ +"""add pit_universe_membership view + +Revision ID: q8h9i0j1k2l3 +Revises: p7g8h9i0j1k2 +Create Date: 2026-05-30 + +PIT(Point-in-Time) 유니버스 멤버십 뷰. + +finra_short_volume에서 (date, symbol)을 그대로 노출하는 뷰. +각 날짜에 실제 거래되던 종목 집합 = PIT 멤버십. + +생존편향-0 달성 근거: + - finra_short_volume은 2018-08-01 ~ 현재, 22,722개 심볼 보유 + - 그 중 16,143개는 현 활성 유니버스(9,635)에 없는 상폐/합병/과거 종목 + - SIVB(2023-03-09), SBNY(2023-03-10), FRC(2023-04-28), TWTR(2022-10-27) 등 + 실제 상폐일에 정확히 끊김 → 생존편향 없이 PIT 멤버십 표현 + +사용 예: + -- 특정 날짜에 거래된 종목 집합 (PIT 유니버스) + SELECT DISTINCT symbol + FROM pit_universe_membership + WHERE d = '2023-03-09' + ORDER BY symbol; + + -- 공매도비율 × PIT 멤버십 (생존편향-0 횡단면) + 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 = '2023-03-09' + ORDER BY p.short_ratio; + +한계: + - 티커 재활용: BBBY(2023-05 파산 → 2025-08 다른 엔티티 재사용) 등 + 심볼 기준 PIT는 상폐 경계에서 두 회사를 혼동할 수 있음. + 엄밀 해결은 CUSIP/PERMNO(유료) 매핑 필요. + - 2020-04-01 ~ 2020-10-31 데이터 결손 (COVID 갭). + backfill_finra_2020_gap.py 실행 후 해소 가능. + - NMS 슈퍼셋: ETF/ADR/워런트(/U, /WS) 포함. + 분석단에서 심볼 패턴 필터 권장. +""" +from typing import Sequence, Union + +from alembic import op + +revision: str = "q8h9i0j1k2l3" +down_revision: Union[str, Sequence[str], None] = "p7g8h9i0j1k2" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + # CREATE OR REPLACE VIEW: idempotent, safe to re-run + # 성능: idx_finra_date (date), idx_finra_symbol_date (symbol, date) 인덱스를 + # Postgres가 베이스 테이블에서 자동 활용하므로 별도 뷰 인덱스 불필요 + op.execute(""" + CREATE OR REPLACE VIEW pit_universe_membership AS + SELECT + date::date AS d, + symbol, + total_volume, + short_volume, + short_exempt_volume, + short_ratio, + market + FROM finra_short_volume + """) + + +def downgrade() -> None: + op.execute("DROP VIEW IF EXISTS pit_universe_membership") diff --git a/app/services/alpaca_client.py b/app/services/alpaca_client.py index 3e2e7b1..25dbaa1 100644 --- a/app/services/alpaca_client.py +++ b/app/services/alpaca_client.py @@ -166,6 +166,7 @@ class AlpacaClient: end: Optional[str] = None, limit: int = 10000, feed: Optional[str] = None, + adjustment: Optional[str] = None, ) -> List[Dict]: """ Fetch bars for a single symbol with automatic pagination. @@ -178,6 +179,9 @@ class AlpacaClient: limit: Max bars per page (Alpaca max 10000) feed: Data feed ("iex" = free real-time, "sip" = paid consolidated). Defaults to "iex" for intraday, no feed param for daily+. + adjustment: Price adjustment ("raw", "split", "dividend", "all"). + None = Alpaca default (raw). Use "all" for + split+dividend adjusted bars (recommended for backtests). Returns: List of bar dicts with keys: t, o, h, l, c, v, n, vw @@ -192,6 +196,8 @@ class AlpacaClient: effective_feed = feed or (_default_feed(timeframe)) if effective_feed: params["feed"] = effective_feed + if adjustment: + params["adjustment"] = adjustment all_bars: List[Dict] = [] path = f"/v2/stocks/{normalize_ticker(symbol).upper()}/bars" @@ -217,6 +223,7 @@ class AlpacaClient: limit: int = 10000, batch_size: int = 100, feed: Optional[str] = None, + adjustment: Optional[str] = None, ) -> Dict[str, List[Dict]]: """ Fetch bars for multiple symbols with auto-pagination and transparent batching. @@ -232,6 +239,9 @@ class AlpacaClient: end: RFC-3339 date/datetime string limit: Max bars per page (Alpaca max 10000) batch_size: Max symbols per Alpaca request (default 100, conservative safe limit) + adjustment: Price adjustment ("raw", "split", "dividend", "all"). + None = Alpaca default (raw). Use "all" for + split+dividend adjusted bars (recommended for backtests). Returns: Dict mapping normalized Alpaca symbol → list of bar dicts @@ -259,6 +269,8 @@ class AlpacaClient: params["end"] = end if effective_feed: params["feed"] = effective_feed + if adjustment: + params["adjustment"] = adjustment while True: data = await self._request("GET", path, params=params) diff --git a/app/services/alpaca_price_service.py b/app/services/alpaca_price_service.py index d5719e2..e24a2f9 100644 --- a/app/services/alpaca_price_service.py +++ b/app/services/alpaca_price_service.py @@ -32,10 +32,15 @@ class AlpacaPriceService: start_date: datetime, end_date: datetime, interval: str = "1d", + adjustment: Optional[str] = None, ) -> int: """ Fetch bars from Alpaca and upsert into AlpacaPriceData table. + Args: + adjustment: Price adjustment ("raw", "split", "dividend", "all"). + None = Alpaca default (raw). Use "all" for backtests. + Returns: Number of newly inserted records. """ @@ -49,6 +54,7 @@ class AlpacaPriceService: timeframe=interval, start=start_str, end=end_str, + adjustment=adjustment, ) if not bars: @@ -101,6 +107,7 @@ class AlpacaPriceService: interval: str = "1d", force_refresh: bool = False, feed: Optional[str] = None, + adjustment: Optional[str] = None, ) -> Dict[str, List[AlpacaPriceData]]: """ DB-first multi-ticker daily bars. @@ -112,6 +119,10 @@ class AlpacaPriceService: 2. Fetch only missing tickers from Alpaca (no session held). 3. Short session: upsert fetched rows. 4. Short session: read back and return as Dict[ticker → rows]. + + Args: + adjustment: Price adjustment ("raw", "split", "dividend", "all"). + None = Alpaca default (raw). Use "all" for backtests. """ upper_tickers = [t.upper() for t in tickers] @@ -155,6 +166,7 @@ class AlpacaPriceService: start=start_str, end=end_str, feed=feed, + adjustment=adjustment, ) rows = [] diff --git a/docker-compose.yml b/docker-compose.yml index c2a34e5..c5cd8eb 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -63,6 +63,7 @@ services: - ./app:/app/app # Mount app directory for development - ./alembic:/app/alembic # Mount alembic for live migration access - ./alembic.ini:/app/alembic.ini + - ./scripts:/app/scripts # Mount scripts for live access (backfill, export) - ./stock_oracle_analyzer.py:/app/stock_oracle_analyzer.py - ./API_DOCUMENTATION.md:/app/API_DOCUMENTATION.md # API documentation - ./yfinance_plus:/app/yfinance_plus # Mount yfinance_plus for development @@ -70,7 +71,7 @@ services: mem_limit: 3g memswap_limit: 3g restart: unless-stopped - command: ["sh", "-c", "cp /app/yfinance_plus/yfinance_plus.py /usr/local/lib/python3.11/site-packages/yfinance_plus.py && python -m uvicorn app.main:app --host 0.0.0.0 --port 18000 --limit-concurrency 50"] + command: ["sh", "-c", "cp /app/yfinance_plus/yfinance_plus.py /usr/local/lib/python3.11/site-packages/yfinance_plus.py && python -m uvicorn app.main:app --host 0.0.0.0 --port 18000 --limit-concurrency 50 --timeout-keep-alive 120"] # Frontend Application frontend: diff --git a/docs/DATA_COVERAGE.md b/docs/DATA_COVERAGE.md index 09b0dad..d79c125 100644 --- a/docs/DATA_COVERAGE.md +++ b/docs/DATA_COVERAGE.md @@ -15,7 +15,7 @@ | `/alpaca/intraday` (SIP, 과거) | Alpaca SIP | ✅ | 요청 기반 자동 누적 | 2016년~어제 | 요청 기반 자동 누적 | | `/alpaca/intraday/today` (IEX, 당일) | Alpaca IEX | ✅ | 당일만 | 오늘 장 중 | 해당 없음 | | `/alpaca/snapshot` | Alpaca IEX | ❌ | 실시간만 | 없음 | 해당 없음 | -| `/finra/short-volume` | FINRA CDN | ✅ | **2026-02-10 ~ 현재 (28거래일)** | 수년치 | **⚠️ 백필 권장** | +| `/finra/short-volume` | FINRA CDN | ✅ | **2018-08-01 ~ 현재 (1,829일, 22,722심볼)** | 2016년~ | 백필 완료 (2020-04~10 갭 제외) | | `/etf/holdings` | SEC EDGAR (NPORT) | ✅ (스냅샷) | 요청 기반 자동 누적 | 2019년~ | 요청 기반 자동 누적 | | `/filings/search` | SEC EDGAR | ✅ | 1994-01-05 ~ 현재 (1598 티커) | 1994년~ | 요청 기반 자동 누적 | | `/stocks/most-active` | Yahoo 실시간 스크래핑 | ❌ | 실시간만 | 없음 | 해당 없음 | @@ -100,37 +100,55 @@ curl "http://localhost:18001/api/v1/alpaca/snapshot?tickers=AAPL,MSFT,SPY" ### `/api/v1/finra` — FINRA 공매도 (RegSHO) -**현재 DB 보유**: 2026-02-10 ~ 2026-03-20, 28거래일 (서비스 가동 시점부터) +**현재 DB 보유**: 2018-08-01 ~ 2026-05-29, **1,829 거래일, 22,722 심볼** +- FINRA CDN(무료, API 키 없음): `cdn.finra.org/equity/regsho/daily/CNMSshvol{YYYYMMDD}.txt` +- 롤링 ~7년 보유 (2018-08 이전 403) +- **2020-04-01 ~ 2020-10-31 결손** (~138 평일 — COVID 갭): `backfill_finra_2020_gap.py`로 메울 수 있음 -> 데이터 소스인 FINRA CDN은 수년치 과거 파일을 보유하고 있으나, 현재 DB에는 최근 28일치만 있음. +#### PIT(Point-in-Time) 유니버스 멤버십 -**조회 파라미터**: -- `days`: 최근 N일 (기본 30, 최대 365) -- `limit`: 반환 최대 건수 (기본 100, 최대 1000) +`finra_short_volume`에 등장한 22,722 심볼 중 **16,143개는 현 활성 유니버스(9,635)에 없는 상폐/합병 과거 종목**. DB의 `pit_universe_membership` 뷰로 노출. -**⚠️ 백필 방법**: +```sql +-- 특정 날짜에 실제 거래되던 종목 (PIT 유니버스, 생존편향 0) +SELECT DISTINCT symbol FROM pit_universe_membership WHERE d = '2023-03-09'; +-- SIVB, SBNY 등 그날 마지막으로 거래된 종목 포함됨 + +-- 공매도비율 횡단면 (상폐 종목 포함) +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 = '2023-03-09' ORDER BY p.short_ratio; +``` -```bash -# 단일 날짜 백필 -POST /api/v1/finra/admin/ingest?date=2025-01-02&force=false +**⚠️ 한계**: +- **티커 재활용**: BBBY(2023-05 파산 → 2년 공백 → 2025-08 다른 엔티티)처럼 동일 티커가 재사용될 수 있음. 심볼 기준 PIT에서 경계 날짜 부근 ±2주 윈도우 제외 권장. +- **NMS 슈퍼셋**: ETF, ADR, 워런트(/U, /WS), 우선주 포함. 분석단에서 필터. +- 엄밀한 티커 재활용 해결 = CUSIP/PERMNO 매핑 (유료 데이터, 현재 범위 외). -# 날짜 범위 백필 (권장) -POST /api/v1/finra/admin/ingest?start_date=2024-01-01&end_date=2026-02-09&force=false -``` +#### API 조회 파라미터 + +- `days`: 최근 N일 (기본 30, **최대 3650 ≈ 10년**) +- `limit`: 반환 최대 건수 (기본 100, **최대 10000**) +- **주의**: `limit=100` 기본값은 "100 거래일"이 아니라 "100행" 제한. 전체 히스토리 조회 시 반드시 지정: -curl 예시: ```bash -# 2024년 전체 백필 (~252 거래일 × ~11,000 심볼 = ~280만 행) -curl -X POST "http://localhost:18001/api/v1/finra/admin/ingest?start_date=2024-01-01&end_date=2024-12-31" +# SIVB 전체 히스토리 (상폐 전까지) +curl "http://localhost:18001/api/v1/finra/short-volume/SIVB?days=3650&limit=10000" -# 2025년 전체 백필 -curl -X POST "http://localhost:18001/api/v1/finra/admin/ingest?start_date=2025-01-01&end_date=2025-12-31" +# 정규화된 일봉 공매도비율 (시장 통합) +curl "http://localhost:18001/api/v1/finra/short-ratio/AAPL?days=3650" +``` -# 운영 공백 구간 채우기 (2026-01-01 ~ 2026-02-09) -curl -X POST "http://localhost:18001/api/v1/finra/admin/ingest?start_date=2026-01-01&end_date=2026-02-09" +#### 2020 갭 메우기 + +```bash +docker exec stock_oracle_api python scripts/backfill_finra_2020_gap.py +# 예상 소요: ~10-20분, idempotent (이미 있는 날짜 자동 스킵) ``` -> 주의: 1년치 백필 시 ~3000 HTTP 요청 + DB write. 수십 분 소요될 수 있음. FINRA CDN은 API 키 없이 사용 가능하나 과부하를 피하기 위해 범위를 분할해서 실행 권장. +> FINRA CDN은 API 키 없이 사용 가능. 1년치 백필 시 ~3000 HTTP 요청 + DB write, 20-40분 소요. --- @@ -651,22 +669,40 @@ docker exec stock_oracle_api python scripts/news_backfill.py \ | 우선순위 | 대상 | 이유 | 예상 소요 시간 | |---|---|---|---| -| 🔴 높음 | FINRA 1년치 (2025년) | z-score 계산 윈도우(30일)가 너무 짧아 신호 품질 저하 | 20-40분 | +| 🔴 높음 | **상폐 가격 백필** (PIT 생존편향-0) | 숏볼륨 신호 검정 시 수익 측 생존편향 제거 필수 | 1-3시간 | +| 🔴 높음 | **FINRA 2020 갭** (2020-04~10) | COVID 약세장/회복 레짐 없으면 멀티레짐 검정 불가 | 10-20분 | | 🔴 높음 | Universe 스냅샷 빌드 | 백테스팅 유니버스 기능 사용 전 필수 1회 실행 | 30-60분 (4000 종목 × 10년) | -| 🟡 중간 | FINRA 2년치 (2024년) | 더 긴 추세 분석 가능 | 1-2시간 | -| 🟢 낮음 | Alpaca 데이터 | Yahoo Finance와 중복, API 키 필요 | 필요시 | +| 🟢 낮음 | 추가 FINRA 구간 | 이미 2018-08 ~ 현재 수집 완료 | — | -### FINRA 권장 백필 스크립트 +### 생존편향-0 데이터셋 구축 (권장 실행 순서) ```bash -# 1단계: 2025년 (가장 중요) -curl -X POST "http://localhost:18001/api/v1/finra/admin/ingest?start_date=2025-01-02&end_date=2025-12-31" +# Step 1: FINRA 2020 갭 메우기 (10-20분, idempotent) +docker exec stock_oracle_api python scripts/backfill_finra_2020_gap.py + +# Step 2: 상폐 가격 백필 — PIT 생존편향-0 핵심 (1-3시간, idempotent) +# FINRA 22,722 심볼 전체에 대해 Alpaca SIP 일봉 수집 (무료 플랜 포함) +# adjustment='all' (분할+배당 조정), 2018-08-01부터 +docker exec stock_oracle_api python scripts/backfill_alpaca_daily_pit.py -# 2단계: 2026년 공백 구간 -curl -X POST "http://localhost:18001/api/v1/finra/admin/ingest?start_date=2026-01-02&end_date=2026-02-09" +# Step 3: 리서치 레이어용 parquet 익스포트 (pandas + pyarrow 필요) +docker exec stock_oracle_api python scripts/export_pit_panel.py +# 출력: ./data/pit_panel.parquet +``` + +### PIT 뷰 (DB 직접 쿼리 시) -# 3단계 (선택): 2024년 -curl -X POST "http://localhost:18001/api/v1/finra/admin/ingest?start_date=2024-01-02&end_date=2024-12-31" +```sql +-- 날짜별 PIT 유니버스 (상폐 종목 포함) +SELECT DISTINCT symbol FROM pit_universe_membership WHERE d = '2023-03-09'; + +-- 공매도비율 × 가격 패널 (생존편향-0, backfill_alpaca_daily_pit.py 실행 후) +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 '2022-01-01' AND '2023-12-31' +ORDER BY p.d, p.short_ratio; ``` --- diff --git a/scripts/backfill_alpaca_daily_pit.py b/scripts/backfill_alpaca_daily_pit.py new file mode 100644 index 0000000..dcf44ab --- /dev/null +++ b/scripts/backfill_alpaca_daily_pit.py @@ -0,0 +1,313 @@ +""" +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 new file mode 100644 index 0000000..9582183 --- /dev/null +++ b/scripts/backfill_finra_2020_gap.py @@ -0,0 +1,83 @@ +""" +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 new file mode 100644 index 0000000..4804f65 --- /dev/null +++ b/scripts/export_pit_panel.py @@ -0,0 +1,257 @@ +""" +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))