diff --git a/backend/app/api/kline.py b/backend/app/api/kline.py index be3d12f..1bc5d49 100644 --- a/backend/app/api/kline.py +++ b/backend/app/api/kline.py @@ -579,7 +579,10 @@ def refresh_views(request: Request): @router.post("/sync_minute") async def sync_minute(request: Request): - """手动触发分钟 K 同步(全市场)。返回 pipeline job_id 可轮询进度。""" + """手动触发分钟 K 同步(全市场)。返回 pipeline job_id 可轮询进度。 + + body 可选: { "days": int } — 指定拉取天数 (不传则用偏好设置)。 + """ import asyncio from app.services.pipeline_jobs import job_store, release_run_slot, try_acquire_run_slot @@ -594,6 +597,16 @@ async def sync_minute(request: Request): if not _minute_allowed(capset): raise HTTPException(status_code=403, detail="需要 Pro+ 权限") + # 可选 body: { "days": int, "extend": bool } + # days: 拉取天数; extend: 向前扩展模式 (从最早数据往前补) + body = {} + try: + body = await request.json() + except Exception: # noqa: BLE001 + pass + override_days = body.get("days") + extend_flag = body.get("extend") + job_id, is_new = job_store.create() if not is_new: return {"status": "reused", "job_id": job_id} @@ -622,10 +635,21 @@ async def sync_minute(request: Request): pass progress("sync_minute", 10, f"标的池 {len(universe)} 只") - days = get_minute_sync_days() + days = override_days if override_days else get_minute_sync_days() + # extend=1 → 向前扩展; days>=365 也自动向前扩展 + extend_backward = bool(extend_flag) or days >= 365 + + def _on_chunk(done: int, total: int, seg_label: str) -> None: + # 进度映射: 10% (标的池解析完) → 95%, 留 5% 给写入+刷新 + pct = 10 + int((done / max(total, 1)) * 85) + progress("sync_minute", pct, f"拉取分钟K… {done}/{total} 批 [{seg_label}]") def _run(): - return kline_sync.sync_and_persist_minute(universe, repo, capset, days=days) + return kline_sync.sync_and_persist_minute( + universe, repo, capset, days=days, + extend_backward=extend_backward, + on_chunk_done=_on_chunk, + ) written = await loop.run_in_executor(_long_task_executor, _run) @@ -646,6 +670,79 @@ async def sync_minute(request: Request): return {"status": "started", "job_id": job_id} +@router.post("/sync_minute_single") +async def sync_minute_single(request: Request, body: dict): + """手动拉取单只股票的分钟K并落库 (前复权)。 + + body: { "symbol": "000001.SZ" } + 用于个股分时图"获取数据"按钮: 本地无数据时单独拉取并持久化。 + """ + from app.services.preferences import get_minute_sync_days + from app.tickflow.capabilities import Cap + + symbol = body.get("symbol", "").strip() + if not symbol: + raise HTTPException(status_code=400, detail="symbol 不能为空") + + repo = request.app.state.repo + capset = request.app.state.capabilities + + if not _minute_allowed(capset): + raise HTTPException(status_code=403, detail="需要 Pro+ 权限") + + days = get_minute_sync_days() + loop = asyncio.get_event_loop() + + def _run(): + return kline_sync.sync_and_persist_minute([symbol], repo, capset, days=days) + + written = await loop.run_in_executor(_long_task_executor, _run) + + # 刷新视图 + from app.jobs.daily_pipeline import _refresh_single_view + _refresh_single_view(repo, "kline_minute") + + return {"status": "ok", "symbol": symbol, "rows": written} + + +@router.post("/clear_minute") +async def clear_minute(request: Request): + """清空全部分钟K数据 (仅 kline_minute, 不影响其他数据)。 + + 删除 data/kline_minute/ 下所有分区 parquet, 刷新视图。 + 需二次确认: body { "confirm": true }。 + """ + import shutil + + body = await request.json() if request.method == "POST" else {} + if not body.get("confirm"): + raise HTTPException(status_code=400, detail="需传 confirm: true 以确认清空") + + repo = request.app.state.repo + minute_dir = repo.store.data_dir / "kline_minute" + + # 统计待删除行数 (用于返回) + removed = 0 + if minute_dir.exists(): + try: + result = repo.db.execute("SELECT COUNT(*) AS cnt FROM kline_minute").fetchone() + removed = result[0] if result else 0 + except Exception: # noqa: BLE001 + pass + # 仅删 kline_minute 目录, 绝不触碰其他目录 + shutil.rmtree(minute_dir, ignore_errors=True) + + # 刷新视图 (重建空视图) + from app.jobs.daily_pipeline import _refresh_single_view + _refresh_single_view(repo, "kline_minute") + + from app.api.data import invalidate_storage_cache + invalidate_storage_cache() + + logger.info("minute K cleared: %d rows removed", removed) + return {"status": "ok", "removed": removed} + + @router.post("/extend_history") async def extend_history(request: Request): """向前扩展历史日K数据 — 独立于盘后管道。 @@ -886,195 +983,3 @@ async def rebuild_enriched(request: Request): import concurrent.futures as _cf _long_task_executor = _cf.ThreadPoolExecutor(max_workers=2, thread_name_prefix="long-task") - -@router.post("/extend_minute_history") -async def extend_minute_history(request: Request): - """向前扩展分钟K历史数据 — 仅拉数据,不做任何后续处理。 - - body: { "value": int, "unit": "day"|"month" } - - day 单位:1~15 天(所有有分钟K权限的套餐可用) - - month 单位:1~6 月(每月按 30 天计,即最多 180 天)—— 仅 Expert+ 可用 - 返回 job_id,可轮询 /api/pipeline/jobs 查看进度。 - """ - import asyncio - import traceback as _tb - try: - body = await request.json() - value = body.get("value") - unit = body.get("unit", "day") - if not value or value <= 0: - raise HTTPException(status_code=400, detail="value 必须为正整数") - if unit not in ("day", "month"): - raise HTTPException(status_code=400, detail="unit 只支持 day/month") - - repo = request.app.state.repo - capset = request.app.state.capabilities - - from app.tickflow.capabilities import Cap - if not _minute_allowed(capset): - raise HTTPException(status_code=403, detail="需要 Pro+ 权限 (batch minute K-line)") - - # month 单位(按月扩展更长的分钟K历史)仅 Expert+ 开放;Pro 仅可用 day - if unit == "month": - from app.tickflow.policy import tier_label - base_tier = tier_label().split()[0].split("+")[0].strip().lower() - if base_tier != "expert": - raise HTTPException( - status_code=403, - detail="按月扩展分钟K历史需要 Expert 及以上套餐", - ) - - # 计算天数上限:day 最多 15 天;month 最多 6 月(180 天) - from datetime import timedelta - if unit == "month": - total_days = min(value * 30, 180) - else: - total_days = min(value, 15) - - if total_days <= 0: - raise HTTPException(status_code=400, detail="扩展范围无效") - - from app.services.pipeline_jobs import job_store, release_run_slot, try_acquire_run_slot - from app.api.data import invalidate_storage_cache - - job_id, is_new = job_store.create() - if not is_new: - return {"status": "reused", "job_id": job_id} - - async def task() -> None: - if not try_acquire_run_slot(): - job_store.fail(job_id, "已有数据任务在运行(或上一次任务卡死未结束),请稍后再试") - return - loop = asyncio.get_event_loop() - - def progress(stage: str, pct: int, msg: str, - stage_pct: int | None = None, skip_log: bool = False) -> None: - job_store.progress(job_id, stage, pct, msg, - stage_pct=stage_pct, skip_log=skip_log) - - try: - job_store.start(job_id) - # 获取当前最早日期 - earliest = repo.earliest_minute_date() - if not earliest: - # 本地无分钟K数据 → 以今天为基准往前获取 - from datetime import date as _date - latest = _date.today() - else: - latest = earliest - - new_start = latest - timedelta(days=total_days) - if new_start >= latest: - job_store.fail(job_id, "扩展范围无效") - invalidate_storage_cache() - return - - start_str = new_start.strftime("%Y-%m-%d") - end_str = latest.strftime("%Y-%m-%d") - - progress("extend_minute", 5, "解析标的池…") - universe = _resolve_minute_universe(capset, repo) - progress("extend_minute", 8, f"标的池: {len(universe)} 只") - - from app.tickflow.capabilities import Cap - from app.tickflow.rate_limits import resolve_limit - - limit = resolve_limit( - capset, - Cap.KLINE_MINUTE_BATCH, - default_batch=100, - default_rpm=30, - default_rpm_when_unset=False, - ) - - def _run(): - """全部在 executor 线程里完成,避免阻塞事件循环。""" - from app.services.kline_sync import sync_minute_batch - from datetime import datetime as _dt - - def _chunk(cur: int, tot: int) -> None: - progress("extend_minute", 8 + int(85 * cur / tot), - f"分钟K 批次 {cur}/{tot}", stage_pct=int(100 * cur / tot), skip_log=True) - - df = sync_minute_batch( - universe, - start_time=_dt.combine(new_start, _dt.min.time()), - end_time=_dt.combine(latest, _dt.min.time()), - batch_size=limit.batch, rpm=limit.rpm, - on_chunk_done=_chunk, - ) - - written = 0 - day_count = 0 - if not df.is_empty(): - import polars as pl - df = df.with_columns(pl.col("datetime").dt.date().alias("_trade_date")) - for day_df in df.partition_by("_trade_date"): - trade_date = day_df["_trade_date"][0] - out = repo.store.data_dir / "kline_minute" / f"date={trade_date}" / "part.parquet" - out.parent.mkdir(parents=True, exist_ok=True) - if out.exists(): - existing_df = pl.read_parquet(out) - if "datetime" in existing_df.columns: - existing_df = existing_df.filter(pl.col("datetime").is_not_null()) - day_df = pl.concat([existing_df, day_df.drop("_trade_date")]).unique( - subset=["symbol", "datetime"], keep="last", - ) - else: - day_df = day_df.drop("_trade_date") - day_df = day_df.sort("symbol", "datetime") - from app.services.kline_sync import _atomic_write_parquet - _atomic_write_parquet(day_df, out) - written += day_df.height - day_count += 1 - - # 刷新视图 - d = repo.store.data_dir.as_posix() - try: - repo.db.execute( - f"CREATE OR REPLACE VIEW kline_minute AS " - f"SELECT * FROM read_parquet('{d}/kline_minute/**/*.parquet', union_by_name=true)" - ) - except Exception: - pass - return written, day_count - - progress("extend_minute", 10, f"获取分钟K [{start_str} ~ {end_str}]…") - written, day_count = await loop.run_in_executor(_long_task_executor, _run) - - progress("extend_minute", 95, f"分钟K 完成,{day_count} 天") - job_store.succeed(job_id, { - "minute_days": day_count, - "universe_size": len(universe), - "earliest_before": (earliest or latest).isoformat(), - "earliest_after": new_start.isoformat(), - }) - invalidate_storage_cache() - except Exception as e: - logger.exception("extend_minute_history failed: job_id=%s", job_id) - job_store.fail(job_id, str(e)) - invalidate_storage_cache() - finally: - release_run_slot() - - asyncio.create_task(task()) - return {"status": "started", "job_id": job_id} - except HTTPException: - raise - except Exception as e: - logger.error("extend_minute_history error: %s\n%s", e, _tb.format_exc()) - raise HTTPException(status_code=500, detail=str(e)) from e - - -def _resolve_minute_universe(capset, repo) -> list[str]: - """分钟K标的池解析。""" - from app.tickflow.capabilities import Cap - if capset.has(Cap.KLINE_MINUTE_BATCH): - try: - from app.tickflow.pools import get_pool - all_a = get_pool("CN_Equity_A", refresh=True) - if all_a: - return sorted(all_a) - except Exception: - pass - return [] diff --git a/backend/app/jobs/daily_pipeline.py b/backend/app/jobs/daily_pipeline.py index 21cc113..93b1dda 100644 --- a/backend/app/jobs/daily_pipeline.py +++ b/backend/app/jobs/daily_pipeline.py @@ -489,9 +489,10 @@ def run_now( emit("sync_minute", 90, f"获取分钟K [{minute_start} ~ {today}]…") logger.info("sync_minute: [%s ~ %s] start", minute_start, today) minute_symbols = _resolve_minute_symbols(capset) - def _minute_chunk_progress(cur: int, tot: int) -> None: + def _minute_chunk_progress(cur: int, tot: int, seg_label: str = "") -> None: emit("sync_minute", 90 + int(3 * cur / tot), - f"分钟K 批次 {cur}/{tot}", stage_pct=int(100 * cur / tot), skip_log=True) + f"分钟K 批次 {cur}/{tot}" + (f" [{seg_label}]" if seg_label else ""), + stage_pct=int(100 * cur / tot), skip_log=True) written_minute = kline_sync.sync_and_persist_minute( minute_symbols, repo, capset, days=minute_days, on_chunk_done=_minute_chunk_progress, diff --git a/backend/app/services/kline_sync.py b/backend/app/services/kline_sync.py index 0501836..7bee346 100644 --- a/backend/app/services/kline_sync.py +++ b/backend/app/services/kline_sync.py @@ -506,13 +506,16 @@ def sync_minute_batch( count: int | None = None, batch_size: int | None = None, rpm: int | None = None, - on_chunk_done: Callable[[int, int], None] | None = None, + on_chunk_done: Callable[[int, int, str], None] | None = None, ) -> pl.DataFrame: """批量拉取多股分钟 K。 优先使用 start_time / end_time 区间, 确保所有标的覆盖同一时间段。 count 仅作为 fallback 保留。 on_chunk_done(current, total) 每个 chunk 完成后回调。 + + TickFlow count 上限 10000 根/股, 1 天 240 根 → 单次最多约 41 天。 + 当区间超过 35 天时自动按月 (30 天) 分段拉取, 拼接结果。 """ # 自定义数据源分流: minute provider provider_name = preferences.get_minute_data_provider() @@ -526,37 +529,63 @@ def sync_minute_batch( # 未配置 minute → 回退 TickFlow tf = get_client() + + # TickFlow count 上限 10000 根/股, 1 天 240 根 → 单次最多约 41 个交易日。 + # 按 41 交易日 (≈57 自然日) 分段: ≤41 交易日的区间只产生 1 段 (单次拉满)。 + SEG_CHUNK = timedelta(days=57) # 41 交易日 × 7/5 ≈ 57 自然日 + time_segments: list[tuple[datetime, datetime]] = [] + if start_time and end_time: + seg_start = start_time + while seg_start < end_time: + seg_end = min(seg_start + SEG_CHUNK, end_time) + time_segments.append((seg_start, seg_end)) + seg_start = seg_end + else: + time_segments = [(None, None)] # fallback: 用 count 模式 + + total_steps = len(time_segments) * len(chunked(symbols, batch_size)) + step = 0 out: list[pl.DataFrame] = [] - chunks = chunked(symbols, batch_size) - for i, chunk in enumerate(chunks): - sleep_between_batches(i, rpm) - try: - if start_time and end_time: - raw = tf.klines.batch( - chunk, period="1m", - start_time=_datetime_to_ms(start_time), - end_time=_datetime_to_ms(end_time), - count=10000, - as_dataframe=True, show_progress=False, - ) - else: - raw = tf.klines.batch(chunk, period="1m", count=count or 1200, - as_dataframe=True, show_progress=False) - except Exception as e: # noqa: BLE001 - logger.warning("minute batch fetch failed for %d symbols: %s", len(chunk), e) - continue + for seg_idx, (seg_start, seg_end) in enumerate(time_segments): + # 当前的日期段描述 (供进度展示) + if seg_start and seg_end: + seg_label = f"{seg_start.strftime('%m-%d')}~{seg_end.strftime('%m-%d')}" + else: + seg_label = "最新" + seg_total = len(time_segments) + chunks = chunked(symbols, batch_size) + for i, chunk in enumerate(chunks): + sleep_between_batches(step, rpm) + step += 1 + try: + if seg_start and seg_end: + raw = tf.klines.batch( + chunk, period="1m", + start_time=_datetime_to_ms(seg_start), + end_time=_datetime_to_ms(seg_end), + count=10000, + adjust="forward", + as_dataframe=True, show_progress=False, + ) + else: + raw = tf.klines.batch(chunk, period="1m", count=count or 1200, + adjust="forward", + as_dataframe=True, show_progress=False) + except Exception as e: # noqa: BLE001 + logger.warning("minute batch fetch failed for %d symbols: %s", len(chunk), e) + continue - if isinstance(raw, dict): - for sym, sub in raw.items(): - if sub is None or len(sub) == 0: - continue - out.append(_normalize_minute(sub, default_symbol=sym)) - elif raw is not None and len(raw) > 0: - out.append(_normalize_minute(raw)) + if isinstance(raw, dict): + for sym, sub in raw.items(): + if sub is None or len(sub) == 0: + continue + out.append(_normalize_minute(sub, default_symbol=sym)) + elif raw is not None and len(raw) > 0: + out.append(_normalize_minute(raw)) - if on_chunk_done: - on_chunk_done(i + 1, len(chunks)) + if on_chunk_done: + on_chunk_done(step, total_steps, seg_label) if not out: return pl.DataFrame() @@ -575,6 +604,7 @@ def fetch_minute_single(symbol: str, trade_date: date) -> pl.DataFrame: start_time=_datetime_to_ms(start_time), end_time=_datetime_to_ms(end_time), count=10000, + adjust="forward", as_dataframe=True, show_progress=False, ) except Exception as e: @@ -618,6 +648,20 @@ def _latest_minute_datetime(repo: KlineRepository) -> datetime | None: return None +def _earliest_minute_datetime(repo: KlineRepository) -> datetime | None: + """本地分钟 K 数据的最早时间 (用于向前扩展的起点)。""" + try: + res = repo.execute_one("SELECT min(datetime) FROM kline_minute") + if res and res[0]: + d = res[0] + if isinstance(d, datetime): + return d + return datetime.fromisoformat(str(d)) + except Exception: # noqa: BLE001 + pass + return None + + def _cleanup_null_datetime_minute(repo: KlineRepository) -> None: """检测并清除 datetime 全为 null 的旧版分钟 K 数据(迁移用)。""" minute_dir = repo.store.data_dir / "kline_minute" @@ -703,9 +747,10 @@ def sync_and_persist_minute( repo: KlineRepository, capset: CapabilitySet, days: int = 5, - on_chunk_done: Callable[[int, int], None] | None = None, + on_chunk_done: Callable[[int, int, str], None] | None = None, + extend_backward: bool = False, ) -> int: - """同步分钟 K 并存到 Parquet(仅 raw,不前复权)。返回写入行数。 + """同步分钟 K 并存到 Parquet(前复权价格, SDK 端 adjust=qfq)。返回写入行数。 使用 start_time / end_time 区间拉取, 确保所有标的覆盖同一时间段。 on_chunk_done(current, total) 每个 chunk 完成后回调。 @@ -729,13 +774,28 @@ def sync_and_persist_minute( now = datetime.now() - # 计算时间区间: 首次拉取回溯 N 天, 增量从最后数据时间开始 - last_dt = _latest_minute_datetime(repo) - if last_dt: - start_time = last_dt + if extend_backward: + # 向前扩展模式: 从本地最早数据往前补, 叠加已有数据避免缺口。 + earliest_dt = _earliest_minute_datetime(repo) + # 按交易日换算自然日 (7/5 系数) + # ≤41 交易日: 不加余量, 确保落在单段 (57 自然日) 内 → 单次拉满 + # >41 交易日: +10 天余量覆盖节假日 + calendar_days = int(days * 7 / 5) + (10 if days > 41 else 0) + if earliest_dt: + end_time = earliest_dt + start_time = end_time - timedelta(days=calendar_days) + else: + # 本地无数据 → 从今天往前拉 + start_time = now - timedelta(days=calendar_days) + end_time = now else: - start_time = now - timedelta(days=days) - end_time = now + # 默认增量模式: 首次拉取回溯 N 天, 已有数据则从最新时间增量补到今天 + last_dt = _latest_minute_datetime(repo) + if last_dt: + start_time = last_dt + else: + start_time = now - timedelta(days=days) + end_time = now limit = resolve_limit( capset, diff --git a/backend/app/services/pipeline_jobs.py b/backend/app/services/pipeline_jobs.py index aba9f9c..a7dc694 100644 --- a/backend/app/services/pipeline_jobs.py +++ b/backend/app/services/pipeline_jobs.py @@ -253,6 +253,12 @@ class JobStore: logger.warning("reap_stale: 强制取消卡死 job %s (已运行 %.0fs)", jid, elapsed) self.fail(jid, f"超时自动取消 (运行 {int(elapsed)}s, 疑似卡死)") + # 强制释放重任务锁: 卡死的线程无法被中断, 锁永远不会自然释放。 + # job 已标记 failed, 即使僵尸线程后续写入 parquet, 下次拉取会覆盖, 安全。 + try: + _heavy_run_lock.release() + except RuntimeError: + pass def clear(self) -> None: """清空所有任务(内存 + 磁盘文件)。""" diff --git a/frontend/src/components/StockIntradayChart.tsx b/frontend/src/components/StockIntradayChart.tsx index e4fbb0a..411870f 100644 --- a/frontend/src/components/StockIntradayChart.tsx +++ b/frontend/src/components/StockIntradayChart.tsx @@ -32,11 +32,10 @@ export function StockIntradayChart({ }) const fetchMinute = useMutation({ - mutationFn: () => api.extendMinuteHistory(5, 'day'), + mutationFn: () => api.syncMinuteSingle(symbol), onSuccess: () => { qc.invalidateQueries({ queryKey: ['kline-minute', symbol] }) - qc.invalidateQueries({ queryKey: QK.dataStatus }) - qc.invalidateQueries({ queryKey: QK.pipelineJobs }) + qc.invalidateQueries({ queryKey: QK.klineMinute(symbol, date ?? '') }) setMinuteDismissed(false) }, }) @@ -61,7 +60,7 @@ export function StockIntradayChart({ {fetchMinute.isPending ? (
- 正在获取最近5日分钟K… + 正在获取分钟K数据…
) : sourceIsNone ? ( // 数据源确认无此日分钟数据 (停牌/复牌延迟等): 静态提示 + 保留重试 diff --git a/frontend/src/components/data/MinuteSyncConfig.tsx b/frontend/src/components/data/MinuteSyncConfig.tsx index 2a96b1e..31d50f6 100644 --- a/frontend/src/components/data/MinuteSyncConfig.tsx +++ b/frontend/src/components/data/MinuteSyncConfig.tsx @@ -1,11 +1,10 @@ import { useState, useEffect } from 'react' import { useMutation, useQuery, useQueryClient } from '@tanstack/react-query' -import { Loader2 } from 'lucide-react' +import { Loader2, Trash2, Download, Calendar } from 'lucide-react' import { api } from '@/lib/api' import { QK } from '@/lib/queryKeys' -import { isExpertOrAbove } from '@/lib/capability-labels' -export function MinuteSyncConfig({ caps, isRunning, onStart }: { caps: { label: string; capabilities: Record } | undefined; isRunning: boolean; onStart: () => void }) { +export function MinuteSyncConfig({ caps, onJobStart }: { caps: { label: string; capabilities: Record } | undefined; onJobStart?: (jobId: string) => void }) { const qc = useQueryClient() const prefs = useQuery({ queryKey: QK.preferences, @@ -29,8 +28,40 @@ export function MinuteSyncConfig({ caps, isRunning, onStart }: { caps: { label: update.mutate({ enabled: !enabled, days: localDays }) } + const setDays = (v: number) => { + const clamped = Math.max(1, Math.min(30, v)) + setLocalDays(clamped) + update.mutate({ enabled, days: clamped }) + } + + // 清空分钟K数据 (二次确认) + const [confirmClear, setConfirmClear] = useState(false) + const clearMutation = useMutation({ + mutationFn: () => api.clearMinute(), + onSuccess: () => { + setConfirmClear(false) + qc.invalidateQueries({ queryKey: QK.dataStatus }) + }, + }) + + // 手动获取 (两个独立按钮, 各自指定天数, 不影响自动同步偏好) + const [fetchingMode, setFetchingMode] = useState<'' | '40d' | '1y'>('') + const handleFetch = (mode: '40d' | '1y') => { + if (!hasMinuteCap) return + // 两个按钮都用向前扩展模式: 从本地最早数据往前补, 叠加避免缺口 + const fetchDays = mode === '40d' ? 40 : 365 + setFetchingMode(mode) + api.syncMinute(fetchDays, true).then((res) => { + qc.invalidateQueries({ queryKey: QK.pipelineJobs }) + qc.invalidateQueries({ queryKey: QK.dataStatus }) + // 通知主页面跟踪 job 进度 (ActiveJobCard 会显示实时进度+日志) + if (res.job_id && onJobStart) onJobStart(res.job_id) + }).finally(() => setFetchingMode('')) + } + return (
+ {/* 第 1 行: 自动同步开关 + 天数 */}
- {enabled ? '自动同步' : '已关闭'} + {enabled ? '盘后自动同步' : '已关闭'}
- {!hasMinuteCap && ( - - 需 Pro+ - - )} -
- -
- 同步天数
-
+ >− +
{localDays}
+ >+
+ {!hasMinuteCap && ( + 需 Pro+ + )}
-
-
向前扩展历史数据
- -
- -
- A股标的 · 原始数据存储(查询时实时复权) -
-
- ) -} - -function MinuteExtendControls({ hasMinuteCap, tierLabel, isRunning, onStart }: { hasMinuteCap: boolean; tierLabel: string; isRunning: boolean; onStart: () => void }) { - const qc = useQueryClient() - // 月单位(按月扩展更长的分钟K历史)仅 Expert+ 开放;Pro 仅可用"天"(1~15 天) - const canUseMonth = isExpertOrAbove(tierLabel) - const [unit, setUnit] = useState<'day' | 'month'>('day') - const [value, setValue] = useState(5) - const [confirmOpen, setConfirmOpen] = useState(false) - - const dataStatus = useQuery({ - queryKey: QK.dataStatus, - queryFn: api.dataStatus, - }) - // 判断本地是否已有分钟K数据:后端 _safe_aggregate_minute 为避免全表扫描, - // rows 恒为 0,改用 trading_days(分区目录数,真实统计)判断。 - const hasMinuteData = !!(dataStatus.data?.minute?.trading_days) - - const extend = useMutation({ - mutationFn: () => api.extendMinuteHistory(value, unit), - onSuccess: () => { - onStart() - qc.invalidateQueries({ queryKey: QK.pipelineJobs }) - qc.invalidateQueries({ queryKey: QK.dataStatus }) - }, - }) - - // 各单位上限:day 15 天,month 6 月(180 天) - const maxValue = unit === 'month' ? 6 : 15 - - const handleFetch = () => { - if (!hasMinuteData) { - setConfirmOpen(true) - } else { - extend.mutate() - } - } - - // 切换单位时把 value clamp 到新单位的上限 - const switchUnit = (u: 'day' | 'month') => { - if (u === unit) return - setUnit(u) - const max = u === 'month' ? 6 : 15 - setValue(v => Math.min(v, max)) - } - - return ( - <> -
-
- -
- {value} -
- -
- - {canUseMonth ? ( -
- {(['day', 'month'] as const).map(u => ( - - ))} -
- ) : ( - - )} + {/* 第 2 行: 两个手动获取按钮 (40天快速 / 1年分段) */} +
+ +
+ {/* 第 3 行: 清空 */} - {confirmOpen && ( + {/* 说明 */} +
+ A股标的 · 前复权价格 · 均从本地最早数据向前叠加 ·{' '} + 单次拉满约 40 个交易日,{' '} + 1 年按月分段 (速度较慢) +
+ + {/* 清空确认弹窗 */} + {confirmClear && (
-
setConfirmOpen(false)} /> +
!clearMutation.isPending && setConfirmClear(false)} />
-
本地暂无分钟K数据,是否立即获取最近 {value} {unit === 'month' ? '月' : '天'}的分钟K?
+
确认清空分钟K数据?
+
+ 此操作仅删除分钟K (kline_minute) 数据, 不影响日K、复权因子、指标等其他数据。清空后可重新获取。 +
- - +
)} - +
) } diff --git a/frontend/src/lib/api.ts b/frontend/src/lib/api.ts index e3b98c1..7065dc7 100644 --- a/frontend/src/lib/api.ts +++ b/frontend/src/lib/api.ts @@ -1236,8 +1236,21 @@ export const api = { `/api/kline/sync?symbol=${encodeURIComponent(symbol)}&days=${days}`, { method: 'POST' }, ), - syncMinute: () => - request<{ status: string; job_id: string }>('/api/kline/sync_minute', { method: 'POST' }), + syncMinute: (days?: number, extend?: boolean) => + request<{ status: string; job_id: string }>('/api/kline/sync_minute', { + method: 'POST', + body: JSON.stringify({ ...(days ? { days } : {}), ...(extend ? { extend: true } : {}) }), + }), + syncMinuteSingle: (symbol: string) => + request<{ status: string; symbol: string; rows: number }>('/api/kline/sync_minute_single', { + method: 'POST', + body: JSON.stringify({ symbol }), + }), + clearMinute: () => + request<{ status: string; removed: number }>('/api/kline/clear_minute', { + method: 'POST', + body: JSON.stringify({ confirm: true }), + }), extendHistory: (value: number, unit: 'day' | 'month' | 'year') => request<{ status: string; job_id: string }>('/api/kline/extend_history', { method: 'POST', @@ -1248,11 +1261,6 @@ export const api = { method: 'POST', body: JSON.stringify({ start_date: startDate }), }), - extendMinuteHistory: (value: number, unit: 'day' | 'month') => - request<{ status: string; job_id: string }>('/api/kline/extend_minute_history', { - method: 'POST', - body: JSON.stringify({ value, unit }), - }), rebuildEnriched: () => request<{ status: string; job_id: string }>('/api/kline/rebuild_enriched', { method: 'POST', diff --git a/frontend/src/pages/Data.tsx b/frontend/src/pages/Data.tsx index 1f205c5..3e4c980 100644 --- a/frontend/src/pages/Data.tsx +++ b/frontend/src/pages/Data.tsx @@ -1085,7 +1085,7 @@ export function Data() { {openSettings === 'minute' && ( setOpenSettings(null)}> - setOpenSettings(null)} /> + { setActiveJobId(jobId); setOpenSettings(null) }} /> )}