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 @@