From 3401cb56bf8cfe464bae9d6d86a13bcac27cd88e Mon Sep 17 00:00:00 2001 From: jackyu66git Date: Thu, 13 Nov 2025 20:34:36 +0800 Subject: [PATCH] =?UTF-8?q?=E6=B6=88=E7=81=AD=E6=9B=B4=E5=A4=9Abug?= =?UTF-8?q?=EF=BC=8C=E4=BF=AE=E6=94=B9K=E7=BA=BF=E9=A2=9C=E8=89=B2?= =?UTF-8?q?=EF=BC=8Cmacd=E9=A2=9C=E8=89=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- datasvc/app/main.py | 100 +++++++++++++++++++++++++++------------ datasvc/app/storage.py | 90 ++++++++++++++++++++++++++++++++--- web/templates/index.html | 52 ++++++++++---------- 3 files changed, 180 insertions(+), 62 deletions(-) diff --git a/datasvc/app/main.py b/datasvc/app/main.py index 000d8b3..61728f5 100644 --- a/datasvc/app/main.py +++ b/datasvc/app/main.py @@ -942,16 +942,31 @@ async def rest_poll_loop( "mode": "rest_poll", }, ) - await process_candles( - symbol, - timeframe, - candles, - derived_timeframes, - tf_ms, - finalized=True, - allow_verification=True, - closed_flags=closed_flags, - ) + try: + await process_candles( + symbol, + timeframe, + candles, + derived_timeframes, + tf_ms, + finalized=True, + allow_verification=True, + closed_flags=closed_flags, + ) + except asyncio.CancelledError: + raise + except Exception as exc: + logger.exception( + "处理基础周期 K 线失败 [%s %s]", + symbol, + timeframe, + ) + state = fetch_states.get(state_key) + if state: + state.last_error = str(exc) + state.consecutive_errors += 1 + await asyncio.sleep(interval) + continue await asyncio.sleep(interval) except asyncio.CancelledError: raise @@ -1028,7 +1043,7 @@ def build_exchange(): async def fetch_loop(symbol: str, timeframe: str): - """初次通过 REST 补齐历史,随后转入 Binance WebSocket 拉取增量。""" + """初次通过 REST 补齐历史,随后持续轮询/流式拉取增量。""" derived_timeframes = AGGREGATION_TARGETS.get(timeframe, []) global RESAMPLE_WARNING_EMITTED if derived_timeframes and not RESAMPLE_AVAILABLE and not RESAMPLE_WARNING_EMITTED: @@ -1056,27 +1071,52 @@ async def fetch_loop(symbol: str, timeframe: str): extra={"symbol": symbol, "timeframe": timeframe, "timestamp": last_ts, "count": initial_count}, ) - start_since = parse_start_from_ms(START_FROM) - if last_ts is not None: - rewind_since = max(0, last_ts - tf_ms) - initial_since = max(start_since, rewind_since) - else: - initial_since = start_since - - logger.info("启动拉取任务", extra={"symbol": symbol, "timeframe": timeframe}) + start_from = parse_start_from_ms(START_FROM) + backoff = 1.0 try: - await rest_catchup(symbol, timeframe, derived_timeframes, tf_ms, initial_since) - if WS_ENABLED: - await stream_loop(symbol, timeframe, derived_timeframes, tf_ms) - else: - await rest_poll_loop(symbol, timeframe, derived_timeframes, tf_ms) - except asyncio.CancelledError: - logger.info("取消拉取任务", extra={"symbol": symbol, "timeframe": timeframe}) - state = fetch_states.get(state_key) - if state: - state.last_error = "cancelled" - raise + while True: + state = fetch_states[state_key] + state.started_at = datetime.utcnow() + if state.last_candle_ts is not None: + rewind_since = max(0, state.last_candle_ts - tf_ms) + initial_since = max(start_from, rewind_since) + else: + initial_since = start_from + + logger.info( + "启动拉取任务 [%s %s] since=%s", + symbol, + timeframe, + initial_since, + ) + + try: + await rest_catchup(symbol, timeframe, derived_timeframes, tf_ms, initial_since) + if WS_ENABLED: + await stream_loop(symbol, timeframe, derived_timeframes, tf_ms) + else: + await rest_poll_loop(symbol, timeframe, derived_timeframes, tf_ms) + except asyncio.CancelledError: + logger.info("取消拉取任务 [%s %s]", symbol, timeframe) + state.last_error = "cancelled" + raise + except Exception as exc: + state.last_error = str(exc) + state.consecutive_errors += 1 + logger.exception( + "拉取任务异常 [%s %s],%.1f 秒后重启", + symbol, + timeframe, + backoff, + ) + await asyncio.sleep(backoff) + backoff = min(backoff * BACKOFF_BASE, BACKOFF_MAX) + continue + else: + backoff = 1.0 + logger.warning("拉取循环提前结束 [%s %s],1 秒后重启", symbol, timeframe) + await asyncio.sleep(1.0) finally: logger.info("拉取任务退出", extra={"symbol": symbol, "timeframe": timeframe}) diff --git a/datasvc/app/storage.py b/datasvc/app/storage.py index fd231b4..7e4d7fc 100644 --- a/datasvc/app/storage.py +++ b/datasvc/app/storage.py @@ -1,10 +1,15 @@ +import logging import os +import shutil import threading +from datetime import datetime from typing import List, Optional import pandas as pd import pyarrow.dataset as ds +logger = logging.getLogger("datasvc") + _lock = threading.Lock() @@ -24,7 +29,20 @@ def read_candles(base_dir: str, symbol: str, timeframe: str, start: Optional[int p = _path(base_dir, symbol, timeframe) if not os.path.exists(p): return pd.DataFrame(columns=["timestamp", "open", "high", "low", "close", "volume"]) # empty - df = pd.read_parquet(p) + try: + df = pd.read_parquet(p) + except Exception as exc: + with _lock: + backup = _backup_corrupted_file(p) + extra = f",已备份至 {backup}" if backup else "" + logger.warning( + "读取缓存失败,将视为空数据 [%s %s]%s:%s", + symbol, + timeframe, + extra, + exc, + ) + return pd.DataFrame(columns=["timestamp", "open", "high", "low", "close", "volume"]) if start is not None: df = df[df["timestamp"] >= int(start)] if end is not None: @@ -38,19 +56,53 @@ def upsert_candles(base_dir: str, symbol: str, timeframe: str, candles: List[Lis new_df = pd.DataFrame(candles, columns=["timestamp", "open", "high", "low", "close", "volume"]) with _lock: if os.path.exists(p): - old = pd.read_parquet(p) + try: + old = pd.read_parquet(p) + except Exception as exc: + backup = _backup_corrupted_file(p) + extra = f",已备份至 {backup}" if backup else "" + logger.warning( + "读取缓存失败,准备重建文件 [%s %s]%s:%s", + symbol, + timeframe, + extra, + exc, + ) + old = pd.DataFrame(columns=["timestamp", "open", "high", "low", "close", "volume"]) merged = pd.concat([old, new_df], ignore_index=True) merged = merged.drop_duplicates(subset=["timestamp"], keep="last").sort_values("timestamp") else: merged = new_df.sort_values("timestamp") - merged.to_parquet(p, index=False) + temp_path = f"{p}.tmp" + try: + merged.to_parquet(temp_path, index=False) + os.replace(temp_path, p) + finally: + if os.path.exists(temp_path): + try: + os.remove(temp_path) + except OSError: + pass def get_last_timestamp(base_dir: str, symbol: str, timeframe: str) -> Optional[int]: p = _path(base_dir, symbol, timeframe) if not os.path.exists(p): return None - df = pd.read_parquet(p) + try: + df = pd.read_parquet(p) + except Exception as exc: + with _lock: + backup = _backup_corrupted_file(p) + extra = f",已备份至 {backup}" if backup else "" + logger.warning( + "获取最后时间戳失败 [%s %s]%s:%s", + symbol, + timeframe, + extra, + exc, + ) + return None if df.empty: return None return int(df["timestamp"].iloc[-1]) @@ -60,10 +112,36 @@ def read_candle_exact(base_dir: str, symbol: str, timeframe: str, timestamp: int p = _path(base_dir, symbol, timeframe) if not os.path.exists(p): return pd.DataFrame(columns=["timestamp", "open", "high", "low", "close", "volume"]) - dataset = ds.dataset(p, format="parquet") - table = dataset.to_table(filter=ds.field("timestamp") == int(timestamp)) + try: + dataset = ds.dataset(p, format="parquet") + table = dataset.to_table(filter=ds.field("timestamp") == int(timestamp)) + except Exception as exc: + with _lock: + backup = _backup_corrupted_file(p) + extra = f",已备份至 {backup}" if backup else "" + logger.warning( + "读取指定时间 K 线失败 [%s %s]%s:%s", + symbol, + timeframe, + extra, + exc, + ) + return pd.DataFrame(columns=["timestamp", "open", "high", "low", "close", "volume"]) if table.num_rows == 0: return pd.DataFrame(columns=["timestamp", "open", "high", "low", "close", "volume"]) return table.to_pandas() +def _backup_corrupted_file(path: str) -> Optional[str]: + try: + if not os.path.exists(path): + return None + timestamp = datetime.utcnow().strftime("%Y%m%d%H%M%S") + backup_path = f"{path}.corrupted.{timestamp}" + shutil.move(path, backup_path) + return backup_path + except Exception as exc: + logger.warning("备份损坏文件失败 (%s):%s", path, exc) + return None + + diff --git a/web/templates/index.html b/web/templates/index.html index ede1d64..56d7c0e 100644 --- a/web/templates/index.html +++ b/web/templates/index.html @@ -17,7 +17,7 @@