diff --git a/ChanEnum.py b/ChanEnum.py
index 9cc4bf0..76a7b6a 100644
--- a/ChanEnum.py
+++ b/ChanEnum.py
@@ -72,8 +72,27 @@ class Chan_KLC_FX(Enum):
BOTTOM3 = auto()
BOTTOM4 = auto()
UNKNOWN = auto()
-
-
+class Chan_MACD_STATE(Enum):
+ GW = auto()
+ GWK = auto()
+ UP = auto()
+ DOWN = auto()
+ CROSS0 = auto()
+ UNKNOWN = auto()
+ START = auto()
+ NEAR0 = auto()
+class Chan_MACDSEG_DIR(Enum):
+ ABOVE = auto()
+ UNDER = auto()
+class Chan_MACDHISTSET_DIR(Enum):
+ ABOVE = auto()
+ UNDER = auto()
+class Chan_MACDHIST_STATE(Enum):
+ UP = auto()
+ DOWN = auto()
+ PEAK = auto()
+ UNKNOWN = auto()
+
class Chan_BI_DIR(Enum):
UP = auto()
DOWN = auto()
diff --git a/ChanKLC.py b/ChanKLC.py
index 7172096..68f004f 100644
--- a/ChanKLC.py
+++ b/ChanKLC.py
@@ -71,11 +71,11 @@ class ChanKLC():
if self.high >= klu.bbup302 and klu.bbup302 > 0 and (self.klc_fx_type == Chan_KLC_FX.TOP1 or self.klc_fx_type == Chan_KLC_FX.TOP2):
self.bb_out = True
#print(self.end_time, self.high, klu.bbup302, self.klc_fx_type)
- if self.high >= klu.bbup30 and klu.bbup30 > 0:
+ if self.high >= klu.bbup30 and klu.bbup30 > 0 and self.next and (self.next.macd - self.macd) < 0:
self.klc_fx_type = Chan_KLC_FX.TOP4
if self.low <= klu.bblow302 and klu.bblow302 > 0 and (self.klc_fx_type == Chan_KLC_FX.BOTTOM1 or self.klc_fx_type == Chan_KLC_FX.BOTTOM2):
self.bb_out = True
- if self.low <= klu.bblow30 and klu.bblow30 > 0:
+ if self.low <= klu.bblow30 and klu.bblow30 > 0 and self.next and (self.macd - self.next.macd) < 0:
self.klc_fx_type = Chan_KLC_FX.BOTTOM4
if self.fx ==Chan_FX_TYPE.TOP:
if self.macd > 0:
@@ -85,7 +85,7 @@ class ChanKLC():
if self.macd < 0:
self.bb_out = True
def cal_macd_state(self, dir):
-
+ macd_state = 0
return macd_state
def cal_indicators(self):
for index in range(1, len(self.klus)):
diff --git a/ChanKLU.py b/ChanKLU.py
index 8144a64..c92cc38 100644
--- a/ChanKLU.py
+++ b/ChanKLU.py
@@ -1,4 +1,4 @@
-from ChanEnum import Chan_FX_TYPE, Chan_KLU_TYPE, Chan_K_DIR
+from ChanEnum import Chan_FX_TYPE, Chan_KLU_TYPE, Chan_K_DIR, Chan_MACD_STATE, Chan_MACDHIST_STATE
class ChanKLU:
def __init__(self, time, open, high, low, close, volume):
# _time, _close, _open, _high, _low, _extra_info={}
@@ -42,15 +42,24 @@ class ChanKLU:
self.fx_confirmed = False # 分型是否确认
self.klu_type = None
self.cal_klu_min_max()
+ self.range = self.high - self.low
self.body = abs(self.close - self.open)
self.upper_shadow = self.high - max(self.close, self.open)
self.lower_shadow = min(self.close, self.open) - self.low
- self.body_ratio = self.body / self.open
- self.upper_shadow_ratio = self.upper_shadow / self.open
- self.lower_shadow_ratio = self.lower_shadow / self.open
+ self.body_ratio = self.body / self.range
+ self.upper_shadow_ratio = self.upper_shadow / self.body
+ self.lower_shadow_ratio = self.lower_shadow / self.body
self.candle_dir = Chan_K_DIR.CROSS if self.close == self.open else Chan_K_DIR.BULL if self.close > self.open else Chan_K_DIR.BEAR
- self.range = self.high - self.low
self.strength = 0 if self.candle_dir == Chan_K_DIR.CROSS else self.cal_klu_strength()
+
+ self.ema52 = 0
+ self.ema24 = 0
+ self.macd_slop = 0
+ self.signal_slop = 0
+ self.hist_slop = 0
+ self.hist_state = Chan_MACDHIST_STATE.UNKNOWN
+ self.macd_state = Chan_MACD_STATE.UNKNOWN
+ self.macd_hist_gap = 0
#print(self.open, self.close, self.high, self.low, self.candle_dir, self.strength)
def cal_klu_strength(self):
strength = 0
@@ -599,13 +608,35 @@ class ChanKLU:
self.bblow365 = float(item['bblow365']) if 'bblow365' in item and item['bblow365'] else 0
# 设置指标后更新实时分析
self.update_realtime_analysis()
-
+ self.cal_macd_state()
def cal_macd_state(self):
- if self.pre:
- pre_macd_slop = self.macd - self.pre.macd
- pre_signal_slop = self.signal - self.pre.signal
- pre_hist_slop = self.macdhist - self.pre.macdhist
-
+ if self.macd == 0 and self.signal == 0 and self.macdhist == 0:
+ return Chan_MACD_STATE.UNKNOWN
+ if self.pre and self.pre.macd_state != Chan_MACD_STATE.UNKNOWN:
+ self.macd_slop = self.macd - self.pre.macd
+ self.signal_slop = self.signal - self.pre.signal
+ self.hist_slop = self.macdhist - self.pre.macdhist
+ self.macd_hist_gap = abs(self.macdhist - self.macd)
+ if self.pre.signal >= 0 and self.signal < 0 and self.ema52 > self.close and self.next and self.next.signal < 0:
+ self.macd_state = Chan_MACD_STATE.CROSS0
+ elif self.pre.signal <= 0 and self.signal > 0 and self.ema52 < self.close and self.next and self.next.signal > 0:
+ self.macd_state = Chan_MACD_STATE.CROSS0
+ elif self.macd > 0 and self.signal > 0 and self.macdhist > 0:
+ if self.signal > self.macdhist:
+ self.macd_state = Chan_MACD_STATE.GW
+ if self.macdhist < self.pre.macdhist and self.macd_slop > 0 and self.signal_slop > 0 and self.macd_slop < self.pre.macd_slop and self.signal_slop < self.pre.signal_slop and self.macd_hist_gap > self.pre.macd_hist_gap:
+ self.macd_state = Chan_MACD_STATE.GWK
+ elif self.macd_slop > 0 and self.signal_slop > 0:
+ self.macd_state = Chan_MACD_STATE.UP
+ elif self.macd < 0 and self.signal < 0 and self.macdhist < 0:
+ if self.signal < self.macdhist:
+ self.macd_state = Chan_MACD_STATE.GW
+ if self.macdhist > self.pre.macdhist and self.macd_slop < 0 and self.signal_slop < 0 and self.macd_slop > self.pre.macd_slop and self.signal_slop > self.pre.signal_slop and self.macd_hist_gap > self.pre.macd_hist_gap:
+ self.macd_state = Chan_MACD_STATE.GWK
+ elif self.macd_slop < 0 and self.signal_slop < 0:
+ self.macd_state = Chan_MACD_STATE.DOWN
+ else:
+ self.macd_state = Chan_MACD_STATE.START
return self.macd_state
def get_feature_data(self):
diff --git a/ChanMACD.py b/ChanMACD.py
new file mode 100644
index 0000000..76b1fbc
--- /dev/null
+++ b/ChanMACD.py
@@ -0,0 +1,117 @@
+from ChanKLU import ChanKLU
+from ChanEnum import Chan_MACD_STATE, Chan_MACDSEG_DIR, Chan_MACDHISTSET_DIR
+from ChanMACDSeg import ChanMACDSeg
+from ChanMACDUnitTF import ChanMACDUnitTF
+from ChanMACDHistSet import ChanMACDHistSet
+
+class ChanMACD():
+ def __init__(self, klu_list: list[ChanKLU]):
+ self.klu_list = klu_list
+ self.seg_list = []
+ self.unittf_list = []
+ self.histset_list = []
+ self.cal_macd()
+ def cal_macd(self):
+ last_seg = None
+ last_unittf = None
+ last_histset = None
+
+ last_klu = None
+ for klu in self.klu_list:
+ klu.cal_macd_state()
+
+ # 1) 只有当 MACD 已可用(非 UNKNOWN)时,才开始初始化段/单元
+ if last_seg is None:
+ if klu.macd_state != Chan_MACD_STATE.UNKNOWN:
+ # 初始化首段与首单元
+ seg_dir = Chan_MACDSEG_DIR.ABOVE if klu.macd >= 0 else Chan_MACDSEG_DIR.UNDER
+ seg = ChanMACDSeg(klu.time, klu, None, seg_dir)
+ self.seg_list.append(seg)
+ last_seg = seg
+
+ unittf = ChanMACDUnitTF(klu.time, klu, None, None)
+ self.unittf_list.append(unittf)
+ last_unittf = unittf
+
+ # 初始化首个直方图集合(根据当前柱体正负)
+ if klu.macdhist >= 0:
+ histset = ChanMACDHistSet(klu.time, klu, last_unittf, Chan_MACDHISTSET_DIR.ABOVE)
+ else:
+ histset = ChanMACDHistSet(klu.time, klu, last_unittf, Chan_MACDHISTSET_DIR.UNDER)
+ self.histset_list.append(histset)
+ last_histset = histset
+ last_seg.add_histset(histset)
+ last_unittf.add_histset(histset)
+ last_seg.add_unittf(unittf)
+ # 未就绪则继续等下一根;已就绪亦已完成首个结构初始化,继续下一根
+ last_klu = klu
+ continue
+
+ # 2) 过零切段/切单元
+ if klu.macd_state == Chan_MACD_STATE.CROSS0:
+ # 收尾旧段
+ last_seg.end_klu = last_klu
+ last_seg.end_time = last_klu.time
+ # 新段方向取反
+ new_dir = Chan_MACDSEG_DIR.UNDER if last_seg.seg_dir == Chan_MACDSEG_DIR.ABOVE else Chan_MACDSEG_DIR.ABOVE
+ seg = ChanMACDSeg(klu.time, klu, last_seg, new_dir)
+ self.seg_list.append(seg)
+ last_seg.set_next(seg)
+ last_seg = seg
+
+ # 切换 unit tf(以当前histset收尾为边界)
+ unittf = ChanMACDUnitTF(klu.time, klu, last_histset, None)
+ self.unittf_list.append(unittf)
+ if last_unittf:
+ last_unittf.end_klu = last_klu
+ last_unittf.end_time = last_klu.time
+ last_unittf.set_next(unittf)
+ last_unittf = unittf
+
+ # 3) 直方图集合(基于当前 unittf)
+ if klu.macdhist >= 0:
+ if last_histset and last_histset.histset_dir == Chan_MACDHISTSET_DIR.ABOVE:
+ last_histset.add_klu(klu)
+ else:
+ # 结束旧 histset(以前一根结束更合理)
+ if last_histset:
+ end_klu = getattr(klu, 'pre', None) or klu
+ last_histset.end_klu = end_klu
+ last_histset.end_time = end_klu.time
+ histset = ChanMACDHistSet(klu.time, klu, last_unittf, Chan_MACDHISTSET_DIR.ABOVE)
+ self.histset_list.append(histset)
+ if last_histset:
+ last_histset.set_next(histset)
+ last_histset = histset
+ if last_seg:
+ last_seg.add_histset(histset)
+ if last_unittf:
+ last_unittf.add_histset(histset)
+ else:
+ if last_histset and last_histset.histset_dir == Chan_MACDHISTSET_DIR.UNDER:
+ last_histset.add_klu(klu)
+ else:
+ if last_histset:
+ end_klu = getattr(klu, 'pre', None) or klu
+ last_histset.end_klu = end_klu
+ last_histset.end_time = end_klu.time
+ histset = ChanMACDHistSet(klu.time, klu, last_unittf, Chan_MACDHISTSET_DIR.UNDER)
+ self.histset_list.append(histset)
+ if last_histset:
+ last_histset.set_next(histset)
+ last_histset = histset
+ if last_seg:
+ last_seg.add_histset(histset)
+ if last_unittf:
+ last_unittf.add_histset(histset)
+ last_klu = klu
+ # 循环结束后,收尾当前打开的结构
+ if last_seg and getattr(last_seg, 'end_klu', None) is None:
+ last_seg.end_klu = last_klu
+ last_seg.end_time = last_klu.time
+ if last_unittf and getattr(last_unittf, 'end_klu', None) is None:
+ last_unittf.end_klu = last_klu
+ last_unittf.end_time = last_klu.time
+ if last_histset and getattr(last_histset, 'end_klu', None) is None:
+ last_histset.end_klu = last_klu
+ last_histset.end_time = last_klu.time
\ No newline at end of file
diff --git a/ChanMACDHistSet.py b/ChanMACDHistSet.py
new file mode 100644
index 0000000..48802c4
--- /dev/null
+++ b/ChanMACDHistSet.py
@@ -0,0 +1,15 @@
+
+
+class ChanMACDHistSet():
+ def __init__(self, start_time, start_klu, hist_set_dir, pre_histset):
+ self.klu_list = []
+ self.klu_list.append(start_klu)
+ self.ref_klu = None
+ self.hist_set_dir = hist_set_dir
+ self.next = None
+ self.pre = pre_histset
+ def set_next(self, next_histset):
+ self.next = next_histset
+ def add_klu(self, klu):
+ self.klu_list.append(klu)
+ klu.set_histset(self)
\ No newline at end of file
diff --git a/ChanMACDSeg.py b/ChanMACDSeg.py
new file mode 100644
index 0000000..44673ef
--- /dev/null
+++ b/ChanMACDSeg.py
@@ -0,0 +1,27 @@
+
+
+
+class ChanMACDSeg():
+ def __init__(self, start_time, start_klu, pre_seg, seg_dir):
+ self.start_time = start_time
+ self.end_time = None
+ self.start_klu = start_klu
+ self.end_klu = None
+ self.klu_list = []
+ self.klu_list.append(start_klu)
+ self.unittf_list = []
+ self.hist_set = []
+ self.seg_dir = seg_dir
+ self.pre = pre_seg
+ self.next = None
+ def set_next(self, next_seg):
+ self.next = next_seg
+ def add_klu(self, klu):
+ self.klu_list.append(klu)
+ klu.set_seg(self)
+ def add_unittf(self, unittf):
+ self.unittf_list.append(unittf)
+ unittf.set_next(self)
+ def add_histset(self, histset):
+ self.hist_set.append(histset)
+ histset.set_next(self)
\ No newline at end of file
diff --git a/ChanMACDUnitTF.py b/ChanMACDUnitTF.py
new file mode 100644
index 0000000..90a48a0
--- /dev/null
+++ b/ChanMACDUnitTF.py
@@ -0,0 +1,24 @@
+
+
+
+class ChanMACDUnitTF():
+ def __init__(self, start_time, start_klu, start_histset, pre_unittf, dir):
+ self.start_time = start_time
+ self.end_time = None
+ self.start_klu = start_klu
+ self.end_klu = None
+ self.klu_list = []
+ self.klu_list.append(start_klu)
+ self.histset_list = []
+ self.histset_list.append(start_histset)
+ self.next = None
+ self.pre = pre_unittf
+ self.uinttf_dir = dir
+ def set_next(self, next_unittf):
+ self.next = next_unittf
+ def add_histset(self, histset):
+ self.histset_list.append(histset)
+ histset.set_next(self)
+ def add_klu(self, klu):
+ self.klu_list.append(klu)
+ klu.set_unittf(self)
\ No newline at end of file
diff --git a/datasvc/Dockerfile b/datasvc/Dockerfile
new file mode 100644
index 0000000..44f15fb
--- /dev/null
+++ b/datasvc/Dockerfile
@@ -0,0 +1,17 @@
+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
+
+EXPOSE 9000
+
+CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "9000"]
+
+
diff --git a/datasvc/README.md b/datasvc/README.md
new file mode 100644
index 0000000..7a40f33
--- /dev/null
+++ b/datasvc/README.md
@@ -0,0 +1,48 @@
+# Local Data Service (REST + WebSocket)
+
+一键部署、跨平台的本地行情数据服务。默认抓取 Binance 永续合约 `BTC/USDT:USDT, ETH/USDT:USDT` 的 `1m/5m/15m/1h` K 线,增量写入本地 Parquet 并通过 WebSocket 推送。
+
+## 快速开始(方式B:已安装 Docker)
+
+```bash
+cd user_data/Chan/datasvc
+docker compose up -d
+```
+
+- REST: http://localhost:9000/api/candles?symbol=BTC/USDT:USDT&tf=1m
+- WS: ws://localhost:9000/ws?symbol=BTC/USDT:USDT&tf=1m&since=1690000000000
+- Swagger: http://localhost:9000/docs
+
+## 环境变量(docker-compose.yml)
+- EXCHANGE: 交易所,默认 binance
+- SYMBOLS: 逗号分隔交易对
+- TIMEFRAMES: 逗号分隔周期
+- START_DAYS: 首次启动回补最近 N 天
+- POLL_FACTOR: 轮询因子,间隔=周期毫秒*factor
+- DATA_DIR: 容器内数据目录(已映射到 `./data`)
+
+## 数据位置
+- 本地缓存:`user_data/Chan/datasvc/data/{timeframe}/{symbol}.parquet`
+
+## 常用命令
+```bash
+docker compose logs -f
+
+docker compose down
+```
+
+## 接口说明
+- GET /api/candles
+ - 参数:symbol, tf, start(ms), end(ms)
+ - 返回:[{timestamp, open, high, low, close, volume}]
+- WS /ws
+ - 参数:symbol, tf, since(ms)
+ - 消息:
+ - snapshot: 初始快照数组
+ - upsert: 单根K线增量(尾部修正)
+
+## 注意
+- 默认未带交易所 API Key,仅公共行情。
+- 如需更多交易对/周期,修改 `docker-compose.yml` 后重启。
+
+
diff --git a/datasvc/app/main.py b/datasvc/app/main.py
new file mode 100644
index 0000000..85a5c78
--- /dev/null
+++ b/datasvc/app/main.py
@@ -0,0 +1,207 @@
+import os
+import asyncio
+import json
+from datetime import datetime, timedelta
+from typing import Dict, List, Optional
+
+import ccxt
+import pandas as pd
+from fastapi import FastAPI, WebSocket, WebSocketDisconnect, Query
+from fastapi.responses import JSONResponse
+from fastapi.middleware.cors import CORSMiddleware
+
+from .storage import (
+ ensure_storage,
+ read_candles,
+ upsert_candles,
+ get_last_timestamp,
+)
+
+
+DATA_DIR = os.environ.get("DATA_DIR", "/data")
+EXCHANGE = os.environ.get("EXCHANGE", "binance")
+SYMBOLS = [s.strip() for s in os.environ.get("SYMBOLS", "BTC/USDT:USDT,ETH/USDT:USDT").split(",") if s.strip()]
+TIMEFRAMES = [t.strip() for t in os.environ.get("TIMEFRAMES", "1m,5m,15m,1h").split(",") if t.strip()]
+START_FROM = os.environ.get("START_FROM", "2025-01-01") # 首次启动拉取起始日期(UTC)
+POLL_FACTOR = float(os.environ.get("POLL_FACTOR", "0.5")) # 轮询间隔 = tf_ms * factor
+
+ensure_storage(DATA_DIR)
+
+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:
+ table = {
+ "1m": 60_000,
+ "3m": 3 * 60_000,
+ "5m": 5 * 60_000,
+ "15m": 15 * 60_000,
+ "30m": 30 * 60_000,
+ "1h": 60 * 60_000,
+ "2h": 2 * 60 * 60_000,
+ "4h": 4 * 60 * 60_000,
+ "1d": 24 * 60 * 60_000,
+ }
+ return table.get(tf, 60_000)
+
+
+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) # 允许 '2025-01-01' 或 '2025-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()
+
+
+def build_exchange():
+ if EXCHANGE.lower() == "binance":
+ return ccxt.binance({"enableRateLimit": True})
+ raise RuntimeError(f"Unsupported EXCHANGE: {EXCHANGE}")
+
+
+async def fetch_loop(symbol: str, timeframe: str):
+ """持续增量抓取并广播。"""
+ exchange = build_exchange()
+ tf_ms = tf_to_ms(timeframe)
+ start_since = parse_start_from_ms(START_FROM)
+ last_ts = get_last_timestamp(DATA_DIR, symbol, timeframe)
+ since = max(start_since, (last_ts + tf_ms) if last_ts else start_since)
+
+ while True:
+ try:
+ candles = exchange.fetch_ohlcv(symbol, timeframe, since=since, limit=1000)
+ if candles:
+ upsert_candles(DATA_DIR, symbol, timeframe, candles)
+ for row in candles[-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],
+ },
+ }
+ await hub.publish(symbol, timeframe, payload)
+ since = candles[-1][0] + tf_ms
+ await asyncio.sleep(max(1.0, tf_ms * POLL_FACTOR / 1000.0))
+ except Exception:
+ await asyncio.sleep(3.0)
+
+
+@app.on_event("startup")
+async def on_start():
+ ensure_storage(DATA_DIR)
+ for s in SYMBOLS:
+ for tf in TIMEFRAMES:
+ asyncio.create_task(fetch_loop(s, tf))
+
+
+@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:
+ df = read_candles(DATA_DIR, 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):
+ await hub.subscribe(websocket, symbol, tf)
+ try:
+ snap = read_candles(DATA_DIR, 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,
+ "timeframes": TIMEFRAMES,
+ }
+
+
diff --git a/datasvc/app/storage.py b/datasvc/app/storage.py
new file mode 100644
index 0000000..a73e45f
--- /dev/null
+++ b/datasvc/app/storage.py
@@ -0,0 +1,57 @@
+import os
+import threading
+from typing import List, Optional
+
+import pandas as pd
+
+
+_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
+ df = pd.read_parquet(p)
+ 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):
+ old = pd.read_parquet(p)
+ 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)
+
+
+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)
+ if df.empty:
+ return None
+ return int(df["timestamp"].iloc[-1])
+
+
diff --git a/datasvc/docker-compose.yml b/datasvc/docker-compose.yml
new file mode 100644
index 0000000..bb586db
--- /dev/null
+++ b/datasvc/docker-compose.yml
@@ -0,0 +1,19 @@
+version: "3.9"
+services:
+ datasvc:
+ build: .
+ container_name: datasvc
+ restart: unless-stopped
+ environment:
+ - EXCHANGE=binance
+ - SYMBOLS=BTC/USDT:USDT,ETH/USDT:USDT
+ - TIMEFRAMES=1m,5m,15m,1h
+ - START_FROM=2025-01-01
+ - POLL_FACTOR=0.5
+ - DATA_DIR=/data
+ - TZ=Asia/Shanghai
+ ports:
+ - "9000:9000"
+ volumes:
+ - ./data:/data
+
diff --git a/datasvc/requirements.txt b/datasvc/requirements.txt
new file mode 100644
index 0000000..da6ce3b
--- /dev/null
+++ b/datasvc/requirements.txt
@@ -0,0 +1,7 @@
+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
+
diff --git a/strategies/ChanLun_BTC_30.py b/strategies/ChanLun_BTC_30.py
index 7671bf1..ff3612d 100644
--- a/strategies/ChanLun_BTC_30.py
+++ b/strategies/ChanLun_BTC_30.py
@@ -67,15 +67,15 @@ class ChanLun_BTC_30(IStrategy):
can_short = True
lev = 1.0
stoploss = -0.3 # 设置为很大的负值,让custom_stoploss来控制
- use_custom_stoploss = False # 启用自定义止损
+ use_custom_stoploss = True # 启用自定义止损
trailing_stop = False
trailing_stop_positive = 0.025
trailing_stop_positive_offset = 0.045
trailing_only_offset_is_reached = False
- # 启用仓位调整功能以支持分批止盈
- position_adjustment_enable = True
+ # 关闭分批止盈/仓位调整
+ position_adjustment_enable = False
startup_candle_count = 2880
time5 = 5
@@ -119,9 +119,14 @@ class ChanLun_BTC_30(IStrategy):
#self.chan.plot_dual(dataframe_5, dataframe_30)
#chanpy_state = self.chanpy.get_bsp_state(dataframe_5)
#dataframe_5['chanpy_state'] = chanpy_state
+ state_list = self.chan.get_klc_state_list(dataframe_5)
+ dataframe_5['state'] = state_list
+ state_list = self.chan.get_klc_state_list(dataframe_15)
+ dataframe_15['state'] = state_list
+ state_list = self.chan.get_klc_state_list(dataframe_30)
+ dataframe_30['state'] = state_list
state_list = self.chan.get_klc_state_list(dataframe_60)
dataframe_60['state'] = state_list
- dataframe_60['fx'] = state_list
#bi_list_1 = self.chan.get_bi_list(dataframe)
#bi_list_5 = self.chan.get_bi_list(dataframe_5)
#bi_list_15 = self.chan.get_bi_list(dataframe_15)
@@ -138,7 +143,7 @@ class ChanLun_BTC_30(IStrategy):
self.last_time = datetime.now()
dataframe = resampled_merge(dataframe, dataframe_3)
dataframe = resampled_merge(dataframe, dataframe_5)
- #dataframe = resampled_merge(dataframe, dataframe_15)
+ dataframe = resampled_merge(dataframe, dataframe_15)
#dataframe = resampled_merge(dataframe, dataframe_30)
dataframe = resampled_merge(dataframe, dataframe_60)
#dataframe = resampled_merge(dataframe, dataframe_4h)
@@ -166,7 +171,8 @@ class ChanLun_BTC_30(IStrategy):
bb120 = ta.BBANDS(df, timeperiod=120, nbdevup=3.0, nbdevdn=3.0, matype=0)
bb30 = ta.BBANDS(df, timeperiod=41, nbdevup=2.3, nbdevdn=2.3, matype=0)
bb302 = ta.BBANDS(df, timeperiod=41, nbdevup=2.0, nbdevdn=2.0, matype=0)
-
+ bb30 = ta.BBANDS(df, timeperiod=20, nbdevup=2.0, nbdevdn=2.0, matype=0)
+ bb302 = ta.BBANDS(df, timeperiod=20, nbdevup=2.0, nbdevdn=2.0, matype=0)
# 计算布林带中轨(移动平均线)
bb30_middle = ta.SMA(df, timeperiod=90)
@@ -230,162 +236,63 @@ class ChanLun_BTC_30(IStrategy):
new_exitprice = proposed_rate - 50
return new_exitprice
- def adjust_trade_position1(self, trade: Trade, current_time: datetime,
+ def adjust_trade_position(self, trade: Trade, current_time: datetime,
current_rate: float, current_profit: float,
min_stake: Optional[float], max_stake: float,
current_entry_rate: float, current_exit_rate: float,
current_entry_profit: float, current_exit_profit: float,
**kwargs) -> Optional[float]:
- """
- 基于布林带的分批止盈逻辑
- """
- # 获取当前数据
- dataframe, _ = self.dp.get_analyzed_dataframe(trade.pair, self.timeframe)
- if dataframe is None or len(dataframe) == 0:
- return None
-
- last_candle = dataframe.iloc[-1]
-
- # 获取布林带数据
- bb30_middle = last_candle['bbmiddle30']
- bb30_upper = last_candle['bbup30']
- bb30_lower = last_candle['bblow30']
- bb302_upper = last_candle['bbup302']
- bb302_lower = last_candle['bblow302']
-
- # 获取交易的状态标记
- first_tp_triggered = trade.get_custom_data(key="first_tp_triggered", default=False)
- second_tp_triggered = trade.get_custom_data(key="second_tp_triggered", default=False)
-
- if trade.is_short:
- # 做空逻辑
- if not first_tp_triggered and current_rate <= bb30_middle:
- # 第一次止盈:价格跌到bb30中轨,止盈50%
- logger.info(f"做空第一次止盈触发:价格{current_rate} <= BB30中轨{bb30_middle}")
- trade.set_custom_data(key="first_tp_triggered", value=True)
- trade.set_custom_data(key="new_stoploss", value=trade.open_rate) # 设置止损为开仓价
- return -(trade.amount * 0.5) # 减少50%仓位
-
- elif first_tp_triggered and not second_tp_triggered and current_rate <= bb302_lower:
- # 第二次止盈:继续跌到bb302下轨,止盈剩余仓位的60%
- logger.info(f"做空第二次止盈触发:价格{current_rate} <= BB302下轨{bb302_lower}")
- trade.set_custom_data(key="second_tp_triggered", value=True)
- trade.set_custom_data(key="new_stoploss", value=bb30_middle) # 移动止损到bb30中轨
- remaining_amount = trade.amount * 0.5 # 剩余50%
- return -(remaining_amount * 0.6) # 减少剩余仓位的60%
-
- else:
- # 做多逻辑
- if not first_tp_triggered and current_rate >= bb30_middle:
- # 第一次止盈:价格涨到bb30中轨,止盈50%
- logger.info(f"做多第一次止盈触发:价格{current_rate} >= BB30中轨{bb30_middle}")
- trade.set_custom_data(key="first_tp_triggered", value=True)
- trade.set_custom_data(key="new_stoploss", value=trade.open_rate) # 设置止损为开仓价
- return -(trade.amount * 0.5) # 减少50%仓位
-
- elif first_tp_triggered and not second_tp_triggered and current_rate >= bb302_upper:
- # 第二次止盈:继续涨到bb302上轨,止盈剩余仓位的60%
- logger.info(f"做多第二次止盈触发:价格{current_rate} >= BB302上轨{bb302_upper}")
- trade.set_custom_data(key="second_tp_triggered", value=True)
- trade.set_custom_data(key="new_stoploss", value=bb30_middle) # 移动止损到bb30中轨
- remaining_amount = trade.amount * 0.5 # 剩余50%
- return -(remaining_amount * 0.6) # 减少剩余仓位的60%
-
+ # 关闭分批止盈,始终不调整仓位
return None
def custom_stoploss(self, pair: str, trade: Trade, current_time: datetime,
current_rate: float, current_profit: float, after_fill: bool,
**kwargs) -> float | None:
"""
- 动态止损逻辑
+ 止损 = 开仓价 ± 1 * ATR(开仓时的ATR)。
+ 多单: 开仓价 - ATR;空单: 开仓价 + ATR。
"""
- # 检查是否有自定义的新止损价格(分批止盈后的动态止损)
- new_stoploss_price = trade.get_custom_data(key="new_stoploss")
- if new_stoploss_price:
- logger.info(f"使用动态止损价格: {new_stoploss_price}")
- return stoploss_from_absolute(new_stoploss_price, current_rate, is_short=trade.is_short)
-
-
- # 如果没有ATR数据,使用固定的5%止损作为备用
- logger.warning(f"未找到开仓时ATR数据,使用默认5%止损")
- return -0.05
+ entry_atr = trade.get_custom_data(key="entry_atr")
+ if entry_atr is None:
+ # 回退:取当前数据的 ATR 估算
+ dataframe, _ = self.dp.get_analyzed_dataframe(trade.pair, self.timeframe)
+ if dataframe is not None and len(dataframe) > 0 and 'atr' in dataframe.columns:
+ entry_atr = float(dataframe.iloc[-1]['atr'])
+ else:
+ # 最保守的回退:5%
+ return -0.05
+
+ if trade.is_short:
+ stop_price = trade.open_rate + float(entry_atr)
+ else:
+ stop_price = trade.open_rate - float(entry_atr)
+ return stoploss_from_absolute(stop_price, current_rate, is_short=trade.is_short)
def custom_exit(self, pair: str, trade: Trade, current_time: datetime, current_rate: float,
current_profit: float, **kwargs):
- """
- 自定义退出逻辑 - 处理最终止盈条件
- """
- # 获取当前数据
- dataframe, _ = self.dp.get_analyzed_dataframe(pair, self.timeframe)
- if dataframe is None or len(dataframe) == 0:
- return None
-
- last_candle = dataframe.iloc[-1]
-
- # 获取布林带数据
- bb30_upper = last_candle['bbup30']
- bb30_lower = last_candle['bblow30']
-
- # 检查是否已经触发过前两次止盈
- first_tp_triggered = trade.get_custom_data(key="first_tp_triggered", default=False)
- second_tp_triggered = trade.get_custom_data(key="second_tp_triggered", default=False)
-
- if trade.is_short:
- # 做空:如果价格跌到bb30下轨,全部止盈
- if first_tp_triggered and second_tp_triggered and current_rate <= bb30_lower:
- logger.info(f"做空最终止盈触发:价格{current_rate} <= BB30下轨{bb30_lower}")
- return "short_final_tp_bb30_lower"
- else:
- # 做多:如果价格涨到bb30上轨,全部止盈
- if first_tp_triggered and second_tp_triggered and current_rate >= bb30_upper:
- logger.info(f"做多最终止盈触发:价格{current_rate} >= BB30上轨{bb30_upper}")
- return "long_final_tp_bb30_upper"
-
- # 原有退出逻辑
- if trade.is_short:
- last_high = trade.get_custom_data(key="entry_candle_high")
- if last_high and current_rate > last_high:
- return "Relay Top FX exit"
- else:
- last_low = trade.get_custom_data(key="entry_candle_low")
- if last_low and current_rate < last_low:
- return "Relay Bottom FX exit"
-
+ # 不做分批止盈/最终止盈处理,退出由策略信号/ROI/止损决定
return None
-
- def confirm_trade_entry1(self, pair: str, order_type: str, amount: float, rate: float,
+
+ def confirm_trade_entry(self, pair: str, order_type: str, amount: float, rate: float,
time_in_force: str, current_time: datetime, entry_tag: str | None,
side: str, **kwargs) -> bool:
- if self.last_trade:
- if self.last_trade.is_short:
- if side == 'short':
- if self.last_trade.open_date + timedelta(minutes=30) > current_time:
- return False
- else:
- return True
- else:
- if side == 'long':
- if self.last_trade.open_date + timedelta(minutes=30) > current_time:
- return True
- else:
- return False
- #if self.last_trade:
- #print(self.last_trade.open_date, current_time, self.last_trade.open_date + timedelta(minutes=self.time5))
- return True
- def custom_stoploss1(self, pair: str, trade: Trade, current_time: datetime,
- current_rate: float, current_profit: float, after_fill: bool,
- **kwargs) -> float | None:
-
- last_high = trade.get_custom_data(key="entry_candle_high")
- last_low = trade.get_custom_data(key="entry_candle_low")
-
- # Convert absolute price to percentage relative to current_rate
- if last_high:
- return stoploss_from_absolute(last_high, current_rate, is_short=trade.is_short)
- if last_low:
- return stoploss_from_absolute(last_low, current_rate, is_short=trade.is_short)
- # return maximum stoploss value, keeping current stoploss price unchanged
- return None
+ """
+ ATR 过滤:atr < 100 不开单。
+ """
+ try:
+ dataframe, _ = self.dp.get_analyzed_dataframe(pair, self.timeframe)
+ if dataframe is None or len(dataframe) == 0:
+ return False
+ last = dataframe.iloc[-1]
+ atr_str = 'resample_{}_atr'.format(self.get_ticker_indicator()*self.time30)
+ atr_val = float(last.get(atr_str, 0) or 0)
+ if atr_val < 100:
+ logger.info(f"ATR过滤:atr={atr_val:.2f} < 100, 拒绝进场 {pair}")
+ return False
+ return True
+ except Exception as e:
+ logger.warning(f"confirm_trade_entry 异常: {e}")
+ return True
def order_filled(self, pair: str, trade: Trade, order: Order, current_time: datetime, **kwargs) -> None:
"""
@@ -400,10 +307,10 @@ class ChanLun_BTC_30(IStrategy):
# Obtain pair dataframe (just to show how to access it)
dataframe, _ = self.dp.get_analyzed_dataframe(trade.pair, self.timeframe)
last_candle = dataframe.iloc[-1].squeeze()
-
+ atr_str = 'resample_{}_atr'.format(self.get_ticker_indicator()*self.time30)
# 保存开仓时的ATR值用于止损计算
if (trade.nr_of_successful_entries == 1) and (order.ft_order_side == trade.entry_side):
- entry_atr = last_candle['atr']
+ entry_atr = last_candle[atr_str] * 4
trade.set_custom_data(key="entry_atr", value=entry_atr)
logger.info(f"保存开仓时ATR值: {entry_atr}")
return None
@@ -427,7 +334,7 @@ class ChanLun_BTC_30(IStrategy):
['enter_long', 'enter_tag']] = (1, 'long_signal_chan')
dataframe.loc[
(
- (dataframe[state_str].shift(shift_time) == "10")
+ (dataframe[state_str].shift(shift_time) == "10")
#(dataframe[fx_str].shift(shift_time) == 1)
#(dataframe[chanpy_state_str].shift(shift_time+30) == -1)
#(dataframe['resample_{}_state'.format(self.get_ticker_indicator()*self.time5)].shift(self.time5) == "-10") &
diff --git a/web/app.py b/web/app.py
index 9261a3d..34233f3 100644
--- a/web/app.py
+++ b/web/app.py
@@ -267,6 +267,18 @@ def add_indicators(df):
df['ma10'] = (ta.MA(df, timeperiod=10)).fillna(0)
df['ma30'] = (ta.EMA(df, timeperiod=30)).fillna(0)
df['ma250'] = (ta.MA(df, timeperiod=250)).fillna(0)
+ # 新增 EMA 指标
+ df['ema5'] = (ta.EMA(df, timeperiod=5)).fillna(0)
+ df['ema10'] = (ta.EMA(df, timeperiod=10)).fillna(0)
+ df['ema24'] = (ta.EMA(df, timeperiod=24)).fillna(0)
+ df['ema52'] = (ta.EMA(df, timeperiod=52)).fillna(0)
+ # 常用SMA 24/52
+ try:
+ df['sma24'] = (ta.SMA(df, timeperiod=24)).fillna(0)
+ df['sma52'] = (ta.SMA(df, timeperiod=52)).fillna(0)
+ except Exception:
+ df['sma24'] = 0
+ df['sma52'] = 0
df['rsi'] = ta.RSI(df, timeperiod=14)
# 计算布林带 (当前周期 - 20周期,2标准差)
@@ -275,9 +287,11 @@ def add_indicators(df):
df['bb_middle'] = bb['middleband'].fillna(0)
df['bb_lower'] = bb['lowerband'].fillna(0)
bb30 = ta.BBANDS(df, timeperiod=41, nbdevup=2.3, nbdevdn=2.3, matype=0)
+ bb30 = ta.BBANDS(df, timeperiod=20, nbdevup=2.0, nbdevdn=2.0, matype=0)
df['bbup30'] = bb30['upperband'].fillna(0)
df['bblow30'] = bb30['lowerband'].fillna(0)
bb302 = ta.BBANDS(df, timeperiod=41, nbdevup=2.0, nbdevdn=2.0, matype=0)
+ bb302 = ta.BBANDS(df, timeperiod=20, nbdevup=2.0, nbdevdn=2.0, matype=0)
df['bbup302'] = bb302['upperband'].fillna(0)
df['bblow302'] = bb302['lowerband'].fillna(0)
# 计算次周期布林带 (14周期,2标准差)
@@ -293,6 +307,12 @@ def add_indicators(df):
df['ma10'] = df['ma10'].fillna(0)
df['ma30'] = df['ma30'].fillna(0)
df['ma250'] = df['ma250'].fillna(0)
+ df['ema5'] = df['ema5'].fillna(0)
+ df['ema10'] = df['ema10'].fillna(0)
+ df['ema24'] = df['ema24'].fillna(0)
+ df['ema52'] = df['ema52'].fillna(0)
+ df['sma24'] = df['sma24'].fillna(0)
+ df['sma52'] = df['sma52'].fillna(0)
df['rsi'] = df['rsi'].fillna(0)
df['avg_volume'] = df['volume'].rolling(10).mean()
# 计算量比,避免产生Infinity值
@@ -312,10 +332,10 @@ def add_indicators(df):
def calculate_macd(df):
"""计算MACD指标"""
- exp1 = df['close'].ewm(span=24, adjust=False).mean()
- exp2 = df['close'].ewm(span=52, adjust=False).mean()
+ exp1 = df['close'].ewm(span=12, adjust=False).mean()
+ exp2 = df['close'].ewm(span=26, adjust=False).mean()
macd = exp1 - exp2
- signal = macd.ewm(span=18, adjust=False).mean()
+ signal = macd.ewm(span=9, adjust=False).mean()
histogram = macd - signal
return {
@@ -1035,6 +1055,228 @@ def clean_dataframe_for_json(df):
return clean_df
+# ====== 趋势判定与趋势筛选(币对) ======
+
+def classify_trend_stage(df):
+ """根据 EMA 斜率与多空排列判断趋势方向与阶段
+ 返回: direction in {"bull","bear","sideways"}, stage in {"early","mid","late"}, strength_score (0-100)
+ """
+ if df is None or len(df) < 60:
+ return "sideways", "early", 0
+
+ # 使用 EMA5/10/24/52
+ closes = df['close'].values
+ ema5 = df['ema5'].values if 'ema5' in df else ta.EMA(df, timeperiod=5)
+ ema10 = df['ema10'].values if 'ema10' in df else ta.EMA(df, timeperiod=10)
+ ema24 = df['ema24'].values if 'ema24' in df else ta.EMA(df, timeperiod=24)
+ ema52 = df['ema52'].values if 'ema52' in df else ta.EMA(df, timeperiod=52)
+
+ # 最近N根用于斜率与排列判定
+ lookback = min(30, len(df) - 1)
+ if lookback <= 5:
+ return "sideways", "early", 0
+
+ # 简单斜率: 最近k根的线性变化率近似
+ def slope(arr, k=10):
+ k = min(k, len(arr) - 1)
+ if k < 2:
+ return 0.0
+ y = arr[-k:]
+ x = np.arange(k)
+ # 最小二乘拟合斜率
+ denom = np.dot(x - x.mean(), x - x.mean())
+ if denom == 0:
+ return 0.0
+ m = np.dot(y - y.mean(), x - x.mean()) / denom
+ return float(m)
+
+ k_slope = 12 # 斜率窗口
+ s5 = slope(ema5, k_slope)
+ s10 = slope(ema10, k_slope)
+ s24 = slope(ema24, k_slope)
+ s52 = slope(ema52, k_slope)
+
+ # 多空排列
+ last5, last10, last24, last52 = ema5[-1], ema10[-1], ema24[-1], ema52[-1]
+ bull_stack = last5 > last10 > last24 > last52
+ bear_stack = last5 < last10 < last24 < last52
+
+ # 波动性与动量增强: MACD 柱体最近均值
+ macdhist = df['macdhist'].values if 'macdhist' in df else calculate_macd(df)['histogram']
+ hist_recent = macdhist[-lookback:]
+ hist_power = float(np.mean(np.abs(hist_recent))) if len(hist_recent) else 0.0
+
+ # 方向
+ if bull_stack and s24 > 0 and s52 > 0:
+ direction = "bull"
+ elif bear_stack and s24 < 0 and s52 < 0:
+ direction = "bear"
+ else:
+ # 用价格相对 EMA52 辅助
+ if closes[-1] > last52 and (s24 + s52) > 0:
+ direction = "bull"
+ elif closes[-1] < last52 and (s24 + s52) < 0:
+ direction = "bear"
+ else:
+ direction = "sideways"
+
+ # 阶段: 依据(斜率大小、与EMA52距离、MACD柱体扩张/收敛)
+ dist52 = float((closes[-1] - last52) / last52) if last52 else 0.0
+ slope_score = max(0.0, (abs(s24) + abs(s52)) * 1000.0) # 归一化
+ dist_score = min(50.0, abs(dist52) * 200.0)
+ hist_score = min(30.0, hist_power * 10.0)
+ strength = float(min(100.0, slope_score + dist_score + hist_score))
+
+ # 简单阶段判定
+ if direction == "sideways":
+ stage = "early"
+ strength = min(strength, 30.0)
+ else:
+ # 查看最近 hist 是否在扩大或收敛
+ if len(hist_recent) >= 6:
+ recent_growth = np.mean(np.abs(hist_recent[-3:])) - np.mean(np.abs(hist_recent[-6:-3]))
+ else:
+ recent_growth = 0.0
+
+ if recent_growth > 0 and abs(dist52) < 0.05:
+ stage = "early"
+ elif recent_growth > 0 and abs(dist52) >= 0.05:
+ stage = "mid"
+ else:
+ stage = "late"
+
+ return direction, stage, strength
+
+
+def load_crypto_symbols(limit=200):
+ """加载常见USDT永续合约交易对,返回列表"""
+ try:
+ markets = exchange.load_markets()
+ symbols = [s for s in markets.keys() if '/USDT' in s and ':USDT' in s]
+ return symbols[:limit]
+ except Exception:
+ return SYMBOLS
+
+
+@app.route('/api/trend_filter', methods=['GET'])
+def trend_filter():
+ """趋势筛选接口(币对)
+ 参数:
+ timeframe: K线周期
+ start_time, end_time: 毫秒时间戳,可选
+ direction: bull/bear/sideways 可选
+ stage: early/mid/late 可选
+ min_strength: 0-100 可选
+ symbols: 逗号分隔列表,可选;不传则自动加载部分USDT币对
+ 返回符合条件的币对与简要统计
+ """
+ timeframe = request.args.get('timeframe', '1h')
+ start_time = request.args.get('start_time')
+ end_time = request.args.get('end_time')
+ want_direction = request.args.get('direction') # 可为 None
+ want_stage = request.args.get('stage') # 可为 None
+ try:
+ min_strength = float(request.args.get('min_strength', '0'))
+ except ValueError:
+ min_strength = 0.0
+
+ symbols_param = request.args.get('symbols')
+ if symbols_param:
+ symbols_list = [s.strip() for s in symbols_param.split(',') if s.strip()]
+ else:
+ symbols_list = load_crypto_symbols(limit=150)
+
+ results = []
+ for sym in symbols_list:
+ try:
+ df = get_crypto_kl_data(sym, timeframe, start_time=start_time, end_time=end_time)
+ if df is None or len(df) < 60:
+ continue
+ df = add_indicators(df)
+ direction, stage, strength = classify_trend_stage(df)
+
+ if want_direction and direction != want_direction:
+ continue
+ if want_stage and stage != want_stage:
+ continue
+ if strength < min_strength:
+ continue
+
+ last_row = df.iloc[-1]
+ results.append({
+ 'symbol': sym,
+ 'time': int(last_row['timestamp']),
+ 'close': float(last_row['close']),
+ 'direction': direction,
+ 'stage': stage,
+ 'strength': float(round(strength, 2)),
+ 'ema5': float(last_row['ema5']),
+ 'ema10': float(last_row['ema10']),
+ 'ema24': float(last_row['ema24']),
+ 'ema52': float(last_row['ema52'])
+ })
+ except Exception:
+ continue
+
+ # 按强度降序
+ results.sort(key=lambda x: x['strength'], reverse=True)
+ return jsonify({
+ 'count': len(results),
+ 'results': results
+ })
+
+
+@app.route('/api/trend_detail', methods=['GET'])
+def trend_detail():
+ """返回单个币对的K线与EMA、用于前端绘制趋势线
+ 参数: symbol, timeframe, start_time, end_time
+ """
+ symbol = request.args.get('symbol')
+ timeframe = request.args.get('timeframe', '1h')
+ start_time = request.args.get('start_time')
+ end_time = request.args.get('end_time')
+ timezone_name = request.args.get('timezone', 'Asia/Shanghai')
+
+ if not symbol:
+ return jsonify({'error': 'symbol不能为空'})
+
+ df = get_crypto_kl_data(symbol, timeframe, start_time=start_time, end_time=end_time)
+ if df is None or len(df) == 0:
+ return jsonify({'error': '获取数据失败'})
+
+ df = add_indicators(df)
+ direction, stage, strength = classify_trend_stage(df)
+
+ # 简单趋势线: 用最近N根收盘价做线性拟合
+ N = min(80, len(df))
+ sub = df.tail(N)
+ y = sub['close'].values
+ x = np.arange(len(y))
+ denom = np.dot(x - x.mean(), x - x.mean())
+ if denom != 0:
+ m = float(np.dot(y - y.mean(), x - x.mean()) / denom)
+ b = float(y.mean() - m * x.mean())
+ else:
+ m, b = 0.0, float(y[-1])
+
+ client_tz = timezone(timezone_name)
+
+ return jsonify({
+ 'symbol': symbol,
+ 'timeframe': timeframe,
+ 'timezone': timezone_name,
+ 'direction': direction,
+ 'stage': stage,
+ 'strength': float(round(strength, 2)),
+ 'kline_data': clean_dataframe_for_json(df)[['timestamp','open','high','low','close','volume','ema5','ema10','ema24','ema52']].to_dict('records'),
+ 'trend_line': {
+ 'offset': int(df.index[-N]),
+ 'slope': m,
+ 'intercept': b,
+ 'length': int(N)
+ }
+ })
+
@app.route('/')
def index():
"""主页"""
diff --git a/web/templates/index.html b/web/templates/index.html
index 765797b..47ae231 100644
--- a/web/templates/index.html
+++ b/web/templates/index.html
@@ -1057,6 +1057,9 @@
+
+
+
@@ -1177,6 +1180,125 @@
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+ 正在筛选,请稍候...
+
+
+
+
+ | 交易对 |
+ 时间 |
+ 方向 |
+ 阶段 |
+ 强度 |
+ 收盘 |
+ EMA5 |
+ EMA10 |
+ EMA24 |
+ EMA52 |
+ 操作 |
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+ | 时间 |
+ 开 |
+ 高 |
+ 低 |
+ 收 |
+ 成交量 |
+ EMA5 |
+ EMA10 |
+ EMA24 |
+ EMA52 |
+
+
+
+
+
+
+
+
+
+
+