diff --git a/ft-app/app/collector.py b/ft-app/app/collector.py index 235676c..6e6e1f9 100644 --- a/ft-app/app/collector.py +++ b/ft-app/app/collector.py @@ -1,9 +1,124 @@ """Data collector — fetch OHLCV from akshare and upsert into daily_bars.""" +import threading from datetime import date from app.database import SessionLocal from app.models import DailyBar, Contract from app.engine.lock_strategy import compute_amp_5d +_sync_progress = { + "running": False, + "contract": "", + "contract_done": 0, + "contract_total": 0, + "date_done": 0, + "date_total": 0, + "rows": 0, + "finished": False, +} +_lock = threading.Lock() + + +def get_sync_progress() -> dict: + with _lock: + return dict(_sync_progress) + + +def start_sync_positions_background(): + """Start sync_position_rankings in a background thread. Returns immediately.""" + with _lock: + if _sync_progress["running"]: + return False + _sync_progress.update({ + "running": True, + "finished": False, + "contract": "", + "contract_done": 0, + "contract_total": 0, + "date_done": 0, + "date_total": 0, + "rows": 0, + }) + + def _run(): + try: + _do_sync_with_progress() + finally: + with _lock: + _sync_progress["running"] = False + _sync_progress["finished"] = True + + threading.Thread(target=_run, daemon=True).start() + return True + + +def _do_sync_with_progress(): + from app.models import PositionRanking + + db = SessionLocal() + try: + active_contracts = ( + db.query(Contract).filter(Contract.is_active == True).all() + ) + with _lock: + _sync_progress["contract_total"] = len(active_contracts) + + for idx, c in enumerate(active_contracts): + with _lock: + _sync_progress.update({ + "contract": c.code, + "contract_done": idx, + "date_done": 0, + "date_total": 0, + }) + + code = c.code.upper() + + existing_dates = { + r[0] for r in + db.query(PositionRanking.date) + .filter(PositionRanking.contract_code == code) + .distinct() + .all() + } + + bar_dates = [ + r[0] for r in + db.query(DailyBar.date) + .filter(DailyBar.contract == code) + .order_by(DailyBar.date) + .all() + ] + + missing = [d for d in bar_dates if d not in existing_dates] + with _lock: + _sync_progress["date_total"] = len(missing) + + for date_idx, d in enumerate(missing): + date_str = d.strftime("%Y%m%d") + rankings = fetch_position_rankings(code, date_str) + inserted = 0 + for r in rankings: + db.add(PositionRanking( + contract_code=code, + institution=r["institution"], + data_type=r["data_type"], + date=d, + rank=r["rank"], + value=r["value"], + change=r["change"], + )) + inserted += 1 + db.flush() + with _lock: + _sync_progress["date_done"] = date_idx + 1 + _sync_progress["rows"] += inserted + + db.commit() + with _lock: + _sync_progress["contract_done"] = len(active_contracts) + finally: + db.close() + def fetch_contract_bars(contract_code: str, start_date: str | None = None) -> list[dict]: """Fetch daily OHLCV for a single contract from akshare. diff --git a/ft-app/app/routers/admin.py b/ft-app/app/routers/admin.py index d705e74..d6d1600 100644 --- a/ft-app/app/routers/admin.py +++ b/ft-app/app/routers/admin.py @@ -1,9 +1,12 @@ from fastapi import APIRouter, Depends, Form, Request -from fastapi.responses import HTMLResponse, RedirectResponse +from fastapi.responses import HTMLResponse, RedirectResponse, JSONResponse from sqlalchemy.orm import Session from app.database import get_db from app.models import Product, Contract, DailyBar, PositionRanking -from app.collector import sync_active_contracts, sync_one_contract, sync_position_rankings +from app.collector import ( + sync_active_contracts, sync_one_contract, sync_position_rankings, + start_sync_positions_background, get_sync_progress, +) router = APIRouter(prefix="/admin", tags=["admin"]) @@ -147,10 +150,15 @@ def sync_single(contract_code: str): @router.post("/sync-positions") def sync_positions(request: Request): - results = sync_position_rankings() - total = sum(results.values()) - print(f"[sync] Position rankings: {total} rows across {len(results)} contracts") - return RedirectResponse(f"/admin/?tab=sync&pos_synced={total}", status_code=303) + started = start_sync_positions_background() + if not started: + return JSONResponse({"error": "同步正在进行中,请等待完成"}) + return JSONResponse({"started": True}) + + +@router.get("/sync-status") +def sync_status(): + return JSONResponse(get_sync_progress()) @router.post("/sync/product/{product_id}") diff --git a/ft-app/app/templates/admin.html b/ft-app/app/templates/admin.html index 7baee2c..77e97cc 100644 --- a/ft-app/app/templates/admin.html +++ b/ft-app/app/templates/admin.html @@ -34,18 +34,24 @@
-
- -
+ {% if request.query_params.get('synced') %} ✓ 行情 {{ request.query_params.synced }} 条 {% endif %} - {% if request.query_params.get('pos_synced') %} - ✓ 持仓 {{ request.query_params.pos_synced }} 条 - {% endif %} + {{ total_contracts }} 个合约 + + {% if products %} {% for p in products %}
@@ -196,6 +202,64 @@