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, )