diff --git a/data_provider/README.md b/data_provider/README.md new file mode 100644 index 0000000..f8ce8b2 --- /dev/null +++ b/data_provider/README.md @@ -0,0 +1,82 @@ +# Chan 数据提供商 (Chan Data Provider) + +从 **Binance 期货** 交易所拉取加密货币 K 线数据,提供 HTTP + WebSocket 数据服务。 + +## 功能 + +- **多交易对**:支持 BTC, ETH, SOL, DOGE 等 9 个交易对 +- **多时间周期**:基础周期 1m/1h/1d/1w,可合成 30+ 种衍生周期(如 5m, 15m, 4h 等) +- **本地缓存**:CSV 持久化到磁盘,重启快速加载 +- **断线恢复**:交易所连接中断时记录断点,自动补拉缺失数据 +- **实时推送**:WebSocket 订阅最新 K 线更新 +- **内存服务**:启动即加载本地数据,不阻塞服务 + +## 启动 + +```bash +# Docker +docker compose up -d + +# 直接运行 +python main.py + +# 或指定配置 +CONFIG_PATH=./config.json python main.py +``` + +服务默认监听 `http://0.0.0.0:9009`。 + +## 配置 + +编辑 `config.json`: + +```json +{ + "exchange": "binance", + "symbols": ["BTC/USDT:USDT", "ETH/USDT:USDT"], + "start_time": "2024-01-01T00:00:00Z", + "timeframes": ["1m", "1h", "1d", "1w"], + "data_dir": "./data" +} +``` + +| 字段 | 说明 | +|------|------| +| `exchange` | 交易所名称(ccxt 支持即可) | +| `symbols` | 交易对列表 | +| `start_time` | 历史数据起始时间 | +| `timeframes` | 基础周期(从交易所直接拉取) | +| `data_dir` | CSV 数据存储目录 | + +## 可用周期 + +### 基础周期(交易所直接拉取) +`1m`, `1h`, `1d`, `1w` + +### 衍生周期(内存中合成) +| 基础周期 | 可合成的衍生周期 | +|----------|----------------| +| 1m | 2m, 3m, 4m, 5m, 10m, 15m, 20m, 25m, 30m, 45m | +| 1h | 2h, 3h, 4h, 5h, 6h, 7h, 8h, 9h, 10h, 11h, 12h, 16h, 20h | +| 1d | 2d, 3d, 4d, 5d, 6d | +| 1w | 2w, 3w | + +## 数据存储 + +数据以 CSV 格式存储,按时间周期分目录: + +``` +./data/ + 1m/ + binance_BTC_USDT_USDT_1m.csv + binance_ETH_USDT_USDT_1m.csv + ... + 1h/ + ... +``` + +每根 K 线包含:`timestamp`, `datetime`, `open`, `high`, `low`, `close`, `volume`。 + +--- + +API 文档请访问 `http://:9009/api/docs`。 diff --git a/data_provider/api_docs.html b/data_provider/api_docs.html new file mode 100644 index 0000000..2c1c33a --- /dev/null +++ b/data_provider/api_docs.html @@ -0,0 +1,517 @@ + + + + + +Chan 数据提供商 - API 文档 + + + +
+ +
+

Chan 数据提供商

+

加密货币 K 线数据 HTTP + WebSocket API

+
v1.0.0  |  binance  |  port 9009
+
+ + + + +
+

服务信息

+
+
+ GET + / + 服务基本信息 +
+
+

返回服务名称、交易所、交易对列表、可用周期及就绪状态。

+

响应

+
{
+  "service":        "Data Provider",
+  "exchange":        "binance",
+  "symbols":         ["BTC/USDT:USDT", "ETH/USDT:USDT", ...],
+  "base_timeframes":  ["1m", "1h", "1d", "1w"],
+  "derived_timeframes": ["5m", "15m", "4h", ...],
+  "timeframes":       ["1m", "1h", ..., "5m", "15m", ...],
+  "ready":           true
+}
+
+
+
+ + +
+

健康检查

+
+
+ GET + /health + 存活检查 +
+
+

返回服务健康状态,与 / 相同结构,适合负载均衡探测器。

+

响应

+
{
+  "status":  "ok",
+  "exchange": "binance",
+  "symbols":  ["BTC/USDT:USDT", ...],
+  "ready":    true,
+  ...
+}
+
+
+
+ + +
+

可用周期

+
+
+ GET + /timeframes + 列出所有时间周期 +
+
+

返回基础周期(交易所直接拉取)和衍生周期(合成生成)的完整列表。

+

响应

+
{
+  "base_timeframes":    ["1m", "1h", "1d", "1w"],
+  "derived_timeframes": ["5m", "15m", "4h", ...],
+  "timeframes":         ["1m", "1h", ..., "5m", "15m", ...]
+}
+
+
+
+ + +
+

查询 K 线

+
+
+ GET + /api/candles + 获取 OHLCV K 线数据 +
+
+ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + +
参数类型必填说明
symbolstring交易对,如 BTC/USDT:USDT
tfstring时间周期,默认 1m。支持基础及衍生周期
startint可选开始时间戳(毫秒)
endint可选结束时间戳(毫秒)
limitint可选限制返回的 K 线数量(返回最后 N 根)
+ +
若不传 start/end,返回内存中全部数据(可能很多),建议搭配 limit 使用。
+ +

请求示例

+
# 获取 BTC 最近 100 根 5 分钟 K 线
+GET /api/candles?symbol=BTC/USDT:USDT&tf=5m&limit=100
+
+# 指定时间范围
+GET /api/candles?symbol=ETH/USDT:USDT&tf=1h&start=1704067200000&end=1704153600000
+
+# 获取 4 小时周期(衍生周期)
+GET /api/candles?symbol=SOL/USDT:USDT&tf=4h&limit=50
+ +

响应

+

返回 OHLCV 对象数组:

+
[
+  {
+    "timestamp": 1704067200000,
+    "datetime":  "2024-01-01T00:00:00Z",
+    "open":      42850.12,
+    "high":      43100.00,
+    "low":       42780.50,
+    "close":     43050.80,
+    "volume":    125.34
+  },
+  ...
+]
+ +

字段说明

+ + + + + + + + + +
字段类型说明
timestampintUTC 毫秒时间戳
datetimestringISO 8601 格式(末尾 Z)
openfloat开盘价
highfloat最高价
lowfloat最低价
closefloat收盘价
volumefloat成交量
+ +
+
+
+ + +
+

WebSocket 实时推送

+
+
+ WS + /ws + 实时 K 线订阅 +
+
+ +

连接 WebSocket 后,通过 JSON 消息进行订阅管理。服务端在数据更新时主动推送最新 K 线。

+ +

客户端 → 服务端

+ +
+
订阅 K 线
+
{
+  "action":    "subscribe",
+  "symbol":    "BTC/USDT:USDT",
+  "timeframe": "1m"
+}
+
+ +
+
取消订阅
+
{
+  "action":    "unsubscribe",
+  "symbol":    "BTC/USDT:USDT",
+  "timeframe": "1m"
+}
+
+ +
+
心跳 Ping
+
{ "action": "ping" }
+
+ +

服务端 → 客户端

+ +
+
订阅确认
+
{
+  "type":      "subscribed",
+  "symbol":    "BTC/USDT:USDT",
+  "timeframe": "1m"
+}
+
+ +
+
初始快照(订阅后立即推送最近 500 根 K 线)
+
{
+  "type":      "snapshot",
+  "symbol":    "BTC/USDT:USDT",
+  "timeframe": "1m",
+  "data":      [ ... ]
+}
+
+ +
+
K 线更新(增量推送最近 2 根)
+
{
+  "type":      "kline",
+  "symbol":    "BTC/USDT:USDT",
+  "timeframe": "1m",
+  "data":      [ ... ]
+}
+
+ +
+
Pong 响应
+
{ "type": "pong" }
+
+ +
+
错误消息
+
{ "type": "error", "message": "..." }
+
+ +

JavaScript 示例

+
// 连接
+const ws = new WebSocket("ws://localhost:9009/ws");
+
+ws.onopen = () => {
+  // 订阅 BTC 1m K 线
+  ws.send(JSON.stringify({
+    action: "subscribe",
+    symbol: "BTC/USDT:USDT",
+    timeframe: "1m"
+  }));
+};
+
+ws.onmessage = (event) => {
+  const msg = JSON.parse(event.data);
+  if (msg.type === "kline") {
+    console.log(msg.data); // 最新 K 线数组
+  }
+};
+ +
+
+
+ + +
+

时间周期参考

+

以下是完整的周期对照表:

+ + + + + + + +
基础周期合成衍生周期
1m2m, 3m, 4m, 5m, 10m, 15m, 20m, 25m, 30m, 45m
1h2h, 3h, 4h, 5h, 6h, 7h, 8h, 9h, 10h, 11h, 12h, 16h, 20h
1d2d, 3d, 4d, 5d, 6d
1w2w, 3w
+ +

衍生周期由对应基础周期的 K 线通过 OHLCV 聚合合成,查询方式与基础周期完全一致。

+
+ +
+ Chan Data Provider — Built with FastAPI + ccxt + pandas +
+ +
+ + + + diff --git a/data_provider/homepage.html b/data_provider/homepage.html new file mode 100644 index 0000000..c3814f2 --- /dev/null +++ b/data_provider/homepage.html @@ -0,0 +1,417 @@ + + + + + +Chan 数据提供商 + + + +
+ +
+
+ +
+

Chan 数据提供商

+
加密货币 K 线数据服务
+
+
+
+
+ + +
+
服务状态
加载中...
+
交易所
+
交易对
+
基础周期
+
衍生周期
+
数据就绪
+
+ + + + + +
+
+

📈 交易对

+ 0 +
+
+
加载中...
+
+
+ + +
+
+

⏱ 时间周期

+ 0 +
+
+
基础周期(交易所直拉)
+
+
衍生周期(内存合成)
+
+
+
+ + +
+
+

⚡ 快速查询

+
+
+
+
+
交易对
+ +
+
+
周期
+ +
+
+
数量
+ +
+ +
+ +
+
+ +
+ + + + + + diff --git a/data_provider/main.py b/data_provider/main.py index 7392878..462ddd3 100644 --- a/data_provider/main.py +++ b/data_provider/main.py @@ -19,6 +19,7 @@ import ccxt # type: ignore import pandas as pd # type: ignore from fastapi import FastAPI, HTTPException, Query, WebSocket, WebSocketDisconnect from fastapi.middleware.cors import CORSMiddleware +from fastapi.responses import HTMLResponse import uvicorn from technical.util import resample_to_interval @@ -405,11 +406,23 @@ class DataProvider: return results def initialize(self) -> None: - """阻塞式启动:加载本地、从倒数第二根或配置起点补历史、写盘并 set _ready。""" - logger.info("开始初始化数据提供商") + """快速启动:仅加载本地磁盘已有数据到内存,然后立即标记就绪,不阻塞服务。""" + logger.info("开始加载本地数据") for symbol in self.symbols: for timeframe in self.timeframes: existing = self._load_local(symbol, timeframe) + with self._lock: + self.data.setdefault(symbol, {})[timeframe] = existing + self._ready.set() + logger.info("本地数据加载完成,服务已就绪") + + def run_initial_history_fetch(self) -> None: + """后台一次性拉取所有 symbol/tf 的历史数据(从本地末根或配置起点到当前),然后落盘。""" + logger.info("开始后台历史数据拉取") + for symbol in self.symbols: + for timeframe in self.timeframes: + with self._lock: + existing = list(self.data.get(symbol, {}).get(timeframe, [])) tf_ms = TIMEFRAME_TO_MS[timeframe] last_ts = existing[-1]["timestamp"] if existing else None if last_ts is not None: @@ -433,10 +446,10 @@ class DataProvider: history = self._fetch_history(symbol, timeframe, fetch_since) merged = self._merge_candles(timeframe, existing, history) with self._lock: - self.data.setdefault(symbol, {})[timeframe] = merged + self.data[symbol][timeframe] = merged self._write_to_disk(symbol, timeframe, merged) - self._ready.set() - logger.info("数据初始化完成") + self._notify_update(symbol, timeframe) + logger.info("后台历史数据拉取完成") def resample_df(self, df: pd.DataFrame, interval: int) -> pd.DataFrame: """将基础周期 DataFrame 聚合为 interval 分钟周期(freqtrade technical.util)。""" @@ -788,8 +801,15 @@ def create_app(provider: DataProvider) -> FastAPI: loop = asyncio.get_running_loop() ws_manager.set_loop(loop) provider.on_update(_on_data_update) + # 快速加载本地数据后立即就绪,不阻塞服务启动 await loop.run_in_executor(None, provider.initialize) provider.start_background_workers() + # 后台拉取历史数据补齐(不阻塞 HTTP/WS 服务) + threading.Thread( + target=provider.run_initial_history_fetch, + name="initial-history-fetch", + daemon=True, + ).start() try: yield finally: @@ -840,18 +860,15 @@ def create_app(provider: DataProvider) -> FastAPI: data = provider.get_klines(symbol=symbol, timeframe=tf, start_time=start, end_time=end, limit=limit) return data - @app.get("/") - async def root() -> Dict[str, object]: - """根路径:服务名、交易所、交易对与可用周期(含 ready 标志)。""" - return { - "service": "Data Provider", - "exchange": provider.exchange_name, - "symbols": provider.symbols, - "base_timeframes": provider.timeframes, - "derived_timeframes": provider.get_derived_timeframes(), - "timeframes": provider.get_available_timeframes(), - "ready": provider.is_ready(), - } + homepage_path = Path(__file__).resolve().parent / "homepage.html" + docs_path = Path(__file__).resolve().parent / "api_docs.html" + + @app.get("/", response_class=HTMLResponse) + async def root(): + """服务主页。""" + if homepage_path.exists(): + return HTMLResponse(content=homepage_path.read_text(encoding="utf-8")) + return HTMLResponse(content="

主页页面未找到

", status_code=404) @app.websocket("/ws") async def websocket_endpoint(ws: WebSocket): @@ -921,6 +938,13 @@ def create_app(provider: DataProvider) -> FastAPI: finally: await ws_manager.disconnect(ws) + @app.get("/api/docs", response_class=HTMLResponse, include_in_schema=False) + async def api_docs(): + """返回自定义 API 文档页面。""" + if docs_path.exists(): + return HTMLResponse(content=docs_path.read_text(encoding="utf-8")) + return HTMLResponse(content="

API 文档页面未找到

", status_code=404) + return app