From 70942e894208b2ab9de265cd7338fc1ef2a4f684 Mon Sep 17 00:00:00 2001 From: jackyu66git Date: Tue, 18 Nov 2025 12:25:22 +0800 Subject: [PATCH] =?UTF-8?q?=E5=88=A0=E9=99=A4=E4=B8=8D=E8=A6=81=E7=9A=84?= =?UTF-8?q?=E4=B8=9C=E8=A5=BF?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .DS_Store | Bin 10244 -> 10244 bytes datasvc/Dockerfile | 26 - datasvc/README.md | 145 ---- datasvc/app/main.py | 1455 ----------------------------------- datasvc/app/storage.py | 172 ----- datasvc/app/worker.py | 27 - datasvc/docker-compose.yml | 19 - datasvc/pairs.json | 5 - datasvc/requirements.txt | 9 - 交易记录/TradingRecord.xlsx | Bin 10158 -> 0 bytes 交易记录/交易规则.docx | Bin 12767 -> 0 bytes 11 files changed, 1858 deletions(-) delete mode 100644 datasvc/Dockerfile delete mode 100644 datasvc/README.md delete mode 100644 datasvc/app/main.py delete mode 100644 datasvc/app/storage.py delete mode 100644 datasvc/app/worker.py delete mode 100644 datasvc/docker-compose.yml delete mode 100644 datasvc/pairs.json delete mode 100644 datasvc/requirements.txt delete mode 100644 交易记录/TradingRecord.xlsx delete mode 100644 交易记录/交易规则.docx diff --git a/.DS_Store b/.DS_Store index e997eb39c2deb35e8a2bab09e3d0bbc0afa49ae2..c253e25066d261d90128ec81ab888c384605e864 100644 GIT binary patch delta 38 ucmZn(XbG6$F8U^hRb&SoBgi_DwviD+_fEMR2Z%&zc@Wiyi)Gcy1HS_}06 delta 132 zcmZn(XbG6$RCU^hRb-ew+wi_APp<;4X_Ir&Kp3=BIb^NAH|NLE*yo15q;7#SGX z>L^qj8XB7EC|DR;)Yfuxh$`z_2gPUS 衍生周期列表由程序自动推导,无需手动写入 `TIMEFRAMES`。 - ---- - -## 2. 启动与关闭 - -### 2.1 Docker 方式 -```bash -cd user_data/Chan/datasvc -docker compose up -d # 启动 -docker compose logs -f # 查看日志 -docker compose down # 关闭 -``` - -### 2.2 本地运行(无 Docker) -```bash -export DATA_DIR=./data -export TIMEFRAMES="1m,1h,1d" - -cd /Users/jack/Project/freqtrade -uvicorn user_data.Chan.datasvc.app.main:app --reload -``` - -关闭时 Ctrl+C 即可,服务会自动取消后台抓取任务并释放资源。 - ---- - -## 3. 数据存储与聚合 - -### 3.1 基础周期 -只会为 `TIMEFRAMES` 声明的基础周期创建抓取任务(例如 `1m / 1h / 1d`)。 - -### 3.2 衍生周期 -启动后自动维护以下聚合: - -| 基础周期 | 自动生成 | -| --- | --- | -| `1m` | `2m, 3m, 4m, 5m, 10m, 15m, 20m, 25m, 30m` | -| `1h` | `2h, 3h, 4h, 6h, 8h, 12h, 16h` | -| `1d` | `2d, 3d, 4d, 5d, 6d` | -| `1w` | `2w` | -| `1M` | `2M, 3M, 6M` | - -聚合过程通过 `technical.util.resample_to_interval` 完成,写入同一 Parquet 数据目录。 -所有周期都可以被 REST/WS 访问。 - -### 3.3 数据目录 -``` -{DATA_DIR}/{timeframe}/{symbol}.parquet -``` - ---- - -## 4. 接口调用 - -### 4.1 健康检查 -``` -GET /health -``` -返回运行状态、基础/衍生周期列表、各抓取任务的最新进度与错误计数,便于监控。 - -### 4.2 REST API -``` -GET /api/candles?symbol=BTC/USDT:USDT&tf=2h&start=1700000000000&end=1700003600000 -``` -参数说明: -- `symbol`:交易对(必须在 `SYMBOLS` 列表中) -- `tf`:时间周期(支持基础或衍生) -- `start` / `end`:毫秒时间戳,可选 - -返回示例: -```json -[ - {"timestamp": 1700000000000, "open": 36000.0, "high": 36120.0, "low": 35980.0, "close": 36050.0, "volume": 125.4}, - ... -] -``` - -### 4.3 WebSocket -``` -ws://localhost:8000/ws?symbol=ETH/USDT:USDT&tf=15m&since=1700000000000 -``` -- 首次连接:收到 `snapshot` 消息(快照数组) -- 后续增量:收到 `upsert` 消息(最新几根K线),以及周期性 `ping` - -消息示例: -```json -{"topic":"candles.ETH/USDT:USDT.15m","type":"snapshot","data":[{"t":1700000000000,"o":2000.0,"h":2005.0,"l":1995.0,"c":2002.5,"v":312.7}, ...]} -{"topic":"candles.ETH/USDT:USDT.15m","type":"upsert","data":{"t":1700000900000,"o":2002.5,"h":2006.0,"l":2000.0,"c":2004.0,"v":120.8}} -``` - ---- - -## 5. 停机与维护 - -- **正常关闭**:`docker compose down` 或 Ctrl+C。服务会等待所有抓取任务结束并关闭 `ccxt` 客户端。 -- **异常恢复**:若网络异常,服务会自动指数退避重试;可通过 `/health` 的 `consecutive_errors` 与 `last_error` 排查。 -- **数据清理**:直接删除 `DATA_DIR` 下对应的 Parquet 文件即可,下次启动会重新回补。 - ---- - -## 6. 常见问题 - -1. **缺少 `technical` 模块** - 聚合周期会跳过,并在日志中提示;先执行 `pip install technical` 再重启。 - -2. **接收不到某个周期的数据** - 确认该周期在 `TIMEFRAMES` 或自动聚合列表中;若是衍生周期,需要确保对应基础周期已在运行。 - -3. **如何新增交易对/周期** - - 交易对:编辑 `pairs.json`,每行一个字符串,保存后重启服务。 - - 周期:修改 `TIMEFRAMES` 环境变量(Docker 或本地启动命令)后重启。 - ---- - -欢迎结合自身策略或可视化前端直接消费本地数据服务。若要集成到其他项目,可直接引用 `/api/candles` 的 JSON 响应或订阅 `/ws` 的实时推送。 diff --git a/datasvc/app/main.py b/datasvc/app/main.py deleted file mode 100644 index f8b17ba..0000000 --- a/datasvc/app/main.py +++ /dev/null @@ -1,1455 +0,0 @@ -import os -import asyncio -import json -import logging -from contextlib import suppress -from dataclasses import dataclass, field -from datetime import datetime, timedelta -from time import time -from pathlib import Path -from collections import defaultdict -from threading import Event, RLock -from typing import Dict, List, Optional, Set, Tuple, Union - -import ccxt -import ccxt.async_support as ccxt_async -import pandas as pd -import websockets -from fastapi import FastAPI, WebSocket, WebSocketDisconnect, Query, HTTPException, status -from fastapi.responses import JSONResponse -from fastapi.middleware.cors import CORSMiddleware - -# docker compose down && docker compose build --no-cache && docker compose up -d -# docker compose down && docker compose build && docker compose up -d - -from technical.util import resample_to_interval # type: ignore - - -from .storage import candle_path, ensure_storage, read_candles, write_candles_snapshot - - -LOG_LEVEL = os.environ.get("LOG_LEVEL", "INFO").upper() -logging.basicConfig( - level=LOG_LEVEL, - format="%(asctime)s %(levelname)s [%(name)s] %(message)s", -) -logger = logging.getLogger("datasvc") -BASE_DIR = Path(__file__).resolve().parent.parent -DEFAULT_PAIRS_FILE = BASE_DIR / "pairs.json" -RESAMPLE_AVAILABLE = resample_to_interval is not None -RESAMPLE_WARNING_EMITTED = False -WS_ENABLED = os.environ.get("WS_ENABLED", "false").lower() in {"1", "true", "yes"} -REST_POLL_INTERVAL = max(1.0, float(os.environ.get("REST_POLL_INTERVAL", "5"))) -REST_POLL_WINDOW = max(1, int(os.environ.get("REST_POLL_WINDOW", "10"))) - -REST_MAX_CONCURRENCY = int(os.environ.get("REST_MAX_CONCURRENCY", "1")) -REST_FETCH_SEMAPHORE = asyncio.Semaphore(max(1, REST_MAX_CONCURRENCY)) - -AGGREGATION_PLAN: Dict[str, List[str]] = { - "1m": ["2m", "3m", "4m", "5m", "10m", "15m", "20m", "25m", "30m"], - "1h": ["2h", "3h", "4h", "6h", "8h", "12h", "16h"], - "1d": ["2d", "3d", "4d", "5d", "6d"], - "1w": ["2w"], - "1M": ["2M", "3M", "6M"], -} - -CandleRow = List[Union[int, float]] - - -CANDLE_COLUMNS = ["timestamp", "open", "high", "low", "close", "volume"] -CANDLES_CACHE: Dict[Tuple[str, str], pd.DataFrame] = {} -CACHE_LOCK = RLock() -BASE_TIMEFRAMES: Set[str] = set() -BASE_DIRTY_VERSION: Dict[Tuple[str, str], int] = {} -BASE_FLUSH_INTERVAL_SECONDS = 600 -CACHE_FILE_MTIME: Dict[Tuple[str, str], float] = {} -IS_ENGINE_PROCESS = os.environ.get("DATASVC_ENGINE") == "1" - - -def _empty_frame() -> pd.DataFrame: - return pd.DataFrame(columns=CANDLE_COLUMNS) - - -def _normalize_dataframe(df: pd.DataFrame) -> pd.DataFrame: - if df.empty: - return _empty_frame() - normalized = df.copy() - missing_columns = [col for col in CANDLE_COLUMNS if col not in normalized.columns] - for column in missing_columns: - normalized[column] = 0.0 if column != "timestamp" else 0 - normalized = normalized[CANDLE_COLUMNS] - normalized["timestamp"] = normalized["timestamp"].astype("int64") - for column in CANDLE_COLUMNS[1:]: - normalized[column] = normalized[column].astype("float64") - normalized = normalized.drop_duplicates(subset=["timestamp"], keep="last").sort_values("timestamp").reset_index(drop=True) - return normalized - - -def preload_candles_cache(symbols: List[str], timeframes: List[str]) -> None: - new_cache: Dict[Tuple[str, str], pd.DataFrame] = {} - new_mtime: Dict[Tuple[str, str], float] = {} - for symbol in symbols: - for timeframe in timeframes: - df = read_candles(DATA_DIR, symbol, timeframe, None, None) - normalized = _normalize_dataframe(df) - new_cache[(symbol, timeframe)] = normalized - try: - mtime = os.path.getmtime(candle_path(DATA_DIR, symbol, timeframe)) - except OSError: - mtime = 0.0 - new_mtime[(symbol, timeframe)] = mtime - with CACHE_LOCK: - CANDLES_CACHE.clear() - CANDLES_CACHE.update(new_cache) - BASE_DIRTY_VERSION.clear() - CACHE_FILE_MTIME.clear() - for key in new_cache: - if key[1] in BASE_TIMEFRAMES: - BASE_DIRTY_VERSION[key] = 0 - CACHE_FILE_MTIME[key] = new_mtime.get(key, 0.0) - - -def refresh_cache_from_disk(symbol: str, timeframe: str) -> None: - if timeframe not in BASE_TIMEFRAMES: - return - if IS_ENGINE_PROCESS: - return - key = (symbol, timeframe) - path = candle_path(DATA_DIR, symbol, timeframe) - try: - mtime = os.path.getmtime(path) - except FileNotFoundError: - with CACHE_LOCK: - if key not in CANDLES_CACHE: - CANDLES_CACHE[key] = _empty_frame() - CACHE_FILE_MTIME[key] = 0.0 - return - except OSError: - return - with CACHE_LOCK: - cached_mtime = CACHE_FILE_MTIME.get(key, 0.0) - if mtime <= cached_mtime: - return - df = read_candles(DATA_DIR, symbol, timeframe, None, None) - normalized = _normalize_dataframe(df) - with CACHE_LOCK: - CANDLES_CACHE[key] = normalized - CACHE_FILE_MTIME[key] = mtime - if timeframe in BASE_TIMEFRAMES: - BASE_DIRTY_VERSION[key] = 0 - - -def update_cache_mtime(symbol: str, timeframe: str) -> None: - path = candle_path(DATA_DIR, symbol, timeframe) - try: - mtime = os.path.getmtime(path) - except OSError: - mtime = time() - with CACHE_LOCK: - CACHE_FILE_MTIME[(symbol, timeframe)] = mtime - - -def rebuild_all_derived_timeframes(symbols: List[str]) -> None: - if not RESAMPLE_AVAILABLE: - return - for symbol in symbols: - for base_tf, targets in AGGREGATION_TARGETS.items(): - if not targets: - continue - resample_and_store(symbol, base_tf, targets) - - -def cache_get(symbol: str, timeframe: str, start: Optional[int] = None, end: Optional[int] = None) -> pd.DataFrame: - refresh_cache_from_disk(symbol, timeframe) - key = (symbol, timeframe) - with CACHE_LOCK: - df = CANDLES_CACHE.get(key) - if df is None: - df = _empty_frame() - result = df - if start is not None: - result = result[result["timestamp"] >= int(start)] - if end is not None: - result = result[result["timestamp"] <= int(end)] - return result.copy() - - -def cache_get_last_timestamp(symbol: str, timeframe: str) -> Optional[int]: - if timeframe in BASE_TIMEFRAMES: - refresh_cache_from_disk(symbol, timeframe) - key = (symbol, timeframe) - with CACHE_LOCK: - df = CANDLES_CACHE.get(key) - if df is None: - CANDLES_CACHE[key] = _empty_frame() - return None - if df.empty: - return None - return int(df["timestamp"].iloc[-1]) - - -def cache_update(symbol: str, timeframe: str, candles: List[CandleRow]) -> Optional[pd.DataFrame]: - if not candles: - return None - new_df = _normalize_dataframe(pd.DataFrame(candles, columns=CANDLE_COLUMNS)) - if new_df.empty: - return None - key = (symbol, timeframe) - with CACHE_LOCK: - existing = CANDLES_CACHE.get(key) - if existing is None or existing.empty: - merged = new_df - else: - merged = pd.concat([existing, new_df], ignore_index=True) - merged = _normalize_dataframe(merged) - with CACHE_LOCK: - CANDLES_CACHE[key] = merged - if timeframe in BASE_TIMEFRAMES: - BASE_DIRTY_VERSION[key] = BASE_DIRTY_VERSION.get(key, 0) + 1 - snapshot = merged.copy() - return snapshot - - -def collect_engine_status() -> dict: - updated_at = datetime.utcnow().replace(microsecond=0).isoformat() + "Z" - tasks = [state.to_payload() for state in fetch_states.values()] - return { - "updated_at": updated_at, - "tasks": tasks, - "queues": {}, - } - - -def write_engine_status_snapshot() -> None: - try: - ENGINE_STATUS_PATH.parent.mkdir(parents=True, exist_ok=True) - snapshot = collect_engine_status() - ENGINE_STATUS_PATH.write_text(json.dumps(snapshot, ensure_ascii=False), encoding="utf-8") - except Exception: - logger.warning("写入引擎状态快照失败", exc_info=True) - - -async def status_flush_worker(stop_event: Event, interval: float = 5.0) -> None: - await asyncio.to_thread(write_engine_status_snapshot) - try: - while not stop_event.is_set(): - await asyncio.sleep(interval) - await asyncio.to_thread(write_engine_status_snapshot) - except asyncio.CancelledError: - raise - - -def load_engine_status_snapshot() -> Optional[dict]: - try: - content = ENGINE_STATUS_PATH.read_text(encoding="utf-8") - except FileNotFoundError: - return None - except Exception: - logger.warning("读取引擎状态快照失败", exc_info=True) - return None - try: - return json.loads(content) - except json.JSONDecodeError: - logger.warning("解析引擎状态快照失败") - return None - - -async def flush_dirty_base_snapshots(force_all: bool = False) -> None: - with CACHE_LOCK: - if force_all: - target_entries = [] - for key in CANDLES_CACHE.keys(): - symbol, timeframe = key - if timeframe in BASE_TIMEFRAMES: - version = BASE_DIRTY_VERSION.get(key, 0) - target_entries.append((key, version)) - else: - target_entries = [(key, version) for key, version in BASE_DIRTY_VERSION.items() if version > 0] - snapshots = {key: CANDLES_CACHE.get(key, _empty_frame()).copy() for key, _ in target_entries} - if not snapshots: - return - failed: Set[Tuple[str, str]] = set() - for key, snapshot in snapshots.items(): - symbol, timeframe = key - try: - await asyncio.to_thread(write_candles_snapshot, DATA_DIR, symbol, timeframe, snapshot) - except Exception: - failed.add(key) - logger.exception( - "基础周期快照写入失败", - extra={"symbol": symbol, "timeframe": timeframe}, - ) - else: - update_cache_mtime(symbol, timeframe) - if not failed: - logger.debug( - "基础周期快照写入完成", - extra={"count": len(snapshots), "force_all": force_all}, - ) - with CACHE_LOCK: - for key, version in target_entries: - if key in failed: - continue - current_version = BASE_DIRTY_VERSION.get(key, 0) - if current_version == version: - BASE_DIRTY_VERSION[key] = 0 - - -async def base_flush_worker(): - try: - logger.info( - "基础周期定时写盘任务已启动", - extra={"interval_seconds": BASE_FLUSH_INTERVAL_SECONDS}, - ) - while True: - await asyncio.sleep(BASE_FLUSH_INTERVAL_SECONDS) - await flush_dirty_base_snapshots() - except asyncio.CancelledError: - raise - finally: - with suppress(Exception): - await flush_dirty_base_snapshots(force_all=True) - - -async def run_engine(stop_event: Optional[Event] = None): - if stop_event is None: - stop_event = Event() - logger.info("数据引擎启动") - fetch_tasks.clear() - try: - await asyncio.to_thread(rebuild_all_derived_timeframes, SYMBOLS) - except Exception: - logger.exception("初始化衍生周期失败,继续启动引擎") - try: - status_task = asyncio.create_task(status_flush_worker(stop_event), name="status::flush") - fetch_tasks.append(status_task) - flush_task = asyncio.create_task(base_flush_worker(), name="flush::base") - fetch_tasks.append(flush_task) - for s in SYMBOLS: - for tf in FETCH_TIMEFRAMES: - fetch_task = asyncio.create_task(fetch_loop(s, tf), name=f"fetch::{s}::{tf}") - fetch_tasks.append(fetch_task) - while not stop_event.is_set(): - await asyncio.sleep(1.0) - finally: - stop_event.set() - if fetch_tasks: - logger.info("数据引擎正在停止") - tasks = list(fetch_tasks) - for task in tasks: - task.cancel() - results = await asyncio.gather(*tasks, return_exceptions=True) - for result in results: - if isinstance(result, Exception) and not isinstance(result, asyncio.CancelledError): - logger.warning("任务停止时出现异常:%s", result) - fetch_tasks.clear() - with suppress(Exception): - await flush_dirty_base_snapshots(force_all=True) - with suppress(Exception): - await asyncio.to_thread(write_engine_status_snapshot) - logger.info("数据引擎已停止") - - -def _split_env_list(value: str) -> List[str]: - return [item.strip() for item in value.split(",") if item.strip()] - - -def _unique_preserve(values: List[str]) -> List[str]: - seen = set() - ordered: List[str] = [] - for item in values: - if item not in seen: - ordered.append(item) - seen.add(item) - return ordered - - -def _load_symbols() -> List[str]: - path = DEFAULT_PAIRS_FILE - if not path.is_file(): - default_symbols = ["BTC/USDT:USDT"] - logger.warning("交易对配置文件不存在,使用默认值", extra={"file": str(path), "symbols": default_symbols}) - return default_symbols - try: - content = path.read_text(encoding="utf-8") - data = json.loads(content) - except Exception: - default_symbols = ["BTC/USDT:USDT"] - logger.exception("读取交易对配置文件失败,使用默认值", extra={"file": str(path), "symbols": default_symbols}) - return default_symbols - - raw_symbols: List[str] = [] - if isinstance(data, list): - raw_symbols = [str(item).strip() for item in data if isinstance(item, str) and item.strip()] - elif isinstance(data, dict): - candidates = data.get("symbols") or data.get("pairs") - if isinstance(candidates, list): - raw_symbols = [str(item).strip() for item in candidates if isinstance(item, str) and item.strip()] - if not raw_symbols: - default_symbols = ["BTC/USDT:USDT"] - logger.warning("交易对配置文件未提供有效列表,使用默认值", extra={"file": str(path), "symbols": default_symbols}) - return default_symbols - - symbols = _unique_preserve(raw_symbols) - logger.info("已从配置文件载入交易对", extra={"file": str(path), "symbols": symbols}) - return symbols - - -def timeframe_to_minutes(tf: str) -> Optional[int]: - if not tf: - return None - unit = tf[-1] - try: - value = int(tf[:-1]) - except ValueError: - return None - multiplier = { - "m": 1, - "h": 60, - "d": 1440, - "w": 10080, - "M": 43200, # 30 天近似 - }.get(unit) - if multiplier is None: - return None - return value * multiplier - - -def binance_stream_symbol(symbol: str) -> str: - try: - base, rest = symbol.split("/", 1) - except ValueError: - cleaned = symbol.replace("/", "").split(":")[0] - return cleaned.lower() - quote = rest.split(":")[0] - return f"{base}{quote}".lower() - - -def build_stream_url(symbol: str, timeframe: str) -> str: - stream_symbol = binance_stream_symbol(symbol) - return f"{BINANCE_WS_BASE}/{stream_symbol}@kline_{timeframe}" - -DATA_DIR = os.environ.get("DATA_DIR", "/data") -EXCHANGE = os.environ.get("EXCHANGE", "binance") -SYMBOLS = _load_symbols() - -_default_timeframes = ["1m", "1h", "1d", "1w", "1M"] -requested_timeframes = _split_env_list(os.environ.get("TIMEFRAMES", ",".join(_default_timeframes))) -if not requested_timeframes: - requested_timeframes = _default_timeframes - -FETCH_TIMEFRAMES = _unique_preserve(requested_timeframes) -AVAILABLE_TIMEFRAMES = list(FETCH_TIMEFRAMES) -for base_tf in FETCH_TIMEFRAMES: - for derived_tf in AGGREGATION_PLAN.get(base_tf, []): - if derived_tf not in AVAILABLE_TIMEFRAMES: - AVAILABLE_TIMEFRAMES.append(derived_tf) -DERIVED_TIMEFRAMES = [tf for tf in AVAILABLE_TIMEFRAMES if tf not in FETCH_TIMEFRAMES] -AGGREGATION_TARGETS = {tf: AGGREGATION_PLAN.get(tf, []) for tf in FETCH_TIMEFRAMES} -BASE_TIMEFRAMES = set(FETCH_TIMEFRAMES) - -START_FROM = os.environ.get("START_FROM", "2025-01-01") # 首次启动拉取起始日期(UTC) -POLL_FACTOR = float(os.environ.get("POLL_FACTOR", "0.5")) # 轮询间隔 = tf_ms * factor -BACKOFF_BASE = float(os.environ.get("BACKOFF_BASE", "2.0")) -BACKOFF_MAX = float(os.environ.get("BACKOFF_MAX", "30.0")) -BINANCE_WS_BASE = os.environ.get("BINANCE_WS_BASE", "wss://fstream.binance.com/ws").rstrip("/") - -BASE_FLUSH_INTERVAL_MINUTES = max(1, int(os.environ.get("BASE_FLUSH_INTERVAL_MINUTES", "10"))) -BASE_FLUSH_INTERVAL_SECONDS = BASE_FLUSH_INTERVAL_MINUTES * 60 - -ENGINE_STATUS_PATH = Path(DATA_DIR) / "engine_status.json" - -VALID_SYMBOLS = set(SYMBOLS) -VALID_TIMEFRAMES = set(AVAILABLE_TIMEFRAMES) - -ensure_storage(DATA_DIR) -preload_candles_cache(SYMBOLS, AVAILABLE_TIMEFRAMES) - - -app = FastAPI(title="Local Data Service", version="0.1.0") -app.add_middleware( - CORSMiddleware, - allow_origins=["*"], - allow_credentials=True, - allow_methods=["*"], - allow_headers=["*"], -) - - -def tf_to_ms(tf: str) -> int: - minutes = timeframe_to_minutes(tf) - if minutes is None: - logger.warning("无法解析时间周期,默认使用 60 秒", extra={"timeframe": tf}) - return 60_000 - return minutes * 60_000 - - -def ensure_symbol_timeframe(symbol: str, timeframe: str) -> None: - if symbol not in VALID_SYMBOLS: - raise HTTPException( - status_code=status.HTTP_400_BAD_REQUEST, - detail=f"symbol 必须为 {sorted(VALID_SYMBOLS)} 之一。", - ) - if timeframe not in VALID_TIMEFRAMES: - raise HTTPException( - status_code=status.HTTP_400_BAD_REQUEST, - detail=f"tf 必须为 {sorted(VALID_TIMEFRAMES)} 之一。", - ) - - -def parse_start_from_ms(val: str) -> int: - """将 START_FROM 解析成毫秒级时间戳。 - 支持两种格式: - - YYYY-MM-DD(UTC 00:00:00) - - 整型毫秒时间戳字符串 - """ - try: - return int(val) - except Exception: - pass - try: - dt = datetime.fromisoformat(val) # 允许 '2022-01-01' 或 '2022-01-01T00:00:00' - except Exception: - # 回退到固定日期 - dt = datetime(2025, 1, 1) - return int(dt.timestamp() * 1000) - - -class Hub: - def __init__(self) -> None: - self.subscribers: Dict[str, List[WebSocket]] = {} - - def topic(self, symbol: str, timeframe: str) -> str: - return f"candles::{symbol}::{timeframe}" - - async def subscribe(self, ws: WebSocket, symbol: str, timeframe: str): - topic = self.topic(symbol, timeframe) - await ws.accept() - self.subscribers.setdefault(topic, []).append(ws) - - def _clean(self, topic: str): - conns = self.subscribers.get(topic, []) - self.subscribers[topic] = [w for w in conns if not w.client_state.name == "DISCONNECTED"] - - async def publish(self, symbol: str, timeframe: str, payload: dict): - topic = self.topic(symbol, timeframe) - conns = self.subscribers.get(topic, []) - if not conns: - return - message = json.dumps(payload, ensure_ascii=False) - dead: List[WebSocket] = [] - for ws in conns: - try: - await ws.send_text(message) - except Exception: - dead.append(ws) - if dead: - self.subscribers[topic] = [w for w in conns if w not in dead] - - -hub = Hub() - -fetch_tasks: List[asyncio.Task] = [] - -engine_runner_stop: Optional[Event] = None -engine_runner_task: Optional[asyncio.Task] = None - - -def resample_and_store(symbol: str, base_timeframe: str, derived_timeframes: List[str]) -> List[Tuple[str, List[CandleRow]]]: - if not derived_timeframes: - return [] - - base_tf_ms = tf_to_ms(base_timeframe) - - updates: List[Tuple[str, List[List[float]]]] = [] - for target_tf in derived_timeframes: - minutes = timeframe_to_minutes(target_tf) - if minutes is None: - logger.warning("无法解析聚合周期", extra={"target_timeframe": target_tf}) - continue - - last_ts = cache_get_last_timestamp(symbol, target_tf) - start_ts: Optional[int] = None - if last_ts is not None and base_tf_ms is not None and base_tf_ms > 0: - buffer_ms = minutes * 60_000 + base_tf_ms - start_ts = max(0, int(last_ts) - buffer_ms) - - base_slice = cache_get(symbol, base_timeframe, start_ts, None) - if base_slice.empty: - continue - base_slice = base_slice.copy() - if "date" not in base_slice.columns: - base_slice["date"] = pd.to_datetime(base_slice["timestamp"], unit="ms", utc=True) - base_slice = ( - base_slice.drop_duplicates(subset=["timestamp"], keep="last") - .sort_values("timestamp") - .reset_index(drop=True) - ) - base_slice["timestamp"] = base_slice["timestamp"].astype("int64") - - try: - if RESAMPLE_AVAILABLE: - derived_df = resample_to_interval(base_slice, minutes) # type: ignore[misc] - else: - derived_df = _fallback_resample_to_interval(base_slice, minutes) - except Exception: - logger.exception( - "聚合周期计算失败", - extra={"symbol": symbol, "base_timeframe": base_timeframe, "target_timeframe": target_tf}, - ) - continue - if derived_df is None or derived_df.empty: - continue - derived_df = derived_df.copy() - if "timestamp" not in derived_df.columns: - if "date" in derived_df.columns: - dates = pd.to_datetime(derived_df["date"], utc=True, errors="coerce") - derived_df["timestamp"] = (dates.astype("int64") // 1_000_000) - elif isinstance(derived_df.index, pd.DatetimeIndex): - idx = derived_df.index - if idx.tz is None: - idx = idx.tz_localize("UTC") - else: - idx = idx.tz_convert("UTC") - derived_df["timestamp"] = (idx.astype("int64") // 1_000_000) - if "timestamp" not in derived_df.columns: - logger.warning( - "聚合结果缺少 timestamp 列,已跳过", - extra={"target_timeframe": target_tf}, - ) - continue - derived_df = derived_df.dropna(subset=["timestamp", "open", "high", "low", "close", "volume"]) - if derived_df.empty: - continue - derived_df["timestamp"] = derived_df["timestamp"].astype("int64") - derived_df = derived_df.sort_values("timestamp") - if last_ts is not None: - derived_df = derived_df[derived_df["timestamp"] > last_ts] - if derived_df.empty: - continue - numpy_rows = derived_df[["timestamp", "open", "high", "low", "close", "volume"]].to_numpy() - records: List[CandleRow] = [] - for ts, o, h, l, c, v in numpy_rows: - records.append( - [ - int(ts), - float(o), - float(h), - float(l), - float(c), - float(v), - ] - ) - if not records: - continue - cache_update(symbol, target_tf, records) - updates.append((target_tf, records[-3:] if len(records) > 3 else records)) - return updates - - -def normalize_candles_for_timeframe(candles: List[CandleRow], tf_ms: int) -> Tuple[List[CandleRow], List[int]]: - if not candles: - return [], [] - normalized_map: Dict[int, CandleRow] = {} - for row in candles: - if not row: - continue - try: - ts = int(row[0]) - o = float(row[1]) - h = float(row[2]) - l = float(row[3]) - c = float(row[4]) - v = float(row[5]) - except (TypeError, ValueError, IndexError): - continue - normalized_map[ts] = [ts, o, h, l, c, v] - ordered_ts = sorted(normalized_map.keys()) - normalized: List[CandleRow] = [] - missing: List[int] = [] - last_ts: Optional[int] = None - for ts in ordered_ts: - normalized.append(normalized_map[ts]) - if last_ts is not None and tf_ms > 0: - delta = ts - last_ts - if delta > tf_ms: - gap_ts = last_ts + tf_ms - while gap_ts < ts: - missing.append(gap_ts) - gap_ts += tf_ms - last_ts = ts - return normalized, missing - - -def compute_live_derived_updates( - symbol: str, - base_timeframe: str, - derived_timeframes: List[str], - base_tf_ms: int, - candles: List[CandleRow], - last_closed_ts: Optional[int], -) -> Dict[str, List[Tuple[CandleRow, bool]]]: - updates: Dict[str, List[Tuple[CandleRow, bool]]] = {} - if not candles or not derived_timeframes or base_tf_ms <= 0: - return updates - - pending_updates: Dict[str, List[CandleRow]] = defaultdict(list) - - derived_ms_map: Dict[str, int] = {} - max_multiplier = 1 - for target_tf in derived_timeframes: - derived_ms = tf_to_ms(target_tf) - if derived_ms is None or derived_ms <= 0 or derived_ms % base_tf_ms != 0: - continue - multiplier = derived_ms // base_tf_ms - derived_ms_map[target_tf] = derived_ms - if multiplier > max_multiplier: - max_multiplier = multiplier - if not derived_ms_map: - return updates - - window_ms = max_multiplier * base_tf_ms - newest_ts = max(int(row[0]) for row in candles if row) - base_start = newest_ts - window_ms + base_tf_ms - if base_start < 0: - base_start = 0 - - base_df = cache_get(symbol, base_timeframe, base_start, newest_ts) - if base_df.empty: - return updates - base_df = base_df.sort_values("timestamp") - - base_rows: List[Tuple[int, float, float, float, float, float]] = [] - for record in candles: - try: - ts = int(record[0]) - if ts < base_start: - continue - base_rows.append( - ( - ts, - float(record[1]), - float(record[2]), - float(record[3]), - float(record[4]), - float(record[5]), - ) - ) - except (TypeError, ValueError, IndexError): - continue - - if base_rows: - temp_df = pd.DataFrame( - base_rows, - columns=["timestamp", "open", "high", "low", "close", "volume"], - ) - base_df = pd.concat([base_df, temp_df], ignore_index=True) - - if base_df.empty: - return updates - - base_df = ( - base_df.drop_duplicates(subset=["timestamp"], keep="last") - .sort_values("timestamp") - .reset_index(drop=True) - ) - - base_df_indexed = base_df.set_index("timestamp", drop=False) - if base_df_indexed.empty: - return updates - - for target_tf, derived_ms in derived_ms_map.items(): - multiplier = derived_ms // base_tf_ms - rows_with_status: List[Tuple[CandleRow, bool]] = [] - latest_available_ts = int(base_df_indexed.index.max()) - candidate_start = max(base_start, int(base_df_indexed.index.min())) - first_bucket = (candidate_start // derived_ms) * derived_ms - if first_bucket < candidate_start: - first_bucket += derived_ms - last_possible_start = latest_available_ts - (multiplier - 1) * base_tf_ms - current_start = first_bucket - while current_start <= last_possible_start: - expected_ts = [current_start + i * base_tf_ms for i in range(multiplier)] - subset = base_df_indexed.reindex(expected_ts) - if subset.isna().any().any(): - current_start += derived_ms - continue - start_ts = current_start - end_ts = start_ts + derived_ms - base_tf_ms - row: CandleRow = [ - start_ts, - float(subset.iloc[0]["open"]), - float(subset["high"].max()), - float(subset["low"].min()), - float(subset.iloc[-1]["close"]), - float(subset["volume"].sum()), - ] - closed = last_closed_ts is not None and last_closed_ts >= end_ts - pending_updates[target_tf].append(row) - rows_with_status.append((row, closed)) - current_start += derived_ms - if rows_with_status: - updates[target_tf] = rows_with_status - for target_tf, rows in pending_updates.items(): - cache_update(symbol, target_tf, rows) - return updates - - -@dataclass -class FetchState: - symbol: str - timeframe: str - started_at: datetime = field(default_factory=datetime.utcnow) - last_fetch_at: Optional[datetime] = None - last_candle_ts: Optional[int] = None - consecutive_errors: int = 0 - last_error: Optional[str] = None - - def to_payload(self) -> dict: - def serialize_dt(dt: Optional[datetime]) -> Optional[str]: - if not dt: - return None - return dt.replace(microsecond=0).isoformat() + "Z" - - return { - "symbol": self.symbol, - "timeframe": self.timeframe, - "started_at": serialize_dt(self.started_at), - "last_fetch_at": serialize_dt(self.last_fetch_at), - "last_candle_ts": self.last_candle_ts, - "consecutive_errors": self.consecutive_errors, - "last_error": self.last_error, - } - - -fetch_states: Dict[Tuple[str, str], FetchState] = {} - - -async def process_candles( - symbol: str, - timeframe: str, - candles: List[CandleRow], - derived_timeframes: List[str], - tf_ms: int, - finalized: bool, - closed_flags: Optional[List[bool]] = None, -) -> None: - if not candles: - return - if closed_flags is None or len(closed_flags) != len(candles): - closed_flags = [finalized] * len(candles) - state_key = (symbol, timeframe) - cache_update(symbol, timeframe, candles) - base_records = list(zip(candles, closed_flags)) - last_closed_ts: Optional[int] = None - for row, is_closed in base_records: - if is_closed: - if last_closed_ts is None or row[0] > last_closed_ts: - last_closed_ts = row[0] - if last_closed_ts is None: - last_closed_ts = candles[-1][0] - tf_ms - logger.info( - "基础周期 K 线更新完成", - extra={ - "symbol": symbol, - "timeframe": timeframe, - "count": len(candles), - "finalized": finalized, - "last_closed_ts": last_closed_ts, - }, - ) - derived_updates: List[Tuple[str, List[CandleRow]]] = [] - live_derived_updates: Dict[str, List[Tuple[CandleRow, bool]]] = {} - if derived_timeframes: - needs_resample = finalized or len(candles) > 1 - if needs_resample: - derived_updates = await asyncio.to_thread( - resample_and_store, - symbol, - timeframe, - derived_timeframes, - ) - if derived_updates: - logger.info( - "衍生周期批量聚合完成", - extra={ - "symbol": symbol, - "base_timeframe": timeframe, - "targets": [item[0] for item in derived_updates], - "origin": "resample" if RESAMPLE_AVAILABLE else "fallback", - }, - ) - live_derived_updates = await asyncio.to_thread( - compute_live_derived_updates, - symbol, - timeframe, - derived_timeframes, - tf_ms, - candles, - last_closed_ts, - ) - if live_derived_updates: - logger.info( - "衍生周期实时聚合完成", - extra={ - "symbol": symbol, - "base_timeframe": timeframe, - "targets": list(live_derived_updates.keys()), - "origin": "live", - }, - ) - for row, is_closed in base_records[-3:]: - payload = { - "topic": f"candles.{symbol}.{timeframe}", - "type": "upsert", - "data": { - "t": row[0], - "o": row[1], - "h": row[2], - "l": row[3], - "c": row[4], - "v": row[5], - "closed": bool(is_closed), - }, - } - await hub.publish(symbol, timeframe, payload) - last_closed_ts_for_derived = last_closed_ts - for target_tf, rows in derived_updates: - if not rows: - continue - target_tf_ms = tf_to_ms(target_tf) - for row in rows: - ts = int(row[0]) - o, h, l, c, v = map(float, row[1:]) - if target_tf_ms and target_tf_ms > 0: - derived_closed = last_closed_ts_for_derived is not None and last_closed_ts_for_derived >= ts + target_tf_ms - tf_ms - else: - derived_closed = last_closed_ts_for_derived is not None and last_closed_ts_for_derived >= ts - derived_state_key = (symbol, target_tf) - derived_state = fetch_states.get(derived_state_key) - if derived_state is None: - derived_state = FetchState(symbol=symbol, timeframe=target_tf) - fetch_states[derived_state_key] = derived_state - derived_state.last_fetch_at = datetime.utcnow() - derived_state.last_candle_ts = ts - derived_state.consecutive_errors = 0 - derived_state.last_error = None - payload = { - "topic": f"candles.{symbol}.{target_tf}", - "type": "upsert", - "data": { - "t": ts, - "o": o, - "h": h, - "l": l, - "c": c, - "v": v, - "closed": bool(derived_closed), - }, - } - await hub.publish(symbol, target_tf, payload) - if live_derived_updates: - for target_tf, items in live_derived_updates.items(): - if not items: - continue - for row, derived_closed in items: - ts = int(row[0]) - derived_state_key = (symbol, target_tf) - derived_state = fetch_states.get(derived_state_key) - if derived_state is None: - derived_state = FetchState(symbol=symbol, timeframe=target_tf) - fetch_states[derived_state_key] = derived_state - derived_state.last_fetch_at = datetime.utcnow() - derived_state.last_candle_ts = ts - derived_state.consecutive_errors = 0 - derived_state.last_error = None - payload = { - "topic": f"candles.{symbol}.{target_tf}", - "type": "upsert", - "data": { - "t": ts, - "o": float(row[1]), - "h": float(row[2]), - "l": float(row[3]), - "c": float(row[4]), - "v": float(row[5]), - "closed": bool(derived_closed), - }, - } - await hub.publish(symbol, target_tf, payload) - state = fetch_states.get(state_key) - if state: - state.last_fetch_at = datetime.utcnow() - state.last_candle_ts = candles[-1][0] - state.consecutive_errors = 0 - state.last_error = None - # 验证逻辑已移除,衍生周期的缺口依赖轮询补齐 - - -async def rest_catchup( - symbol: str, - timeframe: str, - derived_timeframes: List[str], - tf_ms: int, - start_since: int, -) -> None: - state_key = (symbol, timeframe) - state = fetch_states[state_key] - exchange = build_exchange() - since = start_since - backoff = 1.0 - gap_retry: Dict[int, int] = {} - logger.info("开始 REST 补齐历史", extra={"symbol": symbol, "timeframe": timeframe, "since": since}) - try: - while True: - now_ms = int(datetime.utcnow().timestamp() * 1000) - if since >= now_ms - tf_ms: - break - try: - async with REST_FETCH_SEMAPHORE: - candles = await exchange.fetch_ohlcv(symbol, timeframe, since=since, limit=1000) - except asyncio.CancelledError: - raise - except (ccxt.NetworkError, ccxt.ExchangeNotAvailable, ccxt.RequestTimeout) as exc: - logger.warning( - "历史补齐网络异常,准备重试", - extra={"symbol": symbol, "timeframe": timeframe, "error": str(exc)}, - ) - state.last_error = str(exc) - state.consecutive_errors += 1 - backoff = min(backoff * BACKOFF_BASE, BACKOFF_MAX) - await asyncio.sleep(backoff) - continue - except Exception as exc: - logger.exception( - "历史补齐发生异常,准备重试", - extra={"symbol": symbol, "timeframe": timeframe}, - ) - state.last_error = str(exc) - state.consecutive_errors += 1 - backoff = min(backoff * BACKOFF_BASE, BACKOFF_MAX) - await asyncio.sleep(backoff) - continue - if not candles: - break - candles, missing_ts = normalize_candles_for_timeframe(candles, tf_ms) - if not candles: - since += tf_ms - backoff = 1.0 - await asyncio.sleep(0.2) - continue - await process_candles( - symbol, - timeframe, - candles, - derived_timeframes, - tf_ms, - finalized=True, - closed_flags=[True] * len(candles), - ) - state.consecutive_errors = 0 - state.last_error = None - backoff = 1.0 - if missing_ts: - gap_start = missing_ts[0] - attempts = gap_retry.get(gap_start, 0) + 1 - gap_retry[gap_start] = attempts - if attempts <= 3: - logger.warning( - "检测到缺失 K 线,准备回补", - extra={ - "symbol": symbol, - "timeframe": timeframe, - "missing_from": gap_start, - "missing_to": missing_ts[-1], - "attempt": attempts, - }, - ) - since = gap_start - await asyncio.sleep(0.2) - continue - logger.error( - "缺失 K 线多次回补失败,已跳过", - extra={ - "symbol": symbol, - "timeframe": timeframe, - "missing_from": gap_start, - "missing_to": missing_ts[-1], - }, - ) - gap_retry.pop(gap_start, None) - else: - gap_retry.clear() - - since = candles[-1][0] + tf_ms - - now_ms = int(datetime.utcnow().timestamp() * 1000) - lag = now_ms - since - if lag > tf_ms * 10: - await asyncio.sleep(0.2) - else: - await asyncio.sleep(max(1.0, tf_ms * POLL_FACTOR / 1000.0)) - finally: - with suppress(Exception): - await exchange.close() - logger.info("REST 补齐完成", extra={"symbol": symbol, "timeframe": timeframe, "latest": state.last_candle_ts}) - - -async def rest_poll_loop( - symbol: str, - timeframe: str, - derived_timeframes: List[str], - tf_ms: int, -) -> None: - state_key = (symbol, timeframe) - window = max(REST_POLL_WINDOW, 1) - interval = max(REST_POLL_INTERVAL, 1.0) - exchange = build_exchange() - try: - while True: - state = fetch_states.get(state_key) - latest_ts = state.last_candle_ts if state else None - if latest_ts is None or latest_ts <= 0: - since = parse_start_from_ms(START_FROM) - else: - since = max(0, latest_ts - (window - 1) * tf_ms) - try: - async with REST_FETCH_SEMAPHORE: - candles = await exchange.fetch_ohlcv( - symbol, - timeframe, - since=since, - limit=max(window + 2, window), - ) - except asyncio.CancelledError: - raise - except (ccxt.NetworkError, ccxt.ExchangeNotAvailable, ccxt.RequestTimeout) as exc: - logger.warning( - "实时轮询网络异常,准备重试", - extra={"symbol": symbol, "timeframe": timeframe, "error": str(exc)}, - ) - await asyncio.sleep(interval) - continue - except Exception as exc: - logger.exception( - "实时轮询发生异常", - extra={"symbol": symbol, "timeframe": timeframe}, - ) - await asyncio.sleep(interval) - continue - candles, _ = normalize_candles_for_timeframe(candles, tf_ms) - if candles: - closed_flags = [True] * len(candles) - logger.info( - "轮询拉取基础周期完成", - extra={ - "symbol": symbol, - "timeframe": timeframe, - "count": len(candles), - "since": since, - "mode": "rest_poll", - }, - ) - try: - await process_candles( - symbol, - timeframe, - candles, - derived_timeframes, - tf_ms, - finalized=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 - finally: - with suppress(Exception): - await exchange.close() - logger.info("轮询任务退出", extra={"symbol": symbol, "timeframe": timeframe}) - - -async def stream_loop(symbol: str, timeframe: str, derived_timeframes: List[str], tf_ms: int): - state_key = (symbol, timeframe) - url = build_stream_url(symbol, timeframe) - while True: - try: - async with websockets.connect(url, ping_interval=20, ping_timeout=20) as ws: - logger.info("WebSocket 已连接", extra={"symbol": symbol, "timeframe": timeframe, "url": url}) - async for message in ws: - data = json.loads(message) - kline = data.get("k") - if not kline: - continue - is_closed = bool(kline.get("x")) - row: CandleRow = [ - int(kline["t"]), - float(kline["o"]), - float(kline["h"]), - float(kline["l"]), - float(kline["c"]), - float(kline["v"]), - ] - await process_candles( - symbol, - timeframe, - [row], - derived_timeframes, - tf_ms, - finalized=is_closed, - closed_flags=[is_closed], - ) - except asyncio.CancelledError: - logger.info("取消 WebSocket 任务", extra={"symbol": symbol, "timeframe": timeframe}) - raise - except Exception as exc: - logger.warning( - "WebSocket 连接异常,准备重连", - extra={"symbol": symbol, "timeframe": timeframe, "error": str(exc)}, - ) - state = fetch_states.get(state_key) - start_since = None - if state and state.last_candle_ts: - start_since = state.last_candle_ts + tf_ms - if start_since: - await rest_catchup(symbol, timeframe, derived_timeframes, tf_ms, start_since) - await asyncio.sleep(5.0) - - -def build_exchange(): - if EXCHANGE.lower() == "binance": - return ccxt_async.binance( - { - "enableRateLimit": True, - "timeout": 20_000, - "options": { - "adjustForTimeDifference": True, - "defaultType": "future", - "defaultSubType": "linear", - "defaultMarket": "future", - "defaultSettle": "USDT", - }, - } - ) - raise RuntimeError(f"Unsupported EXCHANGE: {EXCHANGE}") - - -async def fetch_loop(symbol: str, timeframe: str): - """初次通过 REST 补齐历史,随后持续轮询/流式拉取增量。""" - derived_timeframes = AGGREGATION_TARGETS.get(timeframe, []) - global RESAMPLE_WARNING_EMITTED - if derived_timeframes and not RESAMPLE_AVAILABLE and not RESAMPLE_WARNING_EMITTED: - logger.warning( - "缺少 technical.util.resample_to_interval 模块,聚合时间周期生成已跳过", - extra={"timeframe": timeframe}, - ) - RESAMPLE_WARNING_EMITTED = True - - tf_ms = tf_to_ms(timeframe) - state_key = (symbol, timeframe) - last_ts = cache_get_last_timestamp(symbol, timeframe) - fetch_states[state_key] = FetchState(symbol=symbol, timeframe=timeframe, last_candle_ts=last_ts) - initial_sync_flushed = False - - start_from = parse_start_from_ms(START_FROM) - backoff = 1.0 - - try: - 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 not initial_sync_flushed: - await flush_dirty_base_snapshots(force_all=True) - logger.info( - "初次同步完成,基础周期数据已写盘", - extra={"symbol": symbol, "timeframe": timeframe}, - ) - initial_sync_flushed = True - 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}) - - -@app.on_event("startup") -async def on_start(): - logger.info("API 服务启动完成") - if IS_ENGINE_PROCESS: - global engine_runner_stop, engine_runner_task - if engine_runner_task is None or engine_runner_task.done(): - engine_runner_stop = Event() - engine_runner_task = asyncio.create_task(run_engine(engine_runner_stop)) - - -@app.on_event("shutdown") -async def on_shutdown(): - logger.info("API 服务准备退出") - if IS_ENGINE_PROCESS: - global engine_runner_stop, engine_runner_task - if engine_runner_stop is not None: - engine_runner_stop.set() - if engine_runner_task is not None: - with suppress(Exception): - await engine_runner_task - engine_runner_task = None - engine_runner_stop = None - - -@app.get("/health") -async def health(): - now = datetime.utcnow().replace(microsecond=0).isoformat() + "Z" - engine_status = load_engine_status_snapshot() or {"updated_at": None, "tasks": [], "queues": {}} - return { - "status": "ok", - "time": now, - "exchange": EXCHANGE, - "symbols": SYMBOLS, - "base_timeframes": FETCH_TIMEFRAMES, - "derived_timeframes": DERIVED_TIMEFRAMES, - "timeframes": AVAILABLE_TIMEFRAMES, - "engine": engine_status, - "tasks": engine_status.get("tasks", []), - } - - -@app.get("/api/candles") -def api_candles( - symbol: str = Query(..., description="如 BTC/USDT:USDT"), - tf: str = Query("1m", description="时间周期"), - start: Optional[int] = Query(None, description="开始时间戳(ms)"), - end: Optional[int] = Query(None, description="结束时间戳(ms)"), -): - try: - ensure_symbol_timeframe(symbol, tf) - df = cache_get(symbol, tf, start, end) - records = df.to_dict("records") if not df.empty else [] - return JSONResponse(records) - except Exception as e: - return JSONResponse({"error": str(e)}, status_code=500) - - -@app.websocket("/ws") -async def ws_endpoint(websocket: WebSocket, symbol: str, tf: str, since: Optional[int] = None): - if symbol not in VALID_SYMBOLS or tf not in VALID_TIMEFRAMES: - await websocket.close(code=status.WS_1008_POLICY_VIOLATION, reason="invalid symbol/timeframe") - return - await hub.subscribe(websocket, symbol, tf) - try: - snap = cache_get(symbol, tf, since, None) - await websocket.send_text( - json.dumps( - { - "topic": f"candles.{symbol}.{tf}", - "type": "snapshot", - "data": [ - {"t": int(r["timestamp"]), "o": r["open"], "h": r["high"], "l": r["low"], "c": r["close"], "v": r["volume"]} - for _, r in snap.iterrows() - ], - }, - ensure_ascii=False, - ) - ) - except Exception: - pass - try: - while True: - await asyncio.sleep(30) - await websocket.send_text(json.dumps({"type": "ping", "ts": int(datetime.utcnow().timestamp() * 1000)})) - except WebSocketDisconnect: - return - - -@app.get("/") -def root(): - return { - "service": "Local Data Service", - "exchange": EXCHANGE, - "symbols": SYMBOLS, - "base_timeframes": FETCH_TIMEFRAMES, - "derived_timeframes": DERIVED_TIMEFRAMES, - "timeframes": AVAILABLE_TIMEFRAMES, - } - - -def _fallback_resample_to_interval(df: pd.DataFrame, minutes: int) -> pd.DataFrame: - if df.empty or minutes <= 0: - return pd.DataFrame(columns=CANDLE_COLUMNS) - working = df.copy() - if "timestamp" not in working.columns: - return pd.DataFrame(columns=CANDLE_COLUMNS) - working["date"] = pd.to_datetime(working["timestamp"], unit="ms", utc=True) - working = working.set_index("date", drop=True) - columns = ["open", "high", "low", "close", "volume"] - for column in columns: - if column not in working.columns: - working[column] = 0.0 - working = working[columns] - rule = f"{minutes}T" - aggregated = working.resample(rule, label="left", closed="left").agg( - { - "open": "first", - "high": "max", - "low": "min", - "close": "last", - "volume": "sum", - } - ) - aggregated = aggregated.dropna(subset=["open", "high", "low", "close"]).reset_index() - aggregated["timestamp"] = (aggregated["date"].astype("int64") // 1_000_000) - aggregated = aggregated.drop(columns=["date"], errors="ignore") - aggregated = aggregated.dropna(subset=["timestamp"]).reset_index(drop=True) - aggregated["timestamp"] = aggregated["timestamp"].astype("int64") - return aggregated[CANDLE_COLUMNS] - - diff --git a/datasvc/app/storage.py b/datasvc/app/storage.py deleted file mode 100644 index f00a2ef..0000000 --- a/datasvc/app/storage.py +++ /dev/null @@ -1,172 +0,0 @@ -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() - - -def ensure_storage(base_dir: str): - os.makedirs(base_dir, exist_ok=True) - - -def _path(base_dir: str, symbol: str, timeframe: str) -> str: - safe_symbol = symbol.replace("/", "_").replace(":", "_") - d = os.path.join(base_dir, timeframe) - os.makedirs(d, exist_ok=True) - return os.path.join(d, f"{safe_symbol}.parquet") - - -def read_candles(base_dir: str, symbol: str, timeframe: str, start: Optional[int], end: Optional[int]) -> pd.DataFrame: - p = _path(base_dir, symbol, timeframe) - if not os.path.exists(p): - return pd.DataFrame(columns=["timestamp", "open", "high", "low", "close", "volume"]) # empty - 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: - df = df[df["timestamp"] <= int(end)] - df = df.sort_values("timestamp") - return df - - -def upsert_candles(base_dir: str, symbol: str, timeframe: str, candles: List[List[float]]): - p = _path(base_dir, symbol, timeframe) - new_df = pd.DataFrame(candles, columns=["timestamp", "open", "high", "low", "close", "volume"]) - with _lock: - if os.path.exists(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") - 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 write_candles_snapshot(base_dir: str, symbol: str, timeframe: str, df: pd.DataFrame): - columns = ["timestamp", "open", "high", "low", "close", "volume"] - if df.empty: - safe_df = pd.DataFrame(columns=columns) - else: - safe_df = df[columns].copy() - safe_df = safe_df.drop_duplicates(subset=["timestamp"], keep="last").sort_values("timestamp").reset_index(drop=True) - p = _path(base_dir, symbol, timeframe) - with _lock: - temp_path = f"{p}.tmp" - try: - safe_df.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 candle_path(base_dir: str, symbol: str, timeframe: str) -> str: - return _path(base_dir, symbol, timeframe) - - -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 - 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]) - - -def read_candle_exact(base_dir: str, symbol: str, timeframe: str, timestamp: int) -> pd.DataFrame: - p = _path(base_dir, symbol, timeframe) - if not os.path.exists(p): - return pd.DataFrame(columns=["timestamp", "open", "high", "low", "close", "volume"]) - 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/datasvc/app/worker.py b/datasvc/app/worker.py deleted file mode 100644 index 7a5863a..0000000 --- a/datasvc/app/worker.py +++ /dev/null @@ -1,27 +0,0 @@ -import asyncio -import signal -from threading import Event - -from .main import logger, run_engine - - -async def _async_main(): - stop_event = Event() - loop = asyncio.get_running_loop() - for sig in (signal.SIGINT, signal.SIGTERM): - try: - loop.add_signal_handler(sig, stop_event.set) - except NotImplementedError: - # 信号处理在某些平台(如 Windows)不可用,忽略即可 - pass - await run_engine(stop_event) - - -def main(): - logger.info("worker 进程启动") - asyncio.run(_async_main()) - - -if __name__ == "__main__": - main() - diff --git a/datasvc/docker-compose.yml b/datasvc/docker-compose.yml deleted file mode 100644 index 2c3cc02..0000000 --- a/datasvc/docker-compose.yml +++ /dev/null @@ -1,19 +0,0 @@ -services: - datasvc: - build: . - container_name: datasvc - restart: unless-stopped - environment: - - EXCHANGE=binance - - TIMEFRAMES=1m,1h,1d,1w,1M - - START_FROM=2025-01-01 - - POLL_FACTOR=0.5 - - DATA_DIR=/data - - TZ=Asia/Shanghai - - WS_ENABLED=false - - DATASVC_ENGINE=1 - ports: - - "9000:9000" - volumes: - - ./data:/data - diff --git a/datasvc/pairs.json b/datasvc/pairs.json deleted file mode 100644 index ffe4f03..0000000 --- a/datasvc/pairs.json +++ /dev/null @@ -1,5 +0,0 @@ -[ - "BTC/USDT:USDT", - "ETH/USDT:USDT", - "SOL/USDT:USDT" -] diff --git a/datasvc/requirements.txt b/datasvc/requirements.txt deleted file mode 100644 index 346bb9a..0000000 --- a/datasvc/requirements.txt +++ /dev/null @@ -1,9 +0,0 @@ -fastapi==0.111.0 -uvicorn[standard]==0.29.0 -ccxt==4.4.27 -pandas==2.2.2 -pyarrow==16.1.0 -orjson==3.10.3 -technical==1.5.0 -websockets==12.0 - diff --git a/交易记录/TradingRecord.xlsx b/交易记录/TradingRecord.xlsx deleted file mode 100644 index d8d7106d99476c3cad1ee5e0578700ae0447463e..0000000000000000000000000000000000000000 GIT binary patch literal 0 HcmV?d00001 literal 10158 zcmeHtWm_EC)^_8K28SSNT!I95w*+k>Y3EAAWjaa*u09?Le>z- z-!Cye9NfIkHX{yE@8UzJT|lUrJH589S@i0Nu3wfdH|Of*xMXj58%K~LQZ z_49zDBmu)H3N{+Ve>RP(T$;M>?z_nP+rWfEGa(QBosWX;YoV6GbodP6h22qrQ4k%*Q3EG4 z8|P=NKhOW;=6^9E|MKZ&iHhLgxp^Dbe=)Tbk1ZkZCjO$CQq9+2cJWbtOfC)SQae2v zwi-zg!gIe?-`md%OF}U_{gmfxTxC(X_=408Zsnn=*N(4Ip3*s`$T*g*_TacqolTvm zKbP}haBYoyTKb{ro&3Nuwamn^R5kVpt0oC5UOrhUp>TSDPOqZ&SL4eHxG9O3`{kik zjl4NKNh2A4Gbshzn4;nQ3VY+}1pUs&7IPKe{Wg?mS0w5$E%~jgjI*5ts5}fzZ8}b* zGTQO(yjheo`&Fs9a4$JVp7+t_ocZb0aUBn4dGzq2c9-@a4FpDkioU|!)jvs+hyvY1 zgxwKE002k;Ksa}sXTS5r&E5%OY;OB|1#3>?g#!Jhr^KFSgYt$Wxoq%T75gT~q% zIee5u9nHs7(`++*{T-lu3dSh9KHm-@ssaNY@pd)rFS;(CL^Ln>xs_kifJiI|Ol8pm z%UsP&AV`Ow28XD_J1pgGvS^VIkY+b$G@CW86V?k2--v`483)04FWTp0%qxdcIttJ6 z6=xpCn`}eoaz%_4ostS3jgN3sh)&>#sRoEb8z~~+#SK3|noF#?gaz-$5;o8YIgUnq zk)#6J=roRW=@W3incWP%IVT&DNV?Jbd_0swT*LjiEmIU(vInoPnBW_qc5ssco{UN9 zTMQo8K>qp^R;Na51`N-q%d>r7MrzsZA3FPMXE#M*OaD&{Ay=t(Lt*#{gvCaWU{Am> z^jGwhsSnw$^I*3T&AMT`Xfnt_nb92y8+=Cr15__s2G|X0o{r1lw&D{A0?!#E;eGtG zqq5OYZRg*|jZxb24}8>euZ7I!4LtU_HUCb!bF^Pl9OFgzF|rX44>ViQIUl{L?j^{+ zJo1u>Ek*{=K=uwzF>wz6#J0~D#pU#|dgG#DsP;Pqk!2g*RdP-na zsGL@$?j#Fhu7@TET?Pq%!+vh4MtPGA@RJB}3qwitC!t-^>CfEZhTO?w?M$W8ygA%+ zA`X}YZg*Xu-D?5i86_3pf{?jm=+NWrby>F*F^wv#EXokQeJ+hikuCU86$GDitJAsn zp?d&EEfqcKeQl_t4e|`@GlK5ofw+5Zs@#h;xc5xjSGvp(s2B^&(V^lxI!u%_hsl;- zQ!-4RkQ1s8vCfQhZy6uQ@|397TTRv_$?D6dfWP1{UnREAV&550YjSI`Hv|v~%zwU8 zAm(SXG#KOudibL}Uj_+g|3EZSfDA~<(?+>7iHBu;Xj6A45q!$eu*hctx{+hEWF95z zL8Bg$@SMm}(I#mvqNnp*7C+iT7&;hDUjb>+sZh>7XS5tve$>00DHr)Zv0WYiIJb09CzGTf-@p-Hb@6&g+sTHaS_}tB7UOhXWNJA`s>l36}qgzcXJv~XwMuS=% z4p~3uYCuiL9RdpH=GklI6jsEL&A?yQAygE}a2{~;aqte6Vw#Iz-F7)}szd8Z?o!gT(*c~SriSd6Yu(PF+ zlbPvD7bhz_3+JCcz&lA%eqD^X?-1QLjLl84l%a;^n1Wfqr3pk)my5}jWM}cQtP#I+ z?|cZYVQFewM|+k4hly#72lI|2UDoLb)83?}r7C4u`qiBum~KQNuioBM;}Y#ei!&Z0 z5gnC649YQ1GM~`7Sq?ra73+|aw&I3RFa+s^cfR4c5bsI@E=p@5KbCL;N#hi9C~F=E zCIp2i*-}@Q27^xosfIV+z3F`njL0S+!iykK(d0%9hN_53AxFuJF;MKw!&S7iGT=_* zM=+dTfIKf#>zf1g3i7LIaHVnSw;mJ4S=Arg8JqqpG<(~rlQ!V@mHt7~x}=iP&ECrPgB zeP@ZiE@*eYUUsvW&$Qe)IcRzit0$5gjT9_jY+I#|xGy@Fws?7y?KQQimwDlxUQY_7 znWpl#8iikyu8R1{lYBMR$;4jus2Z-IbE*me|Ws{?*=>JKrvVgw$}5g008EH z*4xFy#_YF+uG1KdUEslABH9wgeq~xXX2Lw>N|S%)j)?do2%ZzECZF~*U00Z#%F8C{ zi@PAZk3%5SiAUD$xvB4M20eyM-*_77>FZ^Cxf`fmml2UIsVs5Oa+2vI602yJL=M=7 z`n9mM@XM4_)=<5YeX}ZTdt7_FTRA?x$KN`qVvQ}<&(kev@+SYVaU&LOT)!w(&iy45 zOL##3*K6Bc#$ugjSqD+Hz#gJt**)x(_ujy8itE=i?)kWw>^9weoHuaj)GDlR?eoV? zy92C2!!IJP*~V~o&N>Hjd`997^EJv#hLBdSODeQwmkXPZh(s8{FZ^%$^4QLUw)XZ~ zeQCzO%$&x~ovsQcXE6j@sT2!XagdR#_qTfJi^nDrQ>P{~sn0>X4hO(l9N>rh-EzEO zuiKN&L#3kLGmw4e(W&kps3l9Besi=rAA5KjRJ>tT{tZe=101jOW| zn7t&6Bxtq7Sn?XL;4WS<(|vhSmSbLb+&% z`}A-KahQV_@4eEFEu#Xz%fy_j@kb)rSnxr3o;>Oi`@_&J?o1=ehxk(x=KTsAm{(47 zbUSRP%E+}uXrqn?8)2wH@P?}-;30+s`6DHmNv1&#*RI=jmL1I@GttS{q_DH}8l$>SY*3n~gCJUhdAu(_B;yBHHVc%ElylXGlDNw48L^Jt-%v=9*h=pwJp=Z=O)5LyF~;dr zxDb@f_6_YHn}vdZT+|t7g~w`$RtcIRhhx0`M&v3H9C0S@b@D!h5)(XWm9lD3skzKN zJtO^XC}?UiLvMUmx2Lq8Xrw}8_jNV!{mM?`fb=b>Vs*6{E>xUPw}FfmBK*WvXl+O~ zb#WSk55qje8$ zCiKWOT#QD^0FcDxFzs96riG}d?j~~1+C;!Gb~G(Tvw>apfHprbl4#`|%wl=ffj0zO z9b6wXd~ez2WY0=wqSFDude+hTv?0DoE3+RuqxyJz{sBL%j17i2H$fQmuBp2v0vDW25UkIy)%i z0|k8br&JT9k$esv=eNVTc0?C$iog26tNJ z&5|}gItrGnU_xDfqCJ4>>P*Wl*tPdz(tU6c{%&~lQLRVjx;RJo97Q{LqVrBD8`C^) z+d;&WIQIL=-lI24a#-M5R|19x=277%?KBK;lrlv{`4+Zo`PfX4yPnEm29XgXX^ezG zE6;l^KHwv$G?OwXD0eTUn3k>M{PEqCAuW0{bULq%o zF!92?gIJHUsv+9QzPHM(!&CdYpn+xGy5)+N3qj}nxSf|Kv=KwA^5C=Qj5|W7;Au$Q zlX=C-vf+fMN={z4iqYUly zFz@+)JW=-$JeD0%d18Jqwd0B9A)1u*c7P%H%HWB9f*|hi7>2?0;lW#N%x-}V<2Jb? zj+812pG7;wB*!+BPBE!&iHO)PvbFB<8uP-vuhY~lp{*Ji)k7e$FRAu=)7S4!Y9Yr+ zY?x@$I&+|<>ci7z$i*&N*$rn@qNVCx1vsqn)9R&0wg{VFz~`Nvjb(2G%^NqOB=}sR zQj*BsP!mEGF300}>5-=4HRtS4Q}k=S)MSoP@sw+>;r`g{ASyN~d}^|%!OSug_!0*~ zQWu^`Ji|`Gw5Zw|)Yy0u2hY*+DNTjDj3xWK}7UtTRWqfa3!v6?BLi6^o^9 z3>|H8S`EQ7(PmjL`7-Z7TP`^`s`=?=yqvrow@a94RNF`uX-|9jhnu4Zl%K2m+tSuBzEA+4 zDKg8Ad|}e)msMn{R_7Sy)Kq2otT5tVE&l_(<9&e~n(P%V?Nt+-zKpf8oHb6kMOj{a z%3A~Y^{s`ut<+5-Ga`$;f|S{xWrBZNy-p(RHYL~#dntaF&40@jT`bLP&7S>s{>|X+ zX%9t_^5VA>U5jEkyWMiH$I*UW8L>%PqA|^alGQitsc3L=B({*^!PD}esWRl{r%Koh zL&wDsFk4UKP-*J6L1Xkov()7qG#NH>I7BCG{R=vYek+aSM2r zJwXC9nqgV9q?-0!^?G-WM0m)S>H#_B{12mScEKlWUgH-GxVXNSSG0GVM2*h~$291q zHjuR??Q$MCniLGZwZ@!(Yyz&hYieB}d?cV9MH;fm;S#kv?7c|TY>?YgA z8(}J6PrYl_JX9Yf{5oo&SKCBh5*IR)qvlr?3c~ofB8^@t&M5M#J1on*&|-Hc2j4)q zdq#d=U%BxVqAVrr^dXcrq~6-CX`aBeYbzG-D@KrbYa-)_N<3!ip6W{lEh^p(o#+D& z5r99uDvj*gz$XmwD6r;36a)WC&jRfP5Z68wm(K8m8Y8Bh`emg^zfI6t6Ezmc-H3g z`lw>ngc?jz{9 zi6%C^Ql=Vd(IncH*iB%iIt_*Oo$(f_MraBi>W%eC2 z5;vlNuXFrmd$nXVd19eRQaXYxyjq2OuunB`E@?0t9V8u*k4G4^N`J)7E%X6ftLluAjB$;vz`(h(%L)V=-)d!?+)^DKn;n!xv8&6lm%yqhvw5M0+4`XfPCfmAp&PP?}|ifNluS%BU8;YD;7&}g&TeY^et+FaRzCNEwQ;>b~gs7ZIS9p zc;Kz)z~pOR+h=I-Rd~M*@v*n{V~X>}^C4^%NT85WC=d;UXelt9;Qp0VR;iB70W@Pa zgXa+fzKljIUedRiPa_dE^tW!GpSZy}BO}G>9^XI(G8oJ@r2=YEFF?!?Ps{AZ4u`}t z=V=dhXZ@+mhno`f34M`Dvtda3iB9j?`I)$NdW(0!XS%~5TiM{2>J3gx=)oY;r)CK^ zAH+#Y`w!G?jgPQ;d$3c;(LcDn>F@u3$??`0#7_i zv5|r~smg0&6$2~@e0_^x{ke3lF&C{d<(<@bIx>zMXl%9w1nwK>E7GkSVdC@Rq2^@J zO%sB`y7ZRl3bf&L8P6on;&eD=<7b={aV{_CK010>zqTrAmCj?Kj&`YyVH;-)BwA$_ z97?&-?){h^Hpu~QGHxPu;@pXOgaeVLqW=&;NB45HN6_x7-Wu5Wl$r~2g||iAJ>!nt zL_MiX^#0`s*T=Hic9(8ODUXS2!la-AuacBLBeN%@c|m(v3|Rb)+?RY>CNHCzBkO7f#l*ED2yF7=3p}#DwI}&aNY`D;ss58c1E`PZc*^*>CU7}(sr2OH_Xx>l9M%0^X=}GSez;#tW8+B zQyy|HtCK!a2%T&!Z@+LTYPTqLN7mM=EYxu_12MppJeRhwW%z`TE9gc0v%cuO7GvI( z%euV&yx1~jj^Y*|;!tGW)PK&=wso}*F|QY9Kx|dY9o?H22+=RsJ;d}}sv08` zI%=iH#G0S7Q{Lr8t9NSD1zOXZ`Q~aDp(M)*7B!OI<=jy)nHX*``k1x>F07jy{=LrA zrLPT)hk0gFSniJvEB2e(n%oy-+>QcUCrR}55El0nWg@RwHtjnX zq81|zTg?xSa!s%Ju-TJ)rf8aZFdcVnLCT)34=5|dmS#W9)M1&gb#abQFlqAD)G27-BT{?cwcd%F313j^_Vn2g zR~7x#D+3E|TG0xf2jqVXBKzmR-O9sI>;+4#asG;86MHAK|G^lB-~SxZiF)$gY{Vgl zNDq<&J{i6?CeQ|iWIcZTMdcRsfd!UX?|7-^JFf}XQUZb&g<+4knJ0O<|-VMX?=3hUpU>u+8mX zHe)9g=m81T9CAU}4`U3OCWksQ(Jc<3XA>mn9b&C+@YC8@7+F7T4+Y>F zA+_`HxeR(tF>VYiA&9T*UqR(a>St=)m>-j&j6{o2Mv;u(ZAd|<0QYwI@*rfft7xL z|28;y7FgZ?-}g5C0j593f7$7zEdO@}f8TiVhv3gK59WCOvJK-`!M``7|0*~DvjqSD z*7UD>er-tpsR49_k?>b7zoyxLYT+gNP0O!I_pb{69w_`N4gfSf0RVmt8-A7kyMg|z uv^V8nr2jD2zl#6Lu0IvrQ~x&mzi~}j9tnnxpF4ce0c|jkFi-RI>i+?ZbEb0u diff --git a/交易记录/交易规则.docx b/交易记录/交易规则.docx deleted file mode 100644 index ac67805b1dd8a89c696e34c732d7a8cb29884a15..0000000000000000000000000000000000000000 GIT binary patch literal 0 HcmV?d00001 literal 12767 zcmeHt1y>!**6oH62oPL@26uM|5NzWv!QI(71c%`6?!jGxyGw9)3&Gvp-ahAk=bhY~ z``#b;ZjaGDMt9X*U0q|%UbSW|ImuU0=m2QI8vp=61X# zWGG((ox&arwi~3%Q%rx$a~m3@yfQ>I9Eg~VMbCGIzCc`)HR%i z!0!>0<;P9y>Z`<6G8-z{S{PhJ?XRftu~lrHMQhvHM=)wE$d()Y5 z?ul|yc1ynY_@@y#*!)RypCVum9S9v*p z!fE&bv{-?SHhq=lay_;^znB)1J8F9gVV1MO-?v%;DjPNX0btsj^3qBbn>87Oxw8EF z5&!%4(u3vgD*)i-;IIUBh1Z9MfR@m7)?QV^S7H2@be@h;d^Q_I!x2fCL#!rJJ;#7map=uZ*V>r;! zw?RoG9cWwB!q-~`@MgqRW77tt)j1iq%5yhKo-ALJnPz6q#WdVtig`ua0`Qwg-+96r z^M6ZZp^J*_TCLLRQQD7fU=BB@x}wDaF*1;)MeCQs;heTi?7QZ)7xV`E)52!K1~WXk z8*8GfKktNWzvE{=AyTW+S+~Yb(-OD<2d01WW7e=*{t`SHOknol0p373S=$*h{EH-p z)&`E2;OO;7*!qhw5a8Gaw)4OHD2*8cI}G^W?f>kXw57IVs9pSy6SK@a zxeUhMausB!E4Co>x3`3K)+wt$d^xONaq}q>@!PzF4exPEsN!oFzK&&-6pq>vQIQN8 zCk!=UE|8-w?#Ssa>d?z5Y1gXCn4*8KBkVh=OtNoKo6I_kplPi~i_F}mB$DNGy#p_0 z{^>QD883iCa0`eG06+md!=GOBtG^toYgtdJV7Fz~ydd8ZfGEsRi51ZAz4Kv<)awrh zYSVj)j|>*9gq$$c3!6MEPnEeTeVG?Ambf%6p@I1luKGkJt+s|A0Q4EP7doh zHzX*+dzWuI*r=$~bDjcy4a9aymuIZR9S=buxXe!I@=R`_ZW&j6KG8gRAJA*wovQj)HOogZ?TMI-Lv2 z@S2jsm}~+!@zX9r;?UCMR3VKKd?nUB!V}qqktZv$^wlYlQP>{Hfd0tNOO7h6H)NwT zSraUIll1Y!7aN$5VHmD&&5*j4(H^~%Eje+#J|qz#0aIhuC3!0h*BI1Ci%fxLS+FO8 z*U}b8J7@InZ%cvGwG~iDPv51sv0p3r^pLj?42l`$3+qBLq5Ckt;j#w2_}8O$g+^cp zDCoi2+TRxYor_blHJm6wi6Dgc7K){urU38WvXJaI>`bf1apRu3lG{uzIHyVU<1-ow z_F~Pnn(}v6ssW`?qOk><=txr9!$D!>8nH{r)3&4s&On){&JSazuIgymNTp7eYs;jkF)GD?U zt4sA(I?URoI7{{W{;wXcluwsDn@nkW%hyWPAG|ukZ7pXW?MrOi8!eQ|arrA=69{*< zx?HHhvNB=OUC&hxv}%`mco(lO5uku z5LY8Ut@bET<*n37vUO!681-126>F|mnk#G^C+qVL(CvHNQC{}dIB_^$Slb@6>lHq3 zuFH&1yo8?JhZ6nmk~fk&z5~0q_=dyB!>R3-f*|Omm|JqkA=qiu=ct+VdC< zsCDxPm2(=X|7_8$zZI zUEdj2mWa$x*|=Jx70FX+FO;{cGo4t$G(>SrY*P8d`leZNYICJjWIHv!`u$jOE42F#)rzpXPeVWH#rCp{d5<8ONWlPlQ77$#qH)mZ0g&Y>Re82JH=Oz`>Fhk@+uCEi3SZBmzgWy8@*}hk!Eb|c&=%$gtbfV&@ zk-f5hs@wxHbKX$*{4$Py7Ca=-PbElxW8}#tf@+qgi+RR=VicOMRBwlx4bqGezwIbJ z4jC(;8+rO9c zG5mRD`eQk9s3sYTD~j9(b50o0IOICAt5?c4DRyo#$G%$Y8`-L0kzbZ!HnFl;%NU@W zX$K3H7A5qyaSK`SEmO~C452}#cp72!K7lw3-BnMdmM?6xk1~H-@pDsgQxmKG`p>5% zZqX$8u9yf7-cPppX<_q8qDFw!wyShM@IY6BMliuL6}B=ZN1=R1 zjFWIO$#eg>u0!mPgDu}yp=2N=XwW}Ze2)r`B6o~Ov}miB&cwslag%o-M@lN?upIJi zKg5GIXhB2`Z)k5P^i&_G0T))lM%UlnUHf(CPCKW0wgN)QY6E&iz+hbUU<$yB(hWYAc*T*GCzkaP~9-pQtkmo=0x?cF8V z5CWL)5R0RV(i=Pb0KdQ^lLe-6Hj{?EKQ*S*gAAal9jsIceA(aFOSq=wj+B?-#>WXR zBkJu(A(1hz+R%PoW`%gZ-y513nbGck`7qQUsl#{g*mzyHY!v6`(#XN@aWOWmrv*WEt_4_LA9L~oE%zU8 zzYBLQe4NflRqPH+9M{aAAH?ouja(LCu8LmPIgQq{A?lnW${&6g_6?am8P{EJycOE4 z@A#UxDM*?$(a>Dzn7Hr!J6|d^N{ETdC~5HdJdHjaQAk>I8kTk#26f0h!)r=Vl6ACw*zN^Q()k-zLAb+LF|BmZzqSaIzQw=shx zOE*7Ns*A|;ADatN`i-g@kRFK`g@CB}mgwXs(U)e$~;+x(P{^m0_jUu3B(J8v4RYFtOM zqGN0G{gv4AjRFwa3VG6Ct3b(fYen(KA)0B#a*lS%f~LG+!^~Hhepw2&kumJY`8*z< zkxo5JBF<$TRh=2D24VLcbuefXIwgsB2YXa=1%{1;^`HUmb^vO`8T#IDP$?)tydQrc zFN>eFdblCVgESjGzhLsIUs(TG4ppdhXXRsTm`3falG~&3KrKyX#@6Mh)d6SbP7*c7 z7)L#NW&LzsB2~wEJh2ne4TGJi(K|Y*fb&#kqYfq%JIfKK^~N8gr47;L38E9uYccrp z!X_x{tCV+!8C{NU(iHVbBM$bTjJ+C@yXCm{dg_mt{HKNK%`i3J@6XokL`FId?5G;K ztm($*ACG9xX>T=UYy8~A5{yunnfS zU!4uL{OpWG+>SJ}pC*x2)K!5Q;WI-sz>2mzq5Qd#A)Si?|3;wq#8qGfNg_UhpVGy# zLc#oQ#kJ_*;cCC5Lgl`7QHU`l>}PRiIj>x%yR)s_Adg@@&!UZebjA61yM zC&wkVKb>w6A%;GfP6r&hHqTMk9<1mM9NehqeIDP9zYqVqKzRMCeTIO~A)N%_y6LLJ z;;nifknN3xc24$g{bFAJrxM3j<9fE!CZxo=fb@#CPunvspP&C@ev}}Y=z9UqsNcf> zk{{U{IXIYEnb`jk0Bcm1td=;DJ#=(mI>#@l6Aabt{J8y(gj3ONHI%5XlVd1T-!MRm z_d13=-QX+eA6m7DREcP#nvV`AuiMaj-aO;xt^&hI^peyUQduMlcwN8!N?%P;`1aruIGY_DGsu zWQJ+D_KGthN${;t37uKIfx$~6LPr^i%k2Ug0k)%SGoy2-qm$sQEG7f&7{*C7bE4*)$4-UF@l^?Ybd#7&yF z<_`8A;^Z}O^l)f1uEv5mFN&ijV9=SxB-gAG3I>YW5@Aa?cYfi6j)77)fip?FoSv6A zv%}gZ;^3txRSo8$>zWnOT>EuQ-E)8c*4w30mE%*9H&3}Zv{S6AxgNRDyzlkU18nW! zlC{Y?4Q=7XCt1@s3!+$qd@ge{FG@I7Wn_p#83$rxJ)e2@OGJ=MjSHeagqWDIZW@1o zGc=rVPrEWuGZS=$Nz)@CK5lAxJwI<6ckU_P&@K7a9+{9tgRQ}sL-^))2aoCZScN1P!i^-y@}?>q z6+H4Al3w?p%{ZGfIA*&sNcN1|^}5FSU&RfzgxjBIC<$%l9RbLR*SzNyWVu z({&zaiK#Zn@>L7_<4I<>#p_s6Pm~-JUqmoNU57ot4r*f71a=ExaxFx!_6(52s2v)f z39vg>?Fl4ss_7eQU^}6*kMUk+@!+_SP-5_C>JXlq5ZC8XLyE{9s=280neYOGeLBax zGA0GN50OIsZB;_|q4%+5o)j|zJMWB~SGyIP*J{%7qZSz&)QJ5LAnfy}RxclnP)X$p z$e(G03p|l~y}bx%?Y%IDyI8eGy438hZpLwrO%8kc!HeL3Dkv&2pa^oXCcO>=0HFT0 zOg08@7b)rKgN*)MBgajcFEOGAoWQujeh!4FX<;3l9E{w z=pj{O&5sZiR7{X840NS^yEn!%V4}SC%9Z8j=WlBAM0Kn*ShOKyL6>HPOUgCxHcHNt zjlAo{d)lGDg3#OihIYy}cqiFaXck+T6}%{5z% z%ARx%hEdT1lBK9jRD8$Nl_RrM9kY}m$7h& z%R4m-N9V^_B8{-=F=ItK! z_C|A(O337jUWHRwf{{Zq@qE6%x_gLhOZJvTk?~IN@OnveG|R}a4+fRt7)^#BkhJkK zQ6*QL4=H%ikw~)Da%aFzGyGn+xLz+ zXk^R%`jY{g`wt@{FT#wrXMVw_aQmNgjc)xb>@T~mQ|(Dx;K1@9ysC9NyqyB9&qjhh z3;i#zax&6a{3Z0nO^%pDG9d$gxMiIn+s&0K~G)djSrNUR-0Ay zi7>^Bl51znO{u=RF|2#T=Xgi$GeqpNI^RK2)+%g4HA@r&v8{>JWQ|mN#Xr2e{Z`O| zn2F6SDU_?NFvsq`aMyU((_5Th*Bx3ayzJrKR8E%NHep~EAqy8%i|D+#cfPx(TZrpc zRa`r8b6*&Su&>+bVwng%M&tojGQ=rbJ)jVAXpzf2Gn2sX>vQeIBY0;2NmmbHU#?m3 z#4dw_Ix<)$2Wy>jcGfob40<*;e~=Ztm-Jt09PDsuik((VjBi_!=J~OkI3lWv=OIOU zqx5(6G+y;sj-^RabQy|A&3|2(i>$-_wD37aHZpR0iiBz?HE|SEG;F}R9&-jn59hGR z(RVcWGwa#zN*+D9gH5-!yOvWAm1F;jr!kyQm(;gN7dCPB=SOrEqvT5n9x)afp3xe< z&mZynJoc7mOPel)&DZ=;u#0XH7SW1P6+`ykr)KKCcC@ zTKUaP!V6>)&}b}VbuFf{inoVlqnEn~Vv7yAAn!H}VyD@#!o{o*kp<9X<>>er(|gbD z-^`l6&nUY$cDNjaG`9+;lOzZfg4&GiHulroG*6F zs=93RQ2aVoF2UNbMQHa^^XsVxp2_{76nEN~yWSzKpr0_Lhk{%{C2S?u#WRjE{O0ps(OLB00HI#%CB57u(ta% za{ZSKU?zYMzZiL0@E(8vn)D;0&6m%$kKvE8-#Mq9H)Ngy_!ogS+YkVIDPH&aJ!5l^k%T%p@e`I zo#PH6h745B6lSw^rch7?avx#g@WWx4c&pyjIZ)uH3?fy6pqH295i#5Ey%x_CCuzdK8bBkN*Mh7x{ z`c!2rNma$E{=hJ|5jPQg;t4YNRrRp)( z?WZ&Du>Md^iGYo^p8{WNi(kyDt-(|OPp_R`J$~T_|J!|R0RSX$k3YAp>>XS{Mt>$` zE$K^kOPr_yC-vrQ-gzY}@S=P?w&5&BJfn+j-}+gqIOdRH#Z`U%5^~*>3UmOuq~E3G zK$UYG!H&B~@1LIqSY_|G=H`@uV;>m&21lXYGC#whjOoV*fn3?0-EZ42T#zxyhEU0s z`^KQHTd*JRuTH!kwzg_0*GL(N5x$2^vGU;#588Qh<6)kK6_noQZhx2}lu^p4QfMD8 z=#JZx4ZV{oo#dkFRTYzOL%4BriiBY!KM|CV4Gy~Fz+h0wh-ddqvM&mzrZLvnrZr>G zGG(Ba-EItVk>3rMG3u2}{;rbLDg=-aNh=YQ9+Y)8*e6hnLtI;`^~$qD;Ov69hWNpW z2^(Rn$r8XlxsHw($fd3yQXkijn9a)9nruvMVo)URYD(Pd82vicOWHuk(T;fW-vXEE_@H|9xn+AjOH zoq?d*#=_0v-R-(m{E|+WQK7B8=F0Jb+sTnUl3Hj__HHvUB-5Fh`U=LxS8FQH^Wtoq z4~u#(p9|OhPNALaXHf&Dco)k~=z|%UiF+%2wCml%FPP$&I9oZjqHMqD0ry!s{SQ7W*@+ zrK5T|1blZTp+Lf!@d;5jhevsYbntB^e=$+$ofkW4m0!CG&~F5=&mPeqZW_@K0-uuL zlNx-&nnqAeOyJ29NO_9*x{wKOjR^g zlDTn-Q%|cIQ*}kXNf4Xa?Q;t&nQAWigs>D=WBJg-=QNWs&3!m7eerA;zmBxT&}9~4 zBChlc*fAGO`@x=kUy0~5N{w6tpXIb_s&}C#uxF=q!KATW)L9GJY42j+_~go{rIkEu zDTaYOZ;3*l4UFDT^a@JI^fcv&GY#$Q;HpZ!Om^tse}cxGYG<@z&B4h|0kSHvg`0Mm z<4iOt--)wXfp3=oxS6SW-nf1*VOHSmH1lolX?`V7X*)P*_PcRO?K%;Vu$nV-T1ks%; zEzg+EGed^eADi=wOn{-b3#RIl`-Y1sv7zPyOKX;;hnad7uiUX*1dM%+X$7FylDA5u zWA(Q+ZmPWcXG@>!T={9DKUV-zFBOV!gNv`E zI>rGEN@~S%VQPf1kNvWI&lJ)yV~!|kttgg9Toda=@XhE@Jjo%ssKhba#vxKsheC8G z9M9}3w6dnAH(V9vI#gUJwmp+BwY%?66Ed)LE+2&|ajIIpsx^}~1eNk#-+OZfB82wd z@_#typIIpi+WbC#5;{rXRU@f|D;;>yE;ZkAr9-UBc6Hpj$IbdWZcoF)*l15m9pp|9}g0co3FHiXUiJu z=OMIjVu|?v9Ex~^4oG$HT2WeBcJj!aTVv*fKI~(|E`3-fXw)XN7S3=8^Rx6k#y0b4 zLRr`jw5Ihkm-36PdsLfo&e-HJ8_HI;F{ri-t_cB2bjH@vE>=o75cF?@Ae0JHO&-ry zhBqSn8@|hS&JW1$9n$b;kLY@xo9JC1CW%kg)i=394nKTGSuxvwG{`F6xW=VOYg(Ij z%)ydT>&yX;3f!Q5!mSST!sC(U@D$#0>5ts~0C9(sZcJ>=X)bZD$UPWpiw1c=M8Af?b@_(DQJLllLM;{| zgIhGetcSsE*aH2}h=Bfa$T@;=|9TEaz@hxBo419VwR4m_(Q8YkppUp1#N)u)gI)qK zf-|S<7L1Sx`nApfO^rMe?h#v}ZXAr|wh&36|v5|}7q6Y498!7@?j{ynyTLBK<$%<(yW|M%!?woE)iCY_R> z2%ZZIZbF`p3(EIpp7vG;kFpV)auQn4tDDwBZL+2nr=2&}y~{Y71bS5AjF9BrWJ;Pf z-Nd2xdv;C|kDP8BE;Ep{ickD2ccdU5D=kPl^X=vn&aFr(&k>iHSHhu6M!V}4=E9+x zJG<+|=kD!X~Ry?@l1^}()Q_|=ch%5dWStm9gQ^&?qkXW{X-`LeWaKygaR9DDJw zn-j6&(9+a~j=VTIL6mk2YiR~st1i<&j;$qff!0*KL*WoVA^UM}ulI-lRn9ovgS?)C zP^w-YB1duNmv|ghpIqo9n<64=Rr4UYSR9!!)gmGY8pqqsLkcKy1jcm*A*IS`WDIh) z<#@fk0I1UGs4wwNJBSORu!%U?K6wF!j794|5oy>+b^q0hOv={bCy5~62?i!jf0Y1A zBEbG>l8y8ido$#A3tl%wf%Ir|~|~ zHD=$zYR!PZWN??_>*b~XY1h5Hss9uBy|yz7vD53oE3`_5jbgHj=?3#6WpmBiX==7| z^>fMal7y7pK4nR1g+SR+Pr<1FXY)pJ9(j@5aa1dK8o(w?z#Rt|z$VONa6ZwoiT`Cs z`>)17U$U9_Z;jIb96{3zkn2HtYEd2eqYW$64BTEuFt5DE=WGd~%m&ALHw&D!(ic-P7F^0yt`qKNll!tceyWCQ~m8IIT zcg{C=4X%b69DMcmC0LBhhf-Hn&wFfR&YzGbDrE276ep$LBXnCBHy zFQxyr)`ozj14~Z-xf1@5Wa!u9Kh(v`N&dTmf3Ms96&wprJAbL-{T=x41&zM~tHIgq z|F699cT2w)x&3Wx4_xm5kCM0F;lCF`{0)Z!PuidG-^(I?H}U&6(%&XvG5^Y7eo^** zhySiI{|%?X`5*YdDbBy+e^(R##(U%a5B@i8@plXVE}Q+02LKlF0f65{wBOPHE;Rg# v-hTfJ`Y+kxclhra-`@t(NdB6=|C|5GNxlZV#~(G~@PJORdisXpkGuZ^p#<=2