Throttle OKX amp-stats candle pagination and retry on 429.

Add page pauses, exponential backoff, and cooldown before swap fallback to avoid rate limits.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
dekun
2026-07-23 02:57:24 +08:00
parent 0a9e3aa95c
commit 910c938d0a
5 changed files with 142 additions and 35 deletions
+1
View File
@@ -53,6 +53,7 @@
K 线粒度:**1H**(与整点对齐);价源优先 OKX 指数(ETH-USD / BTC-USD),失败再降级永续标记。 K 线粒度:**1H**(与整点对齐);价源优先 OKX 指数(ETH-USD / BTC-USD),失败再降级永续标记。
近期 K 线接口约仅 **1440** 根(1H≈60 天);更长周期自动续拉 `history-index-candles` / `history-candles` 近期 K 线接口约仅 **1440** 根(1H≈60 天);更长周期自动续拉 `history-index-candles` / `history-candles`
分页带间隔,遇 OKX **429** 会自动退避重试(长周期首次会慢一些)。
--- ---
+70 -18
View File
@@ -7,6 +7,7 @@ from __future__ import annotations
import csv import csv
import io import io
import statistics import statistics
import time
from datetime import date, datetime, timedelta from datetime import date, datetime, timedelta
from typing import Any, Callable, Optional from typing import Any, Callable, Optional
from zoneinfo import ZoneInfo from zoneinfo import ZoneInfo
@@ -401,6 +402,51 @@ def _parse_okx_candle_row(row: list) -> Optional[dict[str, Any]]:
return {"ts": ts, "o": o, "h": h, "l": l, "c": c} return {"ts": ts, "o": o, "h": h, "l": l, "c": c}
def _okx_get_json(
client: httpx.Client,
url: str,
params: dict[str, str],
*,
retries: int = 8,
) -> dict[str, Any]:
"""GET OKX 公共行情;遇 429 指数退避重试."""
last_err: Optional[BaseException] = None
for attempt in range(max(1, int(retries))):
try:
r = client.get(url, params=params)
if r.status_code == 429:
wait = min(12.0, 0.7 * (2**attempt))
time.sleep(wait)
last_err = httpx.HTTPStatusError(
f"429 Too Many Requests for url '{r.url}'",
request=r.request,
response=r,
)
continue
r.raise_for_status()
body = r.json()
if not isinstance(body, dict):
raise RuntimeError("OKX 返回非对象 JSON")
return body
except httpx.HTTPStatusError as exc:
status = exc.response.status_code if exc.response is not None else None
if status == 429 and attempt + 1 < retries:
wait = min(12.0, 0.7 * (2**attempt))
time.sleep(wait)
last_err = exc
continue
raise
except httpx.TransportError as exc:
if attempt + 1 < retries:
time.sleep(min(8.0, 0.5 * (2**attempt)))
last_err = exc
continue
raise
if last_err is not None:
raise last_err
raise RuntimeError("OKX 请求失败")
def fetch_okx_candles( def fetch_okx_candles(
*, *,
url: str, url: str,
@@ -412,10 +458,13 @@ def fetch_okx_candles(
timeout: float = 30.0, timeout: float = 30.0,
history_url: Optional[str] = None, history_url: Optional[str] = None,
max_pages: int = 200, max_pages: int = 200,
page_pause_sec: float = 0.12,
history_page_pause_sec: float = 0.22,
) -> list[dict[str, Any]]: ) -> list[dict[str, Any]]:
"""拉取 [since_ms, until_ms] 覆盖的 K 线(含边界). """拉取 [since_ms, until_ms] 覆盖的 K 线(含边界).
OKX 近期接口约仅 1440 根;更早需 history_* 端点续拉. OKX 近期接口约仅 1440 根;更早需 history_* 端点续拉.
分页带间隔,429 自动退避重试.
""" """
own = client is None own = client is None
client = client or httpx.Client( client = client or httpx.Client(
@@ -428,40 +477,33 @@ def fetch_okx_candles(
after: Optional[str] = None after: Optional[str] = None
active_url = url active_url = url
switched_history = False switched_history = False
empty_streak = 0 for page_i in range(max(20, int(max_pages))):
for _ in range(max(20, int(max_pages))): if page_i > 0:
pause = history_page_pause_sec if switched_history or "history" in active_url else page_pause_sec
if pause > 0:
time.sleep(pause)
params: dict[str, str] = {"instId": inst_id, "bar": bar, "limit": "100"} params: dict[str, str] = {"instId": inst_id, "bar": bar, "limit": "100"}
if after: if after:
params["after"] = after params["after"] = after
r = client.get(active_url, params=params) body = _okx_get_json(client, active_url, params)
r.raise_for_status()
body = r.json()
if str(body.get("code") or "") not in ("0", "0.0", ""): if str(body.get("code") or "") not in ("0", "0.0", ""):
raise RuntimeError(body.get("msg") or f"OKX error {body.get('code')}") raise RuntimeError(body.get("msg") or f"OKX error {body.get('code')}")
data = body.get("data") or [] data = body.get("data") or []
if not data: if not data:
empty_streak += 1
# 近期接口到头 → 切历史端点再试 # 近期接口到头 → 切历史端点再试
if ( if history_url and not switched_history and after is not None:
history_url
and not switched_history
and after is not None
):
active_url = history_url active_url = history_url
switched_history = True switched_history = True
empty_streak = 0 time.sleep(max(history_page_pause_sec, 0.35))
continue continue
break break
empty_streak = 0
oldest_ts = None oldest_ts = None
newest_in_page = None
for row in data: for row in data:
parsed = _parse_okx_candle_row(row) parsed = _parse_okx_candle_row(row)
if not parsed: if not parsed:
continue continue
ts = int(parsed["ts"]) ts = int(parsed["ts"])
oldest_ts = ts if oldest_ts is None else min(oldest_ts, ts) oldest_ts = ts if oldest_ts is None else min(oldest_ts, ts)
newest_in_page = ts if newest_in_page is None else max(newest_in_page, ts)
if ts < since_ms - 3600 * 1000: if ts < since_ms - 3600 * 1000:
continue continue
if ts > until_ms + 3600 * 1000: if ts > until_ms + 3600 * 1000:
@@ -476,6 +518,7 @@ def fetch_okx_candles(
if history_url and not switched_history: if history_url and not switched_history:
active_url = history_url active_url = history_url
switched_history = True switched_history = True
time.sleep(max(history_page_pause_sec, 0.35))
continue continue
break break
after = str(oldest_ts) after = str(oldest_ts)
@@ -488,6 +531,7 @@ def fetch_okx_candles(
): ):
active_url = history_url active_url = history_url
switched_history = True switched_history = True
time.sleep(max(history_page_pause_sec, 0.35))
return [out[k] for k in sorted(out.keys())] return [out[k] for k in sorted(out.keys())]
finally: finally:
if own: if own:
@@ -508,6 +552,7 @@ def fetch_symbol_bars(
bars = fetch_fn(inst_id=meta["index_inst"], since_ms=since_ms, until_ms=until_ms) bars = fetch_fn(inst_id=meta["index_inst"], since_ms=since_ms, until_ms=until_ms)
return bars, f"okx_index:{meta['index_inst']}", meta["index_inst"] return bars, f"okx_index:{meta['index_inst']}", meta["index_inst"]
index_err: Optional[BaseException] = None
try: try:
bars = fetch_okx_candles( bars = fetch_okx_candles(
url=OKX_INDEX_CANDLES, url=OKX_INDEX_CANDLES,
@@ -518,9 +563,12 @@ def fetch_symbol_bars(
) )
if bars: if bars:
return bars, f"okx_index:{meta['index_inst']}", meta["index_inst"] return bars, f"okx_index:{meta['index_inst']}", meta["index_inst"]
except Exception: except Exception as exc:
bars = [] index_err = exc
# 指数侧已触发限频时先冷却,再降级永续,避免连环 429
time.sleep(1.2)
try:
bars = fetch_okx_candles( bars = fetch_okx_candles(
url=OKX_SWAP_CANDLES, url=OKX_SWAP_CANDLES,
history_url=OKX_HISTORY_SWAP_CANDLES, history_url=OKX_HISTORY_SWAP_CANDLES,
@@ -528,8 +576,12 @@ def fetch_symbol_bars(
since_ms=since_ms, since_ms=since_ms,
until_ms=until_ms, until_ms=until_ms,
) )
except Exception as exc:
detail = f"index={index_err}; swap={exc}" if index_err else str(exc)
raise RuntimeError(f"OKX K线拉取失败({detail})") from exc
if not bars: if not bars:
raise RuntimeError("OKX 指数与永续 K 线均无数据") detail = f"index={index_err}" if index_err else "empty"
raise RuntimeError(f"OKX 指数与永续 K 线均无数据({detail})")
return bars, f"okx_swap:{meta['swap_inst']}", meta["swap_inst"] return bars, f"okx_swap:{meta['swap_inst']}", meta["swap_inst"]
+1 -1
View File
@@ -274,7 +274,7 @@
const startHour = Number(el("amp-start-hour")?.value || 16); const startHour = Number(el("amp-start-hour")?.value || 16);
const period = el("amp-period")?.value || "2m"; const period = el("amp-period")?.value || "2m";
const customDays = Number(el("amp-custom-days")?.value || 60); const customDays = Number(el("amp-custom-days")?.value || 60);
setStatus("计算中…(首次拉取 OKX K 线可能需数十秒)"); setStatus("计算中…(长周期会分页拉 OKX,遇限频会自动重试,请稍候)");
try { try {
const data = await apiFetch("/api/amp-stats/compute", { const data = await apiFetch("/api/amp-stats/compute", {
method: "POST", method: "POST",
+1 -1
View File
@@ -1508,7 +1508,7 @@
<script src="/assets/funds.js?v=20260717-funds-scroll-fix"></script> <script src="/assets/funds.js?v=20260717-funds-scroll-fix"></script>
<script src="/assets/dashboard.js?v=20260720-dash-sl-tp"></script> <script src="/assets/dashboard.js?v=20260720-dash-sl-tp"></script>
<script src="/assets/strategy.js?v=9"></script> <script src="/assets/strategy.js?v=9"></script>
<script src="/assets/amp_stats.js?v=4"></script> <script src="/assets/amp_stats.js?v=5"></script>
<script src="/assets/help.js?v=1"></script> <script src="/assets/help.js?v=1"></script>
<script src="/assets/logs.js?v=1"></script> <script src="/assets/logs.js?v=1"></script>
<script src="/assets/ai_review_render.js?v=3"></script> <script src="/assets/ai_review_render.js?v=3"></script>
+61 -7
View File
@@ -14,6 +14,8 @@ from lib.hub.amp_stats_lib import (
window_bounds_for_settlement, window_bounds_for_settlement,
) )
import httpx
TZ = ZoneInfo("Asia/Shanghai") TZ = ZoneInfo("Asia/Shanghai")
@@ -198,11 +200,17 @@ class AmpStatsLibTests(unittest.TestCase):
calls: list[str] = [] calls: list[str] = []
class FakeResp: class FakeResp:
def __init__(self, data): def __init__(self, data, url="https://x", status_code=200):
self._data = data self._data = data
self.status_code = status_code
self.url = url
self.request = httpx.Request("GET", url)
def raise_for_status(self): def raise_for_status(self):
return None if self.status_code >= 400:
raise httpx.HTTPStatusError(
"err", request=self.request, response=self
)
def json(self): def json(self):
return {"code": "0", "data": self._data} return {"code": "0", "data": self._data}
@@ -214,14 +222,14 @@ class AmpStatsLibTests(unittest.TestCase):
# recent: only 2 pages then empty; history continues # recent: only 2 pages then empty; history continues
if "history" not in url: if "history" not in url:
if after is None: if after is None:
return FakeResp([["2000", "1", "2", "0.5", "1.5"], ["1900", "1", "2", "0.5", "1.5"]]) return FakeResp([["2000", "1", "2", "0.5", "1.5"], ["1900", "1", "2", "0.5", "1.5"]], url=url)
if after == "1900": if after == "1900":
return FakeResp([]) # recent exhausted return FakeResp([], url=url) # recent exhausted
return FakeResp([]) return FakeResp([], url=url)
# history # history
if after == "1900": if after == "1900":
return FakeResp([["1800", "1", "2", "0.5", "1.5"], ["1000", "1", "2", "0.5", "1.5"]]) return FakeResp([["1800", "1", "2", "0.5", "1.5"], ["1000", "1", "2", "0.5", "1.5"]], url=url)
return FakeResp([]) return FakeResp([], url=url)
def close(self): def close(self):
return None return None
@@ -234,11 +242,57 @@ class AmpStatsLibTests(unittest.TestCase):
until_ms=3000, until_ms=3000,
client=FakeClient(), client=FakeClient(),
max_pages=10, max_pages=10,
page_pause_sec=0,
history_page_pause_sec=0,
) )
self.assertTrue(any("history-index-candles" in u for u in calls)) self.assertTrue(any("history-index-candles" in u for u in calls))
self.assertGreaterEqual(len(bars), 3) self.assertGreaterEqual(len(bars), 3)
self.assertEqual(bars[0]["ts"], 1000) self.assertEqual(bars[0]["ts"], 1000)
def test_fetch_retries_on_429(self):
from lib.hub.amp_stats_lib import fetch_okx_candles
import httpx as _httpx
hits = {"n": 0}
class FakeResp:
def __init__(self, status_code, data=None):
self.status_code = status_code
self.url = "https://www.okx.com/api/v5/market/history-candles"
self.request = _httpx.Request("GET", self.url)
self._data = data or []
def raise_for_status(self):
if self.status_code >= 400:
raise _httpx.HTTPStatusError("429", request=self.request, response=self)
def json(self):
return {"code": "0", "data": self._data}
class FakeClient:
def get(self, url, params=None):
hits["n"] += 1
if hits["n"] < 3:
return FakeResp(429)
return FakeResp(200, [["1000", "1", "2", "0.5", "1.5"]])
def close(self):
return None
bars = fetch_okx_candles(
url="https://www.okx.com/api/v5/market/candles",
history_url=None,
inst_id="ETH-USDT-SWAP",
since_ms=1000,
until_ms=2000,
client=FakeClient(),
max_pages=3,
page_pause_sec=0,
history_page_pause_sec=0,
)
self.assertGreaterEqual(hits["n"], 3)
self.assertEqual(len(bars), 1)
def test_compute_with_mock_fetch(self): def test_compute_with_mock_fetch(self):
now = datetime(2026, 7, 22, 18, 0, tzinfo=TZ) now = datetime(2026, 7, 22, 18, 0, tzinfo=TZ)