From af26df8895c51e41c198674ed5c4799e67591547 Mon Sep 17 00:00:00 2001 From: I Luk Kim Date: Tue, 14 Apr 2026 02:28:04 -0700 Subject: [PATCH] =?UTF-8?q?feat:=20AlpacaPriceData=EC=97=90=20interval=20?= =?UTF-8?q?=EC=BB=AC=EB=9F=BC=20=EC=B6=94=EA=B0=80=20=E2=80=94=20=EC=9D=BC?= =?UTF-8?q?=EB=B4=89/=EB=B6=84=EB=B4=89=20DB=20=EC=B6=A9=EB=8F=8C=20?= =?UTF-8?q?=ED=95=B4=EA=B2=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - alpaca_price_data 테이블에 interval 컬럼 추가 (VARCHAR, default '1d') - unique constraint: (ticker, date) → (ticker, date, interval) - get_or_fetch_multi_bars(): interval 필터로 coverage 체크 및 DB 조회 - insert rows에 interval 포함 - asyncpg param limit 대비 CHUNK 3000→2900 (파라미터 11개/행) - alembic 마이그레이션: h9b0c1d2e3f4 Co-Authored-By: Claude Sonnet 4.6 --- ...b0c1d2e3f4_add_interval_to_alpaca_price.py | 49 +++++++++++++++++++ app/models/alpaca_price.py | 6 ++- app/services/alpaca_price_service.py | 9 ++-- 3 files changed, 59 insertions(+), 5 deletions(-) create mode 100644 alembic/versions/h9b0c1d2e3f4_add_interval_to_alpaca_price.py diff --git a/alembic/versions/h9b0c1d2e3f4_add_interval_to_alpaca_price.py b/alembic/versions/h9b0c1d2e3f4_add_interval_to_alpaca_price.py new file mode 100644 index 0000000..d292679 --- /dev/null +++ b/alembic/versions/h9b0c1d2e3f4_add_interval_to_alpaca_price.py @@ -0,0 +1,49 @@ +"""add interval column to alpaca_price_data + +Revision ID: h9b0c1d2e3f4 +Revises: g8a9b0c1d2e3 +Create Date: 2026-04-14 +""" +from typing import Sequence, Union + +from alembic import op +import sqlalchemy as sa + +revision: str = "h9b0c1d2e3f4" +down_revision: Union[str, Sequence[str], None] = "g8a9b0c1d2e3" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + # 1. Add interval column (default '1d' for existing rows) + op.add_column( + "alpaca_price_data", + sa.Column("interval", sa.String(10), nullable=False, server_default="1d"), + ) + + # 2. Drop old unique constraint and index + op.drop_constraint("uq_alpaca_price_data", "alpaca_price_data", type_="unique") + op.drop_index("idx_alpaca_price_ticker_date", table_name="alpaca_price_data") + + # 3. Create new unique constraint and index including interval + op.create_unique_constraint( + "uq_alpaca_price_data", "alpaca_price_data", ["ticker", "date", "interval"] + ) + op.create_index( + "idx_alpaca_price_ticker_date", "alpaca_price_data", ["ticker", "date", "interval"] + ) + + +def downgrade() -> None: + op.drop_constraint("uq_alpaca_price_data", "alpaca_price_data", type_="unique") + op.drop_index("idx_alpaca_price_ticker_date", table_name="alpaca_price_data") + + op.create_unique_constraint( + "uq_alpaca_price_data", "alpaca_price_data", ["ticker", "date"] + ) + op.create_index( + "idx_alpaca_price_ticker_date", "alpaca_price_data", ["ticker", "date"] + ) + + op.drop_column("alpaca_price_data", "interval") diff --git a/app/models/alpaca_price.py b/app/models/alpaca_price.py index 961d21e..3d41806 100644 --- a/app/models/alpaca_price.py +++ b/app/models/alpaca_price.py @@ -25,12 +25,14 @@ class AlpacaPriceData(Base): vwap = Column(Float, nullable=True) # Volume-weighted average price (Alpaca-specific) trade_count = Column(Integer, nullable=True) # Number of trades (Alpaca-specific) + interval = Column(String(10), nullable=False, default='1d') # e.g. "1d", "5m", "1h" + # Metadata data_source = Column(String(50), default='ALPACA') created_at = Column(TIMESTAMP(timezone=True), default=lambda: datetime.now(timezone.utc)) updated_at = Column(TIMESTAMP(timezone=True), default=lambda: datetime.now(timezone.utc), onupdate=lambda: datetime.now(timezone.utc)) __table_args__ = ( - UniqueConstraint('ticker', 'date', name='uq_alpaca_price_data'), - Index('idx_alpaca_price_ticker_date', 'ticker', 'date'), + UniqueConstraint('ticker', 'date', 'interval', name='uq_alpaca_price_data'), + Index('idx_alpaca_price_ticker_date', 'ticker', 'date', 'interval'), ) diff --git a/app/services/alpaca_price_service.py b/app/services/alpaca_price_service.py index b5f4113..ab8698f 100644 --- a/app/services/alpaca_price_service.py +++ b/app/services/alpaca_price_service.py @@ -114,6 +114,7 @@ class AlpacaPriceService: .where( and_( AlpacaPriceData.ticker.in_(upper_tickers), + AlpacaPriceData.interval == interval, AlpacaPriceData.date >= start_dt, AlpacaPriceData.date <= end_dt, ) @@ -151,6 +152,7 @@ class AlpacaPriceService: rows.append({ "ticker": original, "date": bar_dt, + "interval": interval, "open": float(bar.get("o", 0)), "high": float(bar.get("h", 0)), "low": float(bar.get("l", 0)), @@ -162,8 +164,8 @@ class AlpacaPriceService: }) if rows: - # Chunk to stay under asyncpg's 32767-param limit (~10 params/row → 3000/chunk) - CHUNK = 3000 + # Chunk to stay under asyncpg's 32767-param limit (~11 params/row → 2900/chunk) + CHUNK = 2900 for i in range(0, len(rows), CHUNK): stmt = pg_insert(AlpacaPriceData).values(rows[i : i + CHUNK]) stmt = stmt.on_conflict_do_nothing(constraint="uq_alpaca_price_data") @@ -173,12 +175,13 @@ class AlpacaPriceService: f"Alpaca multi-bars: stored {len(rows)} rows for {len(need_fetch)} tickers" ) - # Read all data back from DB + # Read back from DB — filtered by interval to avoid mixing 1d/5m/etc. result = await db.execute( select(AlpacaPriceData) .where( and_( AlpacaPriceData.ticker.in_(upper_tickers), + AlpacaPriceData.interval == interval, AlpacaPriceData.date >= start_dt, AlpacaPriceData.date <= end_dt, )