删除不要的东西
This commit is contained in:
@@ -1,26 +0,0 @@
|
|||||||
FROM python:3.11-slim
|
|
||||||
|
|
||||||
ENV PYTHONUNBUFFERED=1 \
|
|
||||||
PIP_NO_CACHE_DIR=1
|
|
||||||
|
|
||||||
WORKDIR /app
|
|
||||||
|
|
||||||
COPY requirements.txt /app/requirements.txt
|
|
||||||
RUN pip install -r /app/requirements.txt
|
|
||||||
|
|
||||||
COPY app /app/app
|
|
||||||
COPY pairs.json /app/pairs.json
|
|
||||||
|
|
||||||
ENV DATA_DIR=/data \
|
|
||||||
EXCHANGE=binance \
|
|
||||||
TIMEFRAMES=1m,1h,1d,1w,1M \
|
|
||||||
START_FROM=2025-01-01 \
|
|
||||||
POLL_FACTOR=0.5
|
|
||||||
|
|
||||||
VOLUME ["/data"]
|
|
||||||
|
|
||||||
EXPOSE 9000
|
|
||||||
|
|
||||||
CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "9000"]
|
|
||||||
|
|
||||||
|
|
||||||
@@ -1,145 +0,0 @@
|
|||||||
# Local Data Service(REST + WebSocket)
|
|
||||||
|
|
||||||
本服务基于 FastAPI + ccxt,自动拉取交易所行情、写入本地 Parquet,同时提供 REST 和 WebSocket 数据访问。
|
|
||||||
自带时间周期聚合能力:只需抓取 `1m / 1h / 1d / 1w / 1M` 等基础周期,即可自动生成 `2m/3m/.../30m`、`2h/3h/.../16h` 等衍生周期。
|
|
||||||
|
|
||||||
---
|
|
||||||
|
|
||||||
## 1. 环境准备
|
|
||||||
|
|
||||||
### 1.1 依赖
|
|
||||||
- Python ≥ 3.10(本地运行方式需要)
|
|
||||||
- `pip install -r requirements.txt`(包含 `fastapi`, `uvicorn`, `ccxt`, `pandas`, `pyarrow`, `technical` 等)
|
|
||||||
- 或者直接使用仓库内的 `docker-compose.yml`
|
|
||||||
|
|
||||||
### 1.2 关键环境变量
|
|
||||||
| 变量 | 说明 | 默认 |
|
|
||||||
| --- | --- | --- |
|
|
||||||
| `DATA_DIR` | 本地 Parquet 存储目录 | `/data` |
|
|
||||||
| `EXCHANGE` | 交易所标识(目前支持 binance) | `binance` |
|
|
||||||
| `TIMEFRAMES` | 基础抓取周期,逗号分隔 | `1m,1h,1d,1w,1M` |
|
|
||||||
| `START_FROM` | 首次启动回补的起始 UTC 时间(ISO 字符串或毫秒时间戳) | `2022-01-01` |
|
|
||||||
| `POLL_FACTOR` | 拉取间隔因子,实际间隔 = 周期毫秒 × factor | `0.5` |
|
|
||||||
| `REST_MAX_CONCURRENCY` | REST 历史拉取并发数 | `4` |
|
|
||||||
| `VERIFY_MAX_CONCURRENCY` | 校验请求并发数 | `2` |
|
|
||||||
| `WS_ENABLED` | 是否启用 Binance WebSocket 增量(`true`/`false`) | `false` |
|
|
||||||
| `REST_POLL_INTERVAL` | 实时轮询 REST 的间隔秒数 | `5` |
|
|
||||||
| `REST_POLL_WINDOW` | 实时轮询时拉取的最新 K 线数量 | `10` |
|
|
||||||
| `BACKOFF_BASE / BACKOFF_MAX` | 异常重试的指数退避参数 | `2.0 / 30.0` |
|
|
||||||
|
|
||||||
> 衍生周期列表由程序自动推导,无需手动写入 `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` 的实时推送。
|
|
||||||
-1455
File diff suppressed because it is too large
Load Diff
@@ -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
|
|
||||||
|
|
||||||
|
|
||||||
@@ -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()
|
|
||||||
|
|
||||||
@@ -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
|
|
||||||
|
|
||||||
@@ -1,5 +0,0 @@
|
|||||||
[
|
|
||||||
"BTC/USDT:USDT",
|
|
||||||
"ETH/USDT:USDT",
|
|
||||||
"SOL/USDT:USDT"
|
|
||||||
]
|
|
||||||
@@ -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
|
|
||||||
|
|
||||||
Binary file not shown.
Binary file not shown.
Reference in New Issue
Block a user