"""同步任务:mairui 历史分时 MA(日 K 级别)。 需求来源:mairui.club/hsdata 提供分时 K + MA 端点(/hsstock/history/{code}.{ex}/d/maN/)。 本任务计算每只股票日 K 级别 MA5/10/20/60 指标,存入 market_data.kline_stock_ma_daily。 算法: - 输入:market_data.kline_stock(已有 ~1170w 行日 K 线) - 按 stock_code 分组,对 close 做 rolling(5/10/20/60).mean() - 首部不足 N 天的行 → MA 留 NULL(不外推) - 全量重算(数据量 5213 只 × 6288 天 ≈ 3300w 行,pandas 处理 ~30s) 数据来源标注: - 本地计算 → source='local_kline_proxy' - mairui API(要付费 licence,当前不可用)→ 未来若升级可走 source='mairui' """ from __future__ import annotations import time from typing import Any import pandas as pd from sqlalchemy import create_engine from app.core.config import settings from app.core.db import ops as db_ops from app.core.sync.base import SyncTask from app.core.utils.logging import get_logger logger = get_logger("sync.mairui_ma_daily") class SyncMairuiMADaily(SyncTask): dataset_id = "mairui_ma_daily" # MA 窗口集合(标准日 K 级别常用 4 个) MA_WINDOWS = [5, 10, 20, 60] def _run( self, *, trigger_source: str = "manual", codes: list[str] | None = None, **kwargs, ) -> dict[str, Any]: t0 = time.time() logger.info("[mairui_ma_daily] 启动日 K MA 计算") # ── 1) 拉 kline_stock ── engine = create_engine(settings.pg_sqlalchemy_url()) df = pd.read_sql( 'SELECT stock_code, trade_date, "close" ' 'FROM market_data.kline_stock ORDER BY stock_code, trade_date', engine, ) if df.empty: return {"status": "error", "message": "kline_stock 为空"} # ── 2) 按 stock_code 分组算 rolling MA ── df["stock_code"] = df["stock_code"].astype(str) df["trade_date"] = pd.to_datetime(df["trade_date"], errors="coerce") df = df.dropna(subset=["trade_date"]) # 只过滤指定 codes(可选,用于增量) if codes: df = df[df["stock_code"].isin(codes)] # 关键步骤:每只股票独立 rolling out_pieces = [] for stock_code, group in df.groupby("stock_code", sort=False): g = group.sort_values("trade_date").copy() for w in self.MA_WINDOWS: g[f"ma{w}"] = g["close"].rolling(window=w, min_periods=w).mean() out_pieces.append(g) out = pd.concat(out_pieces, ignore_index=True) logger.info(f"[mairui_ma_daily] 计算完成: {len(out):,} 行, " f"{out['stock_code'].nunique()} 只, " f"{out['trade_date'].min().date()} ~ {out['trade_date'].max().date()}") # ── 3) 整理为 upsert 行 ── out["trade_date"] = out["trade_date"].dt.strftime("%Y-%m-%d") rows = [] for _, r in out.iterrows(): rows.append({ "stock_code": str(r["stock_code"]), "trade_date": r["trade_date"], "ma5": None if pd.isna(r["ma5"]) else float(r["ma5"]), "ma10": None if pd.isna(r["ma10"]) else float(r["ma10"]), "ma20": None if pd.isna(r["ma20"]) else float(r["ma20"]), "ma60": None if pd.isna(r["ma60"]) else float(r["ma60"]), "source": "local_kline_proxy", }) # ── 4) 分批 upsert(每批 5000 行,避免 SQL 太长)── BATCH = 5000 for i in range(0, len(rows), BATCH): db_ops.upsert_kline_stock_ma_daily_rows(rows[i:i + BATCH]) if (i // BATCH) % 10 == 0: self._progress( message=f"upsert {i + BATCH}/{len(rows)}", current=min(i + BATCH, len(rows)), total=len(rows), ) elapsed = round(time.time() - t0, 1) msg = ( f"MA {self.MA_WINDOWS} 共 {len(rows):,} 行, " f"{out['stock_code'].nunique()} 只, {elapsed}s" ) logger.info(f"[mairui_ma_daily] {msg}") return { "status": "ok", "message": msg, "rows": len(rows), "stocks": int(out["stock_code"].nunique()), "windows": self.MA_WINDOWS, "elapsed_sec": elapsed, }