From ca98967af002f347c7ed2c1cc799be9ff2b23bdf Mon Sep 17 00:00:00 2001 From: fish Date: Sat, 25 Jul 2026 13:53:44 +0800 Subject: [PATCH] =?UTF-8?q?=E7=A7=BB=E9=99=A4=E5=85=A8=E9=83=A8=E5=90=88?= =?UTF-8?q?=E7=BA=A6=E5=90=8C=E6=AD=A5=E6=8C=89=E9=92=AE=EF=BC=8C=E6=94=B9?= =?UTF-8?q?=E4=B8=BA=E5=8D=95=E5=90=88=E7=BA=A6=E5=90=8E=E5=8F=B0=E5=90=8C?= =?UTF-8?q?=E6=AD=A5=E5=B8=A6=E8=BF=9B=E5=BA=A6=E6=9D=A1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 4.7 --- ft-app/app/collector.py | 150 +++++++++++++++++--------------- ft-app/app/routers/admin.py | 29 +++--- ft-app/app/templates/admin.html | 109 ++++++++++------------- 3 files changed, 135 insertions(+), 153 deletions(-) diff --git a/ft-app/app/collector.py b/ft-app/app/collector.py index 6e6e1f9..500de38 100644 --- a/ft-app/app/collector.py +++ b/ft-app/app/collector.py @@ -5,119 +5,127 @@ 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, -} +_progress = {"running": False, "label": "", "done": 0, "total": 0, "finished": False} _lock = threading.Lock() -def get_sync_progress() -> dict: +def get_progress() -> dict: with _lock: - return dict(_sync_progress) + return dict(_progress) -def start_sync_positions_background(): - """Start sync_position_rankings in a background thread. Returns immediately.""" +def _start_bg(target, label): with _lock: - if _sync_progress["running"]: + if _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, - }) + _progress.update(running=True, finished=False, label=label, done=0, total=0) def _run(): try: - _do_sync_with_progress() + target() finally: with _lock: - _sync_progress["running"] = False - _sync_progress["finished"] = True + _progress["running"] = False + _progress["finished"] = True + _progress["done"] = _progress["total"] threading.Thread(target=_run, daemon=True).start() return True -def _do_sync_with_progress(): - from app.models import PositionRanking +def sync_contract_bars_bg(contract_code: str) -> bool: + """Sync bars for one contract in background. Returns True if started.""" + code = contract_code.upper() - db = SessionLocal() - try: - active_contracts = ( - db.query(Contract).filter(Contract.is_active == True).all() - ) - with _lock: - _sync_progress["contract_total"] = len(active_contracts) + def _run(): + db = SessionLocal() + try: + latest = ( + db.query(DailyBar.date) + .filter(DailyBar.contract == code) + .order_by(DailyBar.date.desc()) + .first() + ) + start_date = latest[0].isoformat() if latest else None + bars = fetch_contract_bars(code, start_date) - for idx, c in enumerate(active_contracts): with _lock: - _sync_progress.update({ - "contract": c.code, - "contract_done": idx, - "date_done": 0, - "date_total": 0, - }) + _progress["total"] = len(bars) - code = c.code.upper() + inserted = 0 + min_date = None + for i, bar in enumerate(bars): + existing = ( + db.query(DailyBar) + .filter(DailyBar.contract == code, DailyBar.date == bar["date"]) + .first() + ) + if not existing: + db.add(DailyBar( + contract=code, date=bar["date"], + open=bar["open"], close=bar["close"], + high=bar["high"], low=bar["low"], + )) + inserted += 1 + if min_date is None or bar["date"] < min_date: + min_date = bar["date"] + with _lock: + _progress["done"] = i + 1 + if inserted > 0: + db.flush() + _recompute_amp(db, code, from_date=min_date) + db.commit() + finally: + db.close() + + return _start_bg(_run, f"同步行情 {code}") + + +def sync_positions_bg(contract_code: str) -> bool: + """Sync position rankings for one contract in background. Returns True if started.""" + from app.models import PositionRanking + code = contract_code.upper() + + def _run(): + db = SessionLocal() + try: existing_dates = { r[0] for r in db.query(PositionRanking.date) .filter(PositionRanking.contract_code == code) - .distinct() - .all() + .distinct().all() } - bar_dates = [ r[0] for r in db.query(DailyBar.date) .filter(DailyBar.contract == code) - .order_by(DailyBar.date) - .all() + .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 + with _lock: + _progress["total"] = len(missing) + + inserted = 0 + for i, d in enumerate(missing): + rankings = fetch_position_rankings(code, d.strftime("%Y%m%d")) 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"], + 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 + _progress["done"] = i + 1 - db.commit() - with _lock: - _sync_progress["contract_done"] = len(active_contracts) - finally: - db.close() + db.commit() + finally: + db.close() + + return _start_bg(_run, f"同步持仓 {code}") def fetch_contract_bars(contract_code: str, start_date: str | None = None) -> list[dict]: diff --git a/ft-app/app/routers/admin.py b/ft-app/app/routers/admin.py index d6d1600..7103493 100644 --- a/ft-app/app/routers/admin.py +++ b/ft-app/app/routers/admin.py @@ -4,8 +4,7 @@ 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, - start_sync_positions_background, get_sync_progress, + sync_one_contract, sync_contract_bars_bg, sync_positions_bg, get_progress, ) router = APIRouter(prefix="/admin", tags=["admin"]) @@ -134,31 +133,27 @@ def delete_product(product_id: int, db: Session = Depends(get_db)): return RedirectResponse("/admin/?tab=product", status_code=303) -@router.post("/sync") -def sync_all(request: Request): - results = sync_active_contracts() - total = sum(results.values()) - print(f"[sync] Synced {total} bars across {len(results)} contracts: {results}") - return RedirectResponse(f"/admin/?tab=sync&synced={total}", status_code=303) - - @router.post("/sync/{contract_code}") def sync_single(contract_code: str): count = sync_one_contract(contract_code.upper()) return RedirectResponse(f"/admin/?tab=sync&synced={count}", status_code=303) -@router.post("/sync-positions") -def sync_positions(request: Request): - started = start_sync_positions_background() - if not started: - return JSONResponse({"error": "同步正在进行中,请等待完成"}) - return JSONResponse({"started": True}) +@router.post("/sync/{contract_code}/bg") +def sync_single_bg(contract_code: str): + ok = sync_contract_bars_bg(contract_code.upper()) + return JSONResponse({"ok": ok}) + + +@router.post("/sync-positions/{contract_code}/bg") +def sync_positions_bg_ep(contract_code: str): + ok = sync_positions_bg(contract_code.upper()) + return JSONResponse({"ok": ok}) @router.get("/sync-status") def sync_status(): - return JSONResponse(get_sync_progress()) + return JSONResponse(get_progress()) @router.post("/sync/product/{product_id}") diff --git a/ft-app/app/templates/admin.html b/ft-app/app/templates/admin.html index 77e97cc..1b17d92 100644 --- a/ft-app/app/templates/admin.html +++ b/ft-app/app/templates/admin.html @@ -29,26 +29,17 @@ {# ═══════════════════ Tab: 数据同步 ═══════════════════ #}
-
- 全部合约 -
- -
- - {% if request.query_params.get('synced') %} - ✓ 行情 {{ request.query_params.synced }} 条 - {% endif %} - - {{ total_contracts }} 个合约 +
+ 共 {{ total_contracts }} 个合约
-