Files
market_sync/app/tasks/task_mairui_ma_daily.py
gao b351bd7595 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>
2026-07-03 17:36:49 +08:00

118 lines
4.4 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""同步任务: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,
}