feat: 健康检查加币安 API/WS,订单簿 bids 空防护,自动价差历史取 20 根

- ApiHealthCheckService: 新增币安 API(ping)、币安 WebSocket(5m/15m 连接状态)
- BinanceKlineService: 连接状态追踪 getConnectionStatuses 供健康检查
- CryptoTailOrderbookWsService: book 事件 bids 为空时不再取 [0],避免 Index 0 out of bounds
- BinanceKlineAutoSpreadService + 文档 + i18n: 历史 K 线由 30 根改为 20 根

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
WrBug
2026-02-14 15:27:22 +08:00
parent b50e43c239
commit 9fef4bea59
8 changed files with 120 additions and 18 deletions
@@ -9,7 +9,7 @@ import java.math.RoundingMode
import java.util.concurrent.ConcurrentHashMap
/**
* 自动最小价差:按周期计算。每个周期首次需要时,拉取该周期前的 30 根已收盘 K 线,按方向筛选、IQR 剔除后求平均 × 0.8,缓存 (interval, period)。
* 自动最小价差:按周期计算。每个周期首次需要时,拉取该周期前的 20 根已收盘 K 线,按方向筛选、IQR 剔除后求平均 × 0.8,缓存 (interval, period)。
* 不在保存策略时计算。
*/
@Service
@@ -20,7 +20,7 @@ class BinanceKlineAutoSpreadService(
private val logger = LoggerFactory.getLogger(BinanceKlineAutoSpreadService::class.java)
private val symbol = "BTCUSDC"
private val historyLimit = 30
private val historyLimit = 20
private val autoSpreadCoefficient = BigDecimal("0.8")
private val minSamplesAfterIqr = 3
@@ -16,6 +16,7 @@ import org.springframework.stereotype.Service
import java.math.BigDecimal
import jakarta.annotation.PreDestroy
import java.util.concurrent.ConcurrentHashMap
import java.util.concurrent.atomic.AtomicBoolean
/**
* 币安 K 线 WebSocket:订阅 BTCUSDC 5m/15m,维护当前周期 (open, close),供尾盘策略价差校验使用。
@@ -34,6 +35,8 @@ class BinanceKlineService {
private var ws5m: WebSocket? = null
private var ws15m: WebSocket? = null
private var reconnectJob: Job? = null
private val connected5m = AtomicBoolean(false)
private val connected15m = AtomicBoolean(false)
init {
connectAll()
@@ -45,6 +48,12 @@ class BinanceKlineService {
return openCloseByPeriod[key(intervalSeconds, periodStartUnix)]
}
/** 供 API 健康检查使用:5m / 15m 连接是否正常 */
fun getConnectionStatuses(): Map<String, Boolean> = mapOf(
"5m" to connected5m.get(),
"15m" to connected15m.get()
)
private fun connectAll() {
if (ws5m != null && ws15m != null) return
connectStream("btcusdc@kline_5m") { intervalSec, tMs, openP, closeP ->
@@ -68,7 +77,16 @@ class BinanceKlineService {
else -> 300
}
val request = Request.Builder().url(url).build()
val connectedFlag = when {
streamName.contains("kline_5m") -> connected5m
streamName.contains("kline_15m") -> connected15m
else -> null
}
val ws = client.newWebSocket(request, object : WebSocketListener() {
override fun onOpen(webSocket: WebSocket, response: okhttp3.Response) {
connectedFlag?.set(true)
}
override fun onMessage(webSocket: WebSocket, text: String) {
parseKlineMessage(text, intervalSeconds)?.let { (tMs, o, c) ->
onKline(intervalSeconds, tMs, o, c)
@@ -76,13 +94,19 @@ class BinanceKlineService {
}
override fun onFailure(webSocket: WebSocket, t: Throwable, response: okhttp3.Response?) {
connectedFlag?.set(false)
logger.warn("币安 K 线 WS 异常 $streamName: ${t.message}")
scheduleReconnect()
}
override fun onClosing(webSocket: WebSocket, code: Int, reason: String) {
connectedFlag?.set(false)
if (code != 1000) scheduleReconnect()
}
override fun onClosed(webSocket: WebSocket, code: Int, reason: String) {
connectedFlag?.set(false)
}
})
logger.info("币安 K 线 WS 已连接: $streamName")
return ws
@@ -112,6 +136,8 @@ class BinanceKlineService {
ws15m?.close(1000, "reconnect")
ws5m = null
ws15m = null
connected5m.set(false)
connected15m.set(false)
logger.info("币安 K 线 WS 尝试重连")
connectAll()
}
@@ -127,7 +127,8 @@ class CryptoTailOrderbookWsService(
"book" -> {
val assetId = (json.get("asset_id") as? com.google.gson.JsonPrimitive)?.asString ?: return
val bids = json.get("bids") as? com.google.gson.JsonArray
val firstBid = bids?.get(0) as? com.google.gson.JsonObject
if (bids == null || bids.isEmpty) return
val firstBid = bids.get(0) as? com.google.gson.JsonObject
val bestBid = (firstBid?.get("price") as? com.google.gson.JsonPrimitive)?.asString?.toSafeBigDecimal()
if (bestBid != null) onBestBid(assetId, bestBid)
}
@@ -15,6 +15,7 @@ import org.springframework.context.ApplicationContextAware
import com.wrbug.polymarketbot.service.copytrading.orders.OrderPushService
import com.wrbug.polymarketbot.service.copytrading.monitor.PolymarketActivityWsService
import com.wrbug.polymarketbot.service.copytrading.monitor.UnifiedOnChainWsService
import com.wrbug.polymarketbot.service.binance.BinanceKlineService
import org.springframework.stereotype.Service
import java.util.concurrent.TimeUnit
@@ -76,6 +77,17 @@ class ApiHealthCheckService(
}
}
/**
* 获取 BinanceKlineService(通过 ApplicationContext 避免循环依赖)
*/
private fun getBinanceKlineService(): BinanceKlineService? {
return try {
applicationContext?.getBean(BinanceKlineService::class.java)
} catch (e: BeansException) {
null
}
}
private val logger = LoggerFactory.getLogger(ApiHealthCheckService::class.java)
/**
@@ -91,6 +103,8 @@ class ApiHealthCheckService(
async { checkDataApi() },
async { checkGammaApi() },
async { checkPolygonRpc() },
async { checkBinanceApi() },
async { checkBinanceWebSocket() },
async { checkPolymarketRtdsWebSocket() },
async { checkPolymarketActivityWebSocket() },
async { checkUnifiedOnChainWebSocket() },
@@ -197,6 +211,67 @@ class ApiHealthCheckService(
checkJsonRpcApi("Polygon RPC", rpcUrl)
}
/**
* 检查币安 API(用于 K 线等)
* 使用 /api/v3/ping 端点
*/
private suspend fun checkBinanceApi(): ApiHealthCheckDto = withContext(Dispatchers.IO) {
val url = "https://api.binance.com/api/v3/ping"
checkApi("币安 API", url)
}
/**
* 检查币安 K 线 WebSocket 连接状态(5m / 15m
*/
private suspend fun checkBinanceWebSocket(): ApiHealthCheckDto = withContext(Dispatchers.Default) {
val binanceWsUrl = "wss://stream.binance.com:9443"
try {
val binanceKlineService = getBinanceKlineService()
if (binanceKlineService == null) {
return@withContext ApiHealthCheckDto(
name = "币安 WebSocket",
url = binanceWsUrl,
status = "error",
message = "服务未初始化"
)
}
val statuses = binanceKlineService.getConnectionStatuses()
val total = statuses.size
val connected = statuses.values.count { it }
if (connected == total && total > 0) {
ApiHealthCheckDto(
name = "币安 WebSocket",
url = binanceWsUrl,
status = "success",
message = "连接正常 (5m、15m)"
)
} else if (connected > 0) {
val which = statuses.filter { it.value }.keys.joinToString("")
ApiHealthCheckDto(
name = "币安 WebSocket",
url = binanceWsUrl,
status = "error",
message = "部分连接正常 ($which)"
)
} else {
ApiHealthCheckDto(
name = "币安 WebSocket",
url = binanceWsUrl,
status = "error",
message = "连接断开"
)
}
} catch (e: Exception) {
logger.warn("检查币安 WebSocket 状态失败", e)
ApiHealthCheckDto(
name = "币安 WebSocket",
url = binanceWsUrl,
status = "error",
message = "检查失败:${e.message}"
)
}
}
/**
* 检查 Polymarket RTDS WebSocket 连接状态
* 用于订单推送服务
+12 -12
View File
@@ -19,11 +19,11 @@
### 自动模式计算逻辑
- 通过币安 API 获取**历史 30 根** K 线(与策略周期一致:5m 取 5m K 线,15m 取 15m K 线)。
- 通过币安 API 获取**历史 20 根** K 线(与策略周期一致:5m 取 5m K 线,15m 取 15m K 线)。
- **下单方向 = Down**outcomeIndex = 1):只取「收盘价 < 开盘价」的 K 线,得到价差序列(开盘价 − 收盘价)。
- **下单方向 = Up**outcomeIndex = 0):只取「收盘价 > 开盘价」的 K 线,得到价差序列(收盘价 − 开盘价)。
- **异常值剔除**:对上述价差序列做异常值过滤(见下文「异常值剔除」),再用**剩余样本**求平均价差,乘以系数 **80%** 得到最小价差;后续用该值做 \|收盘价 − 开盘价\| ≥ 该值 的校验。
- **历史数据获取时机**:**在该周期开始时就拉取并计算**,不在保存策略时计算。订单簿 WS 在周期开始时刷新订阅(含每 25 秒或周期切换时的 refreshAndSubscribe),此时对当前周期内所有启用且为 AUTO 的策略,按 (intervalSeconds, periodStartUnix) 预拉该周期前 30 根已收盘 K 线并计算 minSpreadUp/minSpreadDown 写入缓存;该周期内触发时直接用缓存,无需在触发时再调 REST。
- **历史数据获取时机**:**在该周期开始时就拉取并计算**,不在保存策略时计算。订单簿 WS 在周期开始时刷新订阅(含每 25 秒或周期切换时的 refreshAndSubscribe),此时对当前周期内所有启用且为 AUTO 的策略,按 (intervalSeconds, periodStartUnix) 预拉该周期前 20 根已收盘 K 线并计算 minSpreadUp/minSpreadDown 写入缓存;该周期内触发时直接用缓存,无需在触发时再调 REST。
### 异常值剔除
@@ -35,7 +35,7 @@
- 示例:15 组价差,14 组在 50 以内、1 组为 200 → 200 会超出上界被剔除,只用 14 组参与平均。
- **边界与降级**
- 若剔除后剩余样本数过少(如 &lt; 3),则**不剔除**:用全部价差样本求平均 × 0.8。
- 若无满足方向的 K 线(如 30 根里没有 close &lt; open),仍按原文档降级处理(全量 \|close−open\| 或返回 0)。
- 若无满足方向的 K 线(如 20 根里没有 close &lt; open),仍按原文档降级处理(全量 \|close−open\| 或返回 0)。
---
@@ -73,7 +73,7 @@
│ • 计算 effectiveMinSpread
│ - FIXEDeffectiveMinSpread = 策略.minSpreadValue(用户填的固定值) │
│ - AUTOeffectiveMinSpread = 按当前下单方向(outcomeIndex)取「自动计算 │
│ 的最小价差」(见下节;若尚未计算则先拉 30 根历史 K 线并计算、缓存)。 │
│ 的最小价差」(见下节;若尚未计算则先拉 20 根历史 K 线并计算、缓存)。 │
│ • 若 |close open| < effectiveMinSpread → 本轮不下单,等待价差满足。 │
│ • 若 |close open| >= effectiveMinSpread → 通过价差校验,进入 │
│ placeOrderForTrigger(与现有逻辑一致:预签/签名、提交 CLOB 订单、写触发记录)。│
@@ -87,18 +87,18 @@
## 四、自动模式:何时拉历史、如何算、如何用
- **何时拉 30 根历史 K 线并计算**
- **何时拉 20 根历史 K 线并计算**
- **在该周期开始时就预计算**,不在保存策略时计算。
- 订单簿 WS 在**周期开始时**会刷新订阅(`refreshAndSubscribe`:每 25 秒或检测到周期切换时),此时对当前周期内所有启用且 minSpreadMode=AUTO 的策略,按 `(intervalSeconds, periodStartUnix)` 异步拉取该周期前 30 根已收盘 K 线(REST `endTime = periodStartUnix * 1000`),按 Up/Down 分别算 avgSpread × 0.8(含 IQR 剔除)并写入缓存。该周期内后续触发时直接用缓存,**不在触发时再调 REST**。
- 订单簿 WS 在**周期开始时**会刷新订阅(`refreshAndSubscribe`:每 25 秒或检测到周期切换时),此时对当前周期内所有启用且 minSpreadMode=AUTO 的策略,按 `(intervalSeconds, periodStartUnix)` 异步拉取该周期前 20 根已收盘 K 线(REST `endTime = periodStartUnix * 1000`),按 Up/Down 分别算 avgSpread × 0.8(含 IQR 剔除)并写入缓存。该周期内后续触发时直接用缓存,**不在触发时再调 REST**。
- 若某周期未做预计算(如服务刚启动且尚未到刷新时机),触发时仍会按需调用 `computeAndCache` 并缓存,保证逻辑正确。
- 前端「自动最小价差」接口仅作**预览**,实际下单校验不依赖该接口。
- **计算细节**
- 历史 30 根:币安 REST `GET /api/v3/klines?symbol=BTCUSDC&interval=5m|15m&limit=30`(或 31 取前 30 根已收盘),每根格式为 [openTime, open, high, low, close, ...]。
- 历史 20 根:币安 REST `GET /api/v3/klines?symbol=BTCUSDC&interval=5m|15m&limit=20`(或 21 取前 20 根已收盘),每根格式为 [openTime, open, high, low, close, ...]。
- **DownoutcomeIndex=1**:筛选 close < open,价差 = open close,得到价差序列 → **异常值剔除(IQR** → 对剩余价差求平均,再 × 0.8 → minSpreadDown。
- **UpoutcomeIndex=0**:筛选 close > open,价差 = close open,得到价差序列 → **异常值剔除(IQR** → 对剩余价差求平均,再 × 0.8 → minSpreadUp。
- **异常值剔除**:见上文「异常值剔除」;剔除后再平均。若剔除后剩余样本 &lt; 3,则不剔除,用全部价差样本求平均。
- 若无满足方向的 K 线(例如 30 根里没有一根 close < open),可降级:用全部 30 根的 |closeopen| 平均 × 0.8,或返回 0/不校验,具体产品可定。
- 若无满足方向的 K 线(例如 20 根里没有一根 close < open),可降级:用全部 20 根的 |closeopen| 平均 × 0.8,或返回 0/不校验,具体产品可定。
- **触发时使用**
- 当前要下单的是 outcomeIndex0=Up, 1=Down),取对应的 minSpreadUp 或 minSpreadDown 作为 effectiveMinSpread,再与 |close open| 比较。
@@ -110,7 +110,7 @@
| 模块 | 职责 |
|------|------|
| **BinanceKlineService(新)** | 1)订阅币安 WSBTCUSDC 的 5m、15m K 线流(可按需只订阅有策略使用的周期)。<br>2)维护「当前周期」数据:以 periodStartUnix(或 K 线 t 对齐)为 key,存 (open, close)K 线 WS 推送时更新 close,新周期首条推送时更新 open。<br>3)提供 getCurrentOpenClose(symbol, intervalSeconds, periodStartUnix) → (open, close)?,供执行层价差校验使用。 |
| **BinanceKlineAutoSpreadService 或合入上者(新)** | 1)按**周期**拉取:以 periodStartUnix 为界,REST 拉取该周期前的 30 根已收盘 K 线。<br>2)按 Up/Down 得到价差序列 → **IQR 异常值剔除** → 对剩余价差求平均 × 0.8,缓存 (intervalSeconds, periodStartUnix) → (minSpreadUp, minSpreadDown)。<br>3)提供 getAutoMinSpread(intervalSeconds, periodStartUnix, outcomeIndex) 与 computeAndCache(intervalSeconds, periodStartUnix)。**周期开始时**由 CryptoTailOrderbookWsService 在 refreshAndSubscribe 后对当前周期内 AUTO 策略预调 computeAndCache;触发时直接用缓存,未命中时再按需计算。 |
| **BinanceKlineAutoSpreadService 或合入上者(新)** | 1)按**周期**拉取:以 periodStartUnix 为界,REST 拉取该周期前的 20 根已收盘 K 线。<br>2)按 Up/Down 得到价差序列 → **IQR 异常值剔除** → 对剩余价差求平均 × 0.8,缓存 (intervalSeconds, periodStartUnix) → (minSpreadUp, minSpreadDown)。<br>3)提供 getAutoMinSpread(intervalSeconds, periodStartUnix, outcomeIndex) 与 computeAndCache(intervalSeconds, periodStartUnix)。**周期开始时**由 CryptoTailOrderbookWsService 在 refreshAndSubscribe 后对当前周期内 AUTO 策略预调 computeAndCache;触发时直接用缓存,未命中时再按需计算。 |
| **CryptoTailStrategy(实体)** | 新增字段建议:minSpreadModeNONE/FIXED/AUTO)、minSpreadValue(固定时使用;AUTO 时可为空或存上次计算值用于展示)。 |
| **CryptoTailStrategyExecutionService(现有)** | 在 tryTriggerWithPriceFromWs 与 runCycle 分支中,在调用 placeOrderForTrigger 前:若 minSpreadMode != NONE,则取 open/close 与 effectiveMinSpread,校验 \|closeopen\| >= effectiveMinSpread;不通过则 return,不调用 placeOrderForTrigger。 |
| **CryptoTailOrderbookWsService(现有)** | 仍只根据 CLOB bestBid 触发;价差校验在执行层统一做。**新增**refreshAndSubscribe 完成后,对当前周期内所有启用且 minSpreadMode=AUTO 的策略,异步调用 BinanceKlineAutoSpreadService.computeAndCache,在周期开始即预计算最小价差。 |
@@ -169,7 +169,7 @@ sequenceDiagram
### 6.2 自动(AUTO)时序图
自动模式:不在保存策略时计算。**在该周期开始时就预计算**(订单簿 WS 刷新订阅时对该周期内 AUTO 策略异步拉 30 根历史 K 线并计算、缓存);触发时直接用缓存,同一周期内复用。
自动模式:不在保存策略时计算。**在该周期开始时就预计算**(订单簿 WS 刷新订阅时对该周期内 AUTO 策略异步拉 20 根历史 K 线并计算、缓存);触发时直接用缓存,同一周期内复用。
```mermaid
sequenceDiagram
@@ -213,8 +213,8 @@ sequenceDiagram
Orderbook->>Orderbook: refreshAndSubscribe() → buildSubscriptionMap() → newMap
Orderbook->>Orderbook: precomputeAutoMinSpreadForCurrentPeriods(newMap)
Orderbook->>AutoSpread: computeAndCache(intervalSeconds, periodStartUnix) [异步]
AutoSpread->>BinanceREST: GET /api/v3/klines?symbol=BTCUSDC&interval=15m&limit=30&endTime=periodStart*1000
BinanceREST-->>AutoSpread: 30 根已收盘 K 线
AutoSpread->>BinanceREST: GET /api/v3/klines?symbol=BTCUSDC&interval=15m&limit=20&endTime=periodStart*1000
BinanceREST-->>AutoSpread: 20 根已收盘 K 线
AutoSpread->>AutoSpread: 按 Up/Down 拆价差 → IQR 剔除 → 平均×0.8 → 缓存
Note over CLOB_WS,Exec: 同一周期内再次触发(如另一 outcome 或再次 bestBid
+1 -1
View File
@@ -1447,7 +1447,7 @@
"timeWindowStartLEEnd": "Window start must not be greater than end",
"timeWindowExceed": "Time window must not exceed period length",
"minSpreadMode": "Min spread",
"minSpreadModeTip": "Whether to place an order is based on the spread between open and close in the current period. Auto: system computes a suggested spread from the last 30 klines (updated each period); Fixed: you enter a value (e.g. 30), order only when spread ≥ that value; None: no spread check, order when price is in range.",
"minSpreadModeTip": "Whether to place an order is based on the spread between open and close in the current period. Auto: system computes a suggested spread from the last 20 klines (updated each period); Fixed: you enter a value (e.g. 30), order only when spread ≥ that value; None: no spread check, order when price is in range.",
"minSpreadModeNone": "None",
"minSpreadModeFixed": "Fixed",
"minSpreadModeAuto": "Auto",
+1 -1
View File
@@ -1446,7 +1446,7 @@
"timeWindowStartLEEnd": "时间区间开始不能大于结束",
"timeWindowExceed": "时间区间不能超过周期长度",
"minSpreadMode": "最小价差",
"minSpreadModeTip": "根据当前周期开盘价与收盘价的价差决定是否下单。自动:系统按历史 30 根 K 线计算建议价差(每周期更新);固定:您输入一个数值(如 30),仅当价差 ≥ 该值时才下单;无:不校验价差,满足价格区间即下单。",
"minSpreadModeTip": "根据当前周期开盘价与收盘价的价差决定是否下单。自动:系统按历史 20 根 K 线计算建议价差(每周期更新);固定:您输入一个数值(如 30),仅当价差 ≥ 该值时才下单;无:不校验价差,满足价格区间即下单。",
"minSpreadModeNone": "无",
"minSpreadModeFixed": "固定",
"minSpreadModeAuto": "自动",
+1 -1
View File
@@ -1447,7 +1447,7 @@
"timeWindowStartLEEnd": "時間區間開始不能大於結束",
"timeWindowExceed": "時間區間不能超過週期長度",
"minSpreadMode": "最小價差",
"minSpreadModeTip": "依當前週期開盤價與收盤價的價差決定是否下單。自動:系統依歷史 30 根 K 線計算建議價差(每週期更新);固定:您輸入一個數值(如 30),僅當價差 ≥ 該值時才下單;無:不校驗價差,滿足價格區間即下單。",
"minSpreadModeTip": "依當前週期開盤價與收盤價的價差決定是否下單。自動:系統依歷史 20 根 K 線計算建議價差(每週期更新);固定:您輸入一個數值(如 30),僅當價差 ≥ 該值時才下單;無:不校驗價差,滿足價格區間即下單。",
"minSpreadModeNone": "無",
"minSpreadModeFixed": "固定",
"minSpreadModeAuto": "自動",