feat: 新增 mairui 历史分时 MA 日 K 级别 sync task

work #06 (2026-07-03):user 请求加 mairui /hsdata 历史分时 MA 同步(日 K 级别)。

mairui 端点探测:/d/ma, /d/ma5/10/20, /15/ma, /30/ma, /60/ma 端点结构存在
但当前免费 licence 返 数据不存在;基础 K 线 (/d/n, /15/n, /30/n, /60/n) 正常。

策略:本地从 kline_stock 计算(pandas per-stock rolling),写新表
kline_stock_ma_daily,source=local_kline_proxy 标识本地派生。
mairui URL 留作未来升级 licence 后切 API 用。

变更:
- app/core/db/models.py: KlineStockMADaily ORM model
- app/core/db/ops.py: upsert_kline_stock_ma_daily_rows (批量 5000/批)
- app/tasks/task_mairui_ma_daily.py: SyncMairuiMADaily (全量重算)
- app/tasks/__init__.py: 注册到 TASKS dict
- app/core/sync/registry.py: SYNC_DEFINITION (sort_order=90, dep=kline_daily)
- app/core/scheduler/scheduler.py: schedule_mairui_ma_daily @ 16:30 + register_sync_jobs
- bin/daily_sync_check.py: SCHEDULE entry (window_end=17:00)

烟测:11,684,592 行, 5510 只, 1370s (23min),MA5/10/20/60 全部计算。
SH600519 样本:ma5=1194.18 ma10=1191.03 ma20=1211.94 ma60=1296.20 (2026-07-03)

未来优化(不在本 work):
- 增量模式(每日只算最近 1-2 天)→ 23min → 30s
- 升级 mairui licence 切到 /d/maN API
- 加 EMA / BOLL / KDJ 等其他指标

Co-Authored-By: Claude <noreply@anthropic.com>
This commit is contained in:
gao
2026-07-03 17:36:49 +08:00
parent 46ddf8e171
commit b351bd7595
8 changed files with 498 additions and 3 deletions
+118
View File
@@ -0,0 +1,118 @@
"""同步任务: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,
}