feat: add is_trading filter to symbols API + hardcode secret paths + backfill script
- API /api/v1/symbols: add is_trading query param + is_trading field in response - config.py: set default DB_PASSWORD_FILE=/run/secrets/db_pw.txt and ENCRYPTION_KEY_FILE=/run/secrets/enc_key.txt - scripts/backfill_candles.py: new script to fetch 500 historical candles per symbol/timeframe for all is_trading=true symbols - Cleared 2.88M non-trading candles from DB
This commit is contained in:
@@ -0,0 +1,156 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Historical candle backfill — fetch 500 candles for ALL timeframes for is_trading=true symbols."""
|
||||
import asyncio
|
||||
import logging
|
||||
import sys
|
||||
from datetime import datetime, timezone
|
||||
|
||||
from sqlalchemy import and_, select
|
||||
from sqlalchemy.dialects.postgresql import insert as pg_insert
|
||||
from sqlalchemy.orm import joinedload
|
||||
|
||||
sys.path.insert(0, "/app")
|
||||
|
||||
from app.database import async_session_factory
|
||||
from app.exchange.factory import factory as exchange_factory
|
||||
from app.models import Candle, Symbol, Exchange
|
||||
|
||||
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(message)s")
|
||||
logger = logging.getLogger("backfill")
|
||||
|
||||
TIMEFRAMES = ["15m", "30m", "1h", "4h", "1d", "1w", "1M"]
|
||||
FETCH_LIMIT = 500 # Max historical candles per symbol/timeframe
|
||||
CONCURRENT = 8 # Fetch 8 symbols concurrently
|
||||
SEM = asyncio.Semaphore(200) # Max concurrent API calls
|
||||
|
||||
|
||||
async def fetch_and_store(symbol: Symbol, adapter, sem: asyncio.Semaphore):
|
||||
"""Fetch all timeframes for one symbol and store in DB."""
|
||||
exchange_name = symbol.exchange.name
|
||||
all_values = []
|
||||
success_tfs = 0
|
||||
|
||||
async def _fetch_one_tf(tf: str):
|
||||
async with sem:
|
||||
try:
|
||||
candles = await adapter.fetch_ohlcv(
|
||||
symbol=symbol.symbol,
|
||||
timeframe=tf,
|
||||
limit=FETCH_LIMIT,
|
||||
)
|
||||
return tf, candles
|
||||
except Exception as e:
|
||||
err = str(e)
|
||||
if "BadSymbol" in err or "does not have market symbol" in err:
|
||||
logger.warning(" %s/%s not found on %s", symbol.symbol, tf, exchange_name)
|
||||
elif "1M" in err and "invalid" in err.lower():
|
||||
logger.debug(" %s: 1M not supported on %s", symbol.symbol, exchange_name)
|
||||
else:
|
||||
logger.error(" %s %s/%s: %s", exchange_name, symbol.symbol, tf, e)
|
||||
return tf, []
|
||||
|
||||
tasks = [_fetch_one_tf(tf) for tf in TIMEFRAMES]
|
||||
results = await asyncio.gather(*tasks)
|
||||
|
||||
for tf, candles in results:
|
||||
if candles:
|
||||
success_tfs += 1
|
||||
for c in candles:
|
||||
all_values.append({
|
||||
"symbol_id": symbol.id,
|
||||
"timeframe": c.timeframe,
|
||||
"timestamp": c.timestamp,
|
||||
"open": c.open,
|
||||
"high": c.high,
|
||||
"low": c.low,
|
||||
"close": c.close,
|
||||
"volume": c.volume,
|
||||
})
|
||||
|
||||
if all_values:
|
||||
async with async_session_factory() as db:
|
||||
stmt = pg_insert(Candle).values(all_values)
|
||||
stmt = stmt.on_conflict_do_nothing(
|
||||
index_elements=["symbol_id", "timeframe", "timestamp"]
|
||||
)
|
||||
await db.execute(stmt)
|
||||
await db.commit()
|
||||
|
||||
return success_tfs, len(all_values)
|
||||
|
||||
|
||||
async def main():
|
||||
# 1. Load all trading symbols
|
||||
async with async_session_factory() as db:
|
||||
query = (
|
||||
select(Symbol)
|
||||
.options(joinedload(Symbol.exchange))
|
||||
.join(Exchange, Exchange.id == Symbol.exchange_id)
|
||||
.where(
|
||||
and_(
|
||||
Symbol.is_trading == True,
|
||||
Symbol.is_active == True,
|
||||
Exchange.is_active == True,
|
||||
)
|
||||
)
|
||||
.order_by(Symbol.symbol)
|
||||
)
|
||||
result = await db.execute(query)
|
||||
symbols = list(result.scalars().all())
|
||||
|
||||
total_syms = len(symbols)
|
||||
logger.info("Backfilling %d symbols across 5 exchanges × %d TFs (limit=%d)",
|
||||
total_syms, len(TIMEFRAMES), FETCH_LIMIT)
|
||||
logger.info("Estimated: %d × %d × %d = ~%d candles max",
|
||||
total_syms, len(TIMEFRAMES), FETCH_LIMIT,
|
||||
total_syms * len(TIMEFRAMES) * FETCH_LIMIT)
|
||||
|
||||
# 2. Group by exchange
|
||||
by_exchange: dict[str, list] = {}
|
||||
for sym in symbols:
|
||||
by_exchange.setdefault(sym.exchange.name, []).append(sym)
|
||||
|
||||
total_candles = 0
|
||||
start_time = datetime.now(timezone.utc)
|
||||
|
||||
for ex_name, sym_list in by_exchange.items():
|
||||
logger.info("\n=== %s: %d symbols ===", ex_name.upper(), len(sym_list))
|
||||
try:
|
||||
adapter = exchange_factory.create(ex_name)
|
||||
except ValueError:
|
||||
logger.error("Unknown exchange: %s", ex_name)
|
||||
continue
|
||||
|
||||
# Process symbols concurrently (CONCURRENT at a time)
|
||||
sem = asyncio.Semaphore(CONCURRENT * len(TIMEFRAMES))
|
||||
|
||||
batch_tasks = []
|
||||
for sym in sym_list:
|
||||
batch_tasks.append(fetch_and_store(sym, adapter, sem))
|
||||
|
||||
results = await asyncio.gather(*batch_tasks, return_exceptions=True)
|
||||
|
||||
ex_candles = 0
|
||||
ex_success = 0
|
||||
for i, r in enumerate(results):
|
||||
if isinstance(r, Exception):
|
||||
logger.error(" Error: %s", r)
|
||||
else:
|
||||
tfs_ok, n = r
|
||||
if tfs_ok > 0:
|
||||
ex_success += 1
|
||||
ex_candles += n
|
||||
|
||||
total_candles += ex_candles
|
||||
elapsed = (datetime.now(timezone.utc) - start_time).total_seconds()
|
||||
logger.info(" %s DONE: %d/%d symbols, %d candles in %.0fs",
|
||||
ex_name.upper(), ex_success, len(sym_list), ex_candles, elapsed)
|
||||
|
||||
elapsed = (datetime.now(timezone.utc) - start_time).total_seconds()
|
||||
logger.info("\n=== BACKFILL COMPLETE ===")
|
||||
logger.info("Total: %d symbols, %d candles, %.0fs (%.1f min)",
|
||||
total_syms, total_candles, elapsed, elapsed / 60)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(main())
|
||||
Reference in New Issue
Block a user