Compare commits

...

2 Commits

Author SHA1 Message Date
I Luk Kim 964cf2237a fix: parquet export 청킹(OOM 방지) + SQL COPY 방식으로 전환
- finra.py export: fetchall() → 연도별 청킹 + ParquetWriter 스트리밍
  (20M row fetchall OOM 방지)
- pit_panel.parquet 실제 생성: 18.6M rows, 757MB, NVDA $120.82 조정가 확인

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2 months ago
I Luk Kim f5feffb80e fix: adjustment='all' 조정가 적용 + parquet export 엔드포인트
- alpaca.py backfill: on_conflict_do_nothing → on_conflict_do_update
  기존 미조정 행(split 전 $1,208)을 조정가($120)로 덮어씀
  NVDA 2024-06-07 ~$120 통과 기준
- finra.py: POST /admin/export-pit-panel (Background) + GET /download
  pit_universe_membership × alpaca_price_data 조인 → /app/data/pit_panel.parquet
- requirements-api.txt: pyarrow>=14.0.0 추가

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2 months ago

@ -447,7 +447,19 @@ async def _run_pit_backfill(start_str: str, end_str: str, force: bool) -> None:
stmt = pg_insert(AlpacaPriceData).values(
rows_to_insert[i: i + CHUNK_SIZE]
)
stmt = stmt.on_conflict_do_nothing(constraint="uq_alpaca_price_data")
stmt = stmt.on_conflict_do_update(
constraint="uq_alpaca_price_data",
set_={
"open": stmt.excluded.open,
"high": stmt.excluded.high,
"low": stmt.excluded.low,
"close": stmt.excluded.close,
"volume": stmt.excluded.volume,
"vwap": stmt.excluded.vwap,
"trade_count": stmt.excluded.trade_count,
"data_source": stmt.excluded.data_source,
}
)
result = await db.execute(stmt)
total_inserted += result.rowcount
await asyncio.sleep(0)
@ -481,7 +493,7 @@ async def _run_pit_backfill(start_str: str, end_str: str, force: bool) -> None:
"Alpaca SIP 일봉(1d)을 백필합니다. 생존편향-0 수익 계산에 필요.\n\n"
"**특성**:\n"
"- `adjustment=all` (분할+배당 조정) — 상폐 종목은 future-proof\n"
"- DB-first, idempotent (`on_conflict_do_nothing`) — 재실행 안전\n"
"- `on_conflict_do_update`: 기존 행도 조정가로 덮어씀 (split 이전 미조정 데이터 수정)\n"
"- Alpaca SIP 일봉은 무료 플랜에서 2016-01-04부터 제공\n"
"- 기본 시작일: 2018-08-01 (FINRA DB 시작일)\n\n"
"**백그라운드 실행**: 즉시 `started` 응답, 1-3시간 소요.\n"

@ -2,14 +2,20 @@
FINRA Short Sale Volume endpoints
"""
import logging
import os
from datetime import date, datetime, timedelta, timezone
from typing import List, Optional
from fastapi import APIRouter, Depends, HTTPException, Query
from fastapi.responses import Response
from fastapi import APIRouter, BackgroundTasks, Depends, HTTPException, Query
from fastapi.responses import FileResponse, Response
from sqlalchemy import text
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.database import AsyncSessionLocal
logger = logging.getLogger(__name__)
from app.core.config import settings
from app.core.database import get_db
from app.schemas.finra import (
@ -184,6 +190,123 @@ async def get_pit_panel(
}
_EXPORT_PATH = "/app/data/pit_panel.parquet"
async def _run_pit_export(start_str: str, end_str: str) -> None:
"""Background task: export PIT panel to /app/data/pit_panel.parquet."""
try:
import pandas as pd
except ImportError:
logger.error("pandas not installed"); return
try:
import pyarrow # noqa: F401
except ImportError:
logger.error("pyarrow not installed"); return
logger.info(f"PIT export: loading {start_str}{end_str} ...")
from datetime import date as date_type
import pyarrow as pa
import pyarrow.parquet as pq
d_from = date_type.fromisoformat(start_str)
d_to = date_type.fromisoformat(end_str)
COLS = ["date","symbol","short_volume","short_exempt_volume","total_volume",
"short_ratio","open","high","low","close","price_volume","vwap"]
SQL = text("""
SELECT p.d::date AS date, p.symbol,
p.short_volume, p.short_exempt_volume, p.total_volume, p.short_ratio,
a.open, a.high, a.low, a.close,
a.volume AS price_volume, a.vwap
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 :s AND :e
ORDER BY p.d, p.symbol
""")
os.makedirs(os.path.dirname(_EXPORT_PATH), exist_ok=True)
tmp_path = _EXPORT_PATH + ".tmp"
# 연도별 청킹으로 메모리 제한 내 처리
writer = None
total_rows = 0
year = d_from.year
while date_type(year, 1, 1) <= d_to:
chunk_start = max(d_from, date_type(year, 1, 1))
chunk_end = min(d_to, date_type(year, 12, 31))
logger.info(f"PIT export: fetching {chunk_start}{chunk_end} ...")
async with AsyncSessionLocal() as db:
rows = (await db.execute(SQL, {"s": chunk_start, "e": chunk_end})).fetchall()
if rows:
df_chunk = pd.DataFrame(rows, columns=COLS)
df_chunk["date"] = pd.to_datetime(df_chunk["date"])
table = pa.Table.from_pandas(df_chunk, preserve_index=False)
if writer is None:
writer = pq.ParquetWriter(tmp_path, table.schema, compression="snappy")
writer.write_table(table)
total_rows += len(rows)
logger.info(f"PIT export: {year} done — {len(rows):,} rows (total {total_rows:,})")
year += 1
if writer:
writer.close()
os.replace(tmp_path, _EXPORT_PATH)
mb = os.path.getsize(_EXPORT_PATH) / 1_048_576 if os.path.exists(_EXPORT_PATH) else 0
logger.info(
f"PIT export complete: {_EXPORT_PATH} | {mb:.0f} MB | "
f"total_rows={total_rows:,} | cols={COLS}"
)
@router.post(
"/admin/export-pit-panel",
summary="PIT 패널 parquet export → /app/data/pit_panel.parquet",
description=(
"FINRA × Alpaca(adjustment='all') 조인 패널을 parquet으로 저장.\n\n"
"`GET /finra/admin/export-pit-panel/download` 로 다운로드.\n\n"
"백그라운드 실행 — 전체 기간(2018-08~현재) 기준 수 분 ~ 10여 분 소요."
),
)
async def export_pit_panel(
background_tasks: BackgroundTasks,
start_date: date = Query(date(2018, 8, 1), description="시작일"),
end_date: date = Query(None, description="종료일 (기본: 오늘)"),
):
_end = end_date or date.today()
background_tasks.add_task(_run_pit_export, start_date.isoformat(), _end.isoformat())
return {
"status": "started",
"start_date": start_date.isoformat(),
"end_date": _end.isoformat(),
"output": _EXPORT_PATH,
"note": "완료 후 GET /finra/admin/export-pit-panel/download 로 다운로드.",
}
@router.get(
"/admin/export-pit-panel/download",
summary="pit_panel.parquet 다운로드",
)
async def download_pit_panel():
if not os.path.exists(_EXPORT_PATH):
raise HTTPException(
status_code=404,
detail="parquet 파일 없음. POST /finra/admin/export-pit-panel 먼저 실행.",
)
return FileResponse(
_EXPORT_PATH,
media_type="application/octet-stream",
filename="pit_panel.parquet",
)
@router.post(
"/admin/ingest",
response_model=IngestResponse,

@ -23,6 +23,7 @@ beautifulsoup4>=4.12.0
# SEC data processing - using direct API calls and EDGAR downloader
aiohttp>=3.8.0
pandas>=2.0.0,<3.0.0
pyarrow>=14.0.0 # parquet export (pit_panel)
numpy>=1.24.0,<3.0.0
python-dateutil>=2.8.0
sec-edgar-downloader>=5.0.0

Loading…
Cancel
Save