消灭更多bug,修改K线颜色,macd颜色
This commit is contained in:
+70
-30
@@ -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})
|
||||
|
||||
|
||||
+84
-6
@@ -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
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user