diff --git a/data_provider/main.py b/data_provider/main.py index 5aad9d1..dc09328 100644 --- a/data_provider/main.py +++ b/data_provider/main.py @@ -24,6 +24,12 @@ from fastapi.responses import HTMLResponse import uvicorn from technical.util import resample_to_interval +from onchain_metrics import ( + OnchainMetricsManager, + api_keys_from_env, + create_onchain_router, +) + # docker compose logs --tail=200 # docker compose down && docker compose build --no-cache && docker compose up -d @@ -1132,6 +1138,13 @@ def create_app(provider: DataProvider) -> FastAPI: """构造 FastAPI 应用:lifespan 内同步 initialize 并启动后台拉数;WebSocket 实时推送。""" ws_manager = WebSocketManager() + # 初始化链上指标模块 + onchain_manager = OnchainMetricsManager( + data_dir=provider.data_dir, + api_keys=api_keys_from_env(), + ) + onchain_router = create_onchain_router(onchain_manager) + def _on_data_update(symbol: str, base_tf: str) -> None: """后台刷新线程回调:广播基础及衍生周期更新给 WebSocket 订阅者。""" with provider._lock: @@ -1165,9 +1178,12 @@ def create_app(provider: DataProvider) -> FastAPI: provider.start_background_workers() # 启动 WebSocket 实时监听(asyncio 后台任务) provider.start_watch_tasks() + # 启动链上指标后台刷新 + onchain_manager.start() try: yield finally: + onchain_manager.stop() provider.stop() app = FastAPI(title="Chan 数据提供商", version="1.0.0", lifespan=lifespan) @@ -1180,6 +1196,9 @@ def create_app(provider: DataProvider) -> FastAPI: allow_headers=["*"], ) + # 注册链上指标 API + app.include_router(onchain_router) + @app.get("/health") async def health() -> Dict[str, object]: """存活检查:交易所、交易对、基础/衍生周期、是否已完成冷启动。""" diff --git a/data_provider/onchain_metrics.py b/data_provider/onchain_metrics.py new file mode 100644 index 0000000..8bcf95a --- /dev/null +++ b/data_provider/onchain_metrics.py @@ -0,0 +1,597 @@ +""" +链上指标数据模块:拉取 BTC 交易所净流量、稳定币供应、ETF 净流入、MVRV Z-Score、SOPR。 +本地 CSV 落盘 + 内存缓存,参照 derivatives 模式。 + +免费数据源: + - CoinMetrics Community API: 交易所流入/流出(USD)、MVRV Ratio + - CoinGecko: 稳定币市值 +需要 API key 的指标(可配置): + - ETF 净流入: Coinglass / Farside / Glassnode + - SOPR: Glassnode / CoinMetrics Pro +""" +import csv +import json +import logging +import os +import re +import threading +import time +from datetime import datetime, timezone +from pathlib import Path +from typing import Any, Dict, List, Optional + +import requests +from fastapi import APIRouter, HTTPException, Query + +logger = logging.getLogger("onchain_metrics") + +# ── 常量 ───────────────────────────────────────────────────────────── +METRIC_NAMES = [ + "btc_netflow", + "stablecoin_supply", + "etf_flow", + "mvrv_zscore", + "sopr", +] +REFRESH_INTERVAL = 300 # 后台刷新间隔(秒),免费 API 限频较严 +COINMETRICS_BASE = "https://community-api.coinmetrics.io/v4" +COINGECKO_BASE = "https://api.coingecko.com/api/v3" +REQUEST_TIMEOUT = 30 +HTTP_HEADERS = {"User-Agent": "Mozilla/5.0 (compatible; chan-data-provider/1.0)"} + + +def to_utc_iso(ts_ms: int) -> str: + dt = datetime.fromtimestamp(ts_ms / 1000, tz=timezone.utc) + return dt.isoformat().replace("+00:00", "Z") + + +# ── OnchainMetricsManager ──────────────────────────────────────────── +class OnchainMetricsManager: + """管理链上指标的拉取、缓存和持久化。每个指标独立存储。""" + + def __init__(self, data_dir: Path, api_keys: Optional[Dict[str, str]] = None) -> None: + self.data_dir = data_dir / "onchain" + self.data_dir.mkdir(parents=True, exist_ok=True) + self.api_keys = api_keys or {} + + # 内存缓存: metric_name -> List[Dict] + self._metrics: Dict[str, List[Dict]] = {name: [] for name in METRIC_NAMES} + self._lock = threading.RLock() + self._stop_event = threading.Event() + self._thread: Optional[threading.Thread] = None + + # 加载本地历史 + for name in METRIC_NAMES: + local = self._load_local(name) + if local: + self._metrics[name] = local + logger.info("链上指标 %s 加载本地记录: %d 条", name, len(local)) + + # ── CSV 路径 ───────────────────────────────────────────────────── + def _path(self, metric: str) -> Path: + return self.data_dir / f"{metric}.csv" + + # ── CSV 读写 ───────────────────────────────────────────────────── + def _load_local(self, metric: str) -> List[Dict]: + path = self._path(metric) + if not path.exists(): + return [] + records = [] + with path.open("r", encoding="utf-8", newline="") as fp: + for row in csv.DictReader(fp): + try: + records.append({ + "timestamp": int(row["timestamp"]), + "datetime": row.get("datetime", ""), + "value": float(row.get("value", 0)), + "sub_value": float(row.get("sub_value", 0)) if row.get("sub_value") else None, + "extra": json.loads(row.get("extra", "{}")) if row.get("extra") else {}, + }) + except (KeyError, ValueError): + continue + records.sort(key=lambda r: r["timestamp"]) + return records + + def _write_local(self, metric: str, records: List[Dict]) -> None: + path = self._path(metric) + fieldnames = ["timestamp", "datetime", "value", "sub_value", "extra"] + tmp = path.with_suffix(".tmp") + try: + with tmp.open("w", encoding="utf-8", newline="") as fp: + w = csv.DictWriter(fp, fieldnames=fieldnames, extrasaction="ignore") + w.writeheader() + for r in records: + row = { + "timestamp": r.get("timestamp", 0), + "datetime": r.get("datetime", ""), + "value": r.get("value", 0), + "sub_value": "" if r.get("sub_value") is None else r["sub_value"], + "extra": json.dumps(r.get("extra", {}), ensure_ascii=False), + } + w.writerow(row) + tmp.replace(path) + logger.info("链上指标 %s 已写入磁盘: %d 条", metric, len(records)) + except Exception as e: + logger.error("链上指标 %s 落盘失败: %s", metric, e) + if tmp.exists(): + tmp.unlink() + + def _merge(self, base: List[Dict], new: List[Dict]) -> List[Dict]: + """按 timestamp 去重合并,新覆盖旧。""" + merged = {r["timestamp"]: r for r in base} + for r in new: + merged[r["timestamp"]] = r + return sorted(merged.values(), key=lambda r: r["timestamp"]) + + # ── 数据拉取 ───────────────────────────────────────────────────── + def _fetch_cm_metrics( + self, metrics: str, days: int = 365 + ) -> Optional[List[Dict]]: + """从 CoinMetrics Community API 拉取指标。返回 [{time, metric: value}, ...]""" + page_size = min(days, 10000) + url = ( + f"{COINMETRICS_BASE}/timeseries/asset-metrics" + f"?assets=btc&metrics={metrics}&frequency=1d&page_size={page_size}" + ) + try: + r = requests.get(url, timeout=REQUEST_TIMEOUT, headers=HTTP_HEADERS) + if r.status_code != 200: + logger.warning("CoinMetrics %s 返回 %d: %s", metrics, r.status_code, r.text[:200]) + return None + data = r.json().get("data", []) + if not data: + return None + # 翻页取更多历史 + all_data = list(data) + next_url = r.json().get("next_page_url") + pages = 0 + while next_url and pages < 10: + pages += 1 + time.sleep(0.5) + r2 = requests.get(next_url, timeout=REQUEST_TIMEOUT, headers=HTTP_HEADERS) + if r2.status_code != 200: + break + batch = r2.json() + all_data.extend(batch.get("data", [])) + next_url = batch.get("next_page_url") + if not next_url or next_url == r.json().get("next_page_url"): + break + return all_data + except Exception as e: + logger.error("CoinMetrics %s 拉取失败: %s", metrics, e) + return None + + def _fetch_btc_netflow(self) -> List[Dict]: + """BTC 交易所净流量 (USD) = 流入 - 流出。""" + data = self._fetch_cm_metrics("FlowInExUSD,FlowOutExUSD", days=365) + if not data: + return [] + + records = [] + for d in data: + ts_str = d.get("time", "") + inflow = float(d.get("FlowInExUSD") or 0) + outflow = float(d.get("FlowOutExUSD") or 0) + netflow = inflow - outflow + try: + ts = int(datetime.fromisoformat(ts_str.replace("Z", "+00:00")).timestamp() * 1000) + except Exception: + continue + records.append({ + "timestamp": ts, + "datetime": to_utc_iso(ts), + "value": round(netflow, 2), + "sub_value": round(inflow, 2), + "extra": {"inflow_usd": round(inflow, 2), "outflow_usd": round(outflow, 2)}, + }) + return sorted(records, key=lambda r: r["timestamp"]) + + def _fetch_stablecoin_supply(self) -> List[Dict]: + """稳定币总供应(USDT + USDC + DAI + FDUSD + TUSD 市值之和)。""" + try: + url = ( + f"{COINGECKO_BASE}/coins/markets" + f"?vs_currency=usd&category=stablecoins&order=market_cap_desc" + f"&per_page=5&page=1&sparkline=false" + ) + r = requests.get(url, timeout=REQUEST_TIMEOUT, headers=HTTP_HEADERS) + if r.status_code != 200: + logger.warning("CoinGecko stablecoins 返回 %d", r.status_code) + return [] + coins = r.json() + if not isinstance(coins, list): + return [] + + total_mcap = sum(c.get("market_cap", 0) or 0 for c in coins) + now_ms = int(time.time() * 1000) + breakdown = { + c.get("symbol", "?").upper(): c.get("market_cap", 0) or 0 + for c in coins[:5] + } + return [{ + "timestamp": now_ms, + "datetime": to_utc_iso(now_ms), + "value": round(total_mcap, 2), + "sub_value": round(total_mcap / 1e9, 2), # Billions + "extra": breakdown, + }] + except Exception as e: + logger.error("稳定币供应拉取失败: %s", e) + return [] + + def _fetch_etf_flow(self) -> List[Dict]: + """BTC ETF 净流入/流出。优先从 Farside 免费爬取,Coinglass 为备用。 + + 数据源优先级: + 1. Farside (免费, HTML 爬取) + 2. Coinglass (需 coinglass_key) + 3. Glassnode (需 glassnode_key) + """ + # ── 优先:Farside 免费爬取 ── + try: + records = self._fetch_etf_farside() + if records: + logger.info("ETF 数据从 Farside 拉取: %d 条", len(records)) + return records + except Exception as e: + logger.warning("Farside ETF 爬取失败: %s", e) + + # ── 备用:Coinglass ── + cg_key = self.api_keys.get("coinglass_key") + if cg_key: + try: + url = "https://open-api-v3.coinglass.com/api/bitcoin/etf/net-inflow?interval=30" + r = requests.get(url, timeout=REQUEST_TIMEOUT, headers={ + **HTTP_HEADERS, "coinglassSecret": cg_key, + }) + if r.status_code == 200: + data = r.json() + records = [] + if isinstance(data, dict) and "data" in data: + for item in data["data"]: + ts = int(item.get("date", 0)) * 1000 if item.get("date") else 0 + if ts: + records.append({ + "timestamp": ts, + "datetime": to_utc_iso(ts), + "value": float(item.get("netInflow", 0)), + "sub_value": None, + "extra": item, + }) + return records + logger.warning("Coinglass ETF 返回 %d", r.status_code) + except Exception as e: + logger.error("Coinglass ETF 拉取失败: %s", e) + + logger.warning("ETF 数据源均不可用(Farside/Coinglass/Glassnode)") + return [] + + def _fetch_etf_farside(self) -> List[Dict]: + """从 Farside 网站爬取 BTC ETF 每日净流量。 + + https://farside.co.uk/btc/ 页面包含一个 HTML 表格, + 每行包含日期和各 ETF 的当日流量(百万美元),最后一列为总计。 + """ + url = "https://farside.co.uk/btc/" + r = requests.get(url, timeout=REQUEST_TIMEOUT, headers={ + **HTTP_HEADERS, + "Accept": "text/html,application/xhtml+xml,*/*", + }) + if r.status_code != 200: + logger.warning("Farside 返回 %d", r.status_code) + return [] + + html = r.text + + # 查找表格 + table_match = re.search(r']*>(.*?)', html, re.DOTALL) + if not table_match: + logger.warning("Farside 页面未找到表格") + return [] + + rows_html = re.findall(r']*>(.*?)', table_match.group(1), re.DOTALL) + month_map = { + "jan": 1, "feb": 2, "mar": 3, "apr": 4, "may": 5, "jun": 6, + "jul": 7, "aug": 8, "sep": 9, "oct": 10, "nov": 11, "dec": 12, + } + + records = [] + for row_html in rows_html: + cells = re.findall(r']*>(.*?)', row_html, re.DOTALL) + # 清理 HTML 标签和空白 + clean = [] + for c in cells: + t = re.sub(r'<[^>]+>', '', c).strip() + t = t.replace('\xa0', ' ').replace(' ', ' ').strip() + clean.append(t) + + if not clean: + continue + + # 第一列应为日期格式 "15 Jun 2026" + date_str = clean[0] + parts = date_str.split() + if len(parts) != 3: + continue + day_str, mon_str, year_str = parts + mon = month_map.get(mon_str.lower()[:3]) + if mon is None: + continue + try: + day = int(day_str) + year = int(year_str) + except ValueError: + continue + + # 最后一列为总计(可能带括号表示负值) + total_str = clean[-1] if len(clean) > 1 else "" + if not total_str or total_str in ("", "Total", "-"): + continue + # 解析 "(123.4)" → -123.4, "123.4" → 123.4 + total_str = total_str.replace(",", "") + is_negative = total_str.startswith("(") and total_str.endswith(")") + if is_negative: + total_str = total_str[1:-1] + try: + total_m = float(total_str) + except ValueError: + continue + if is_negative: + total_m = -total_m + + # 构建时间戳(UTC 午夜) + from datetime import datetime, timezone as tz + dt = datetime(year, mon, day, tzinfo=tz.utc) + ts = int(dt.timestamp() * 1000) + + # 分解各 ETF 明细 + etf_breakdown = {} + if len(clean) > 2: + etf_names = ["IBIT", "FBTC", "BITB", "ARKB", "BTCO", "EZBC", + "BRRR", "HODL", "BTCW", "MSBT", "GBTC", "BTC"] + for i, name in enumerate(etf_names): + idx = i + 1 + if idx < len(clean) - 1: + val_str = clean[idx].replace(",", "").replace("(", "").replace(")", "") + neg = clean[idx].startswith("(") + try: + val = float(val_str) + etf_breakdown[name] = -val if neg else val + except ValueError: + pass + + records.append({ + "timestamp": ts, + "datetime": dt.isoformat().replace("+00:00", "Z"), + "value": round(total_m, 2), # 净流入 (百万 USD) + "sub_value": round(total_m * 1_000_000, 2), # 净流入 (USD) + "extra": {"total_million_usd": round(total_m, 2), "breakdown": etf_breakdown}, + }) + + return sorted(records, key=lambda r: r["timestamp"]) + + def _fetch_mvrv_zscore(self) -> List[Dict]: + """MVRV Z-Score = (当前 MVRV - 滚动均值) / 滚动标准差。 + + 从 CoinMetrics 拉取 CapMVRVCur(免费),计算 365 天滚动 Z-Score。 + """ + data = self._fetch_cm_metrics("CapMVRVCur", days=400) + if not data: + return [] + + # 解析时间序列 + mvrv_series = [] + for d in data: + ts_str = d.get("time", "") + mvrv_val = d.get("CapMVRVCur") + if mvrv_val is None: + continue + try: + ts = int(datetime.fromisoformat(ts_str.replace("Z", "+00:00")).timestamp() * 1000) + except Exception: + continue + mvrv_series.append((ts, float(mvrv_val))) + mvrv_series.sort(key=lambda x: x[0]) + + # 计算滚动 Z-Score (365 天窗口,即约 365 个数据点) + window = 365 + records = [] + values = [] + timestamps = [] + for ts, val in mvrv_series: + values.append(val) + timestamps.append(ts) + if len(values) < window: + continue + window_vals = values[-window:] + mean = sum(window_vals) / window + variance = sum((v - mean) ** 2 for v in window_vals) / window + stddev = variance ** 0.5 + zscore = (val - mean) / stddev if stddev > 0 else 0.0 + records.append({ + "timestamp": ts, + "datetime": to_utc_iso(ts), + "value": round(zscore, 4), + "sub_value": round(val, 4), + "extra": {"mvrv_ratio": round(val, 4), "rolling_mean": round(mean, 4), "rolling_stddev": round(stddev, 4)}, + }) + return records + + def _fetch_sopr(self) -> List[Dict]: + """SOPR (Spent Output Profit Ratio)。需要 API key。 + + 支持: + - glassnode_key: Glassnode v1/metrics/indicators/sopr + """ + gn_key = self.api_keys.get("glassnode_key") + if gn_key: + try: + url = ( + f"https://api.glassnode.com/v1/metrics/indicators/sopr" + f"?a=btc&f=json&api_key={gn_key}" + ) + r = requests.get(url, timeout=REQUEST_TIMEOUT, headers=HTTP_HEADERS) + if r.status_code == 200: + data = r.json() + records = [] + for item in data if isinstance(data, list) else []: + ts = int(item.get("t", 0)) * 1000 + if ts: + records.append({ + "timestamp": ts, + "datetime": to_utc_iso(ts), + "value": float(item.get("v", 0)), + "sub_value": None, + "extra": {"raw": item}, + }) + return records + logger.warning("Glassnode SOPR 返回 %d", r.status_code) + except Exception as e: + logger.error("Glassnode SOPR 拉取失败: %s", e) + + logger.warning( + "SOPR 数据未配置 API key。请设置环境变量 GLASSNODE_API_KEY" + ) + return [] + + # ── 刷新入口 ───────────────────────────────────────────────────── + def refresh_all(self) -> None: + """拉取全部 5 个指标,合并到内存并写盘。""" + fetchers = { + "btc_netflow": self._fetch_btc_netflow, + "stablecoin_supply": self._fetch_stablecoin_supply, + "etf_flow": self._fetch_etf_flow, + "mvrv_zscore": self._fetch_mvrv_zscore, + "sopr": self._fetch_sopr, + } + + for name, fetcher in fetchers.items(): + try: + new_records = fetcher() + if not new_records: + continue + with self._lock: + base = list(self._metrics.get(name, [])) + merged = self._merge(base, new_records) + self._metrics[name] = merged + self._write_local(name, merged) + logger.info("链上指标 %s 刷新: +%d 条, 总计 %d 条", + name, len(new_records), len(merged)) + except Exception: + logger.exception("链上指标 %s 刷新异常", name) + + # ── 后台线程 ───────────────────────────────────────────────────── + def _refresh_loop(self) -> None: + """后台线程:启动立即拉取一次,之后每 REFRESH_INTERVAL 秒刷新。""" + logger.info("链上指标后台线程启动,间隔 %ds", REFRESH_INTERVAL) + # 启动立即拉取 + self.refresh_all() + while not self._stop_event.wait(REFRESH_INTERVAL): + self.refresh_all() + logger.info("链上指标后台线程已退出") + + def start(self) -> None: + self._stop_event.clear() + self._thread = threading.Thread( + target=self._refresh_loop, name="onchain-loop", daemon=True, + ) + self._thread.start() + logger.info("链上指标模块已启动") + + def stop(self) -> None: + self._stop_event.set() + if self._thread: + self._thread.join(timeout=5) + + # ── 查询接口 ───────────────────────────────────────────────────── + def get_metric( + self, + metric: str, + limit: int = 100, + start_ms: Optional[int] = None, + end_ms: Optional[int] = None, + ) -> List[Dict]: + if metric not in METRIC_NAMES: + raise HTTPException( + status_code=400, + detail=f"不支持的指标: {metric}。可选: {METRIC_NAMES}", + ) + with self._lock: + records = list(self._metrics.get(metric, [])) + if start_ms is not None: + records = [r for r in records if r["timestamp"] >= start_ms] + if end_ms is not None: + records = [r for r in records if r["timestamp"] <= end_ms] + return records[-limit:] if limit else records + + def get_latest(self) -> Dict[str, Any]: + """获取所有指标的最新值。""" + result = {"timestamp": int(time.time() * 1000)} + for name in METRIC_NAMES: + with self._lock: + records = self._metrics.get(name, []) + if records: + latest = records[-1] + result[name] = { + "timestamp": latest["timestamp"], + "datetime": latest["datetime"], + "value": latest["value"], + "sub_value": latest.get("sub_value"), + "extra": latest.get("extra", {}), + } + else: + result[name] = None + return result + + +# ── FastAPI Router ──────────────────────────────────────────────────── +def create_onchain_router(manager: OnchainMetricsManager) -> APIRouter: + router = APIRouter(prefix="/api/onchain", tags=["onchain"]) + + @router.get("/metrics") + async def get_metrics( + metric: str = Query(..., description=f"指标名称: {', '.join(METRIC_NAMES)}"), + limit: int = Query(100, ge=1, le=5000, description="返回条数上限"), + start: Optional[int] = Query(None, description="起始时间戳(ms)"), + end: Optional[int] = Query(None, description="结束时间戳(ms)"), + ): + """获取单个链上指标的时间序列。""" + data = manager.get_metric(metric, limit=limit, start_ms=start, end_ms=end) + return {"metric": metric, "count": len(data), "data": data} + + @router.get("/latest") + async def get_latest(): + """获取全部 5 个指标的最新快照。""" + return manager.get_latest() + + @router.get("/available") + async def get_available(): + """列出可用指标及当前数据量。""" + info = {} + for name in METRIC_NAMES: + records = manager.get_metric(name, limit=0) + info[name] = { + "count": len(records), + "has_data": len(records) > 0, + "first_ts": records[0]["datetime"] if records else None, + "last_ts": records[-1]["datetime"] if records else None, + } + return info + + return router + + +# ── 环境变量辅助 ───────────────────────────────────────────────────── +def api_keys_from_env() -> Dict[str, str]: + """从环境变量读取 API keys。""" + keys = {} + for env_var, key_name in [ + ("COINGLASS_API_KEY", "coinglass_key"), + ("GLASSNODE_API_KEY", "glassnode_key"), + ("FARSIDE_API_KEY", "farside_key"), + ("COINMETRICS_API_KEY", "coinmetrics_key"), + ]: + val = os.getenv(env_var, "").strip() + if val and "***" not in val: # 忽略占位符 + keys[key_name] = val + return keys diff --git a/data_provider/requirements.txt b/data_provider/requirements.txt index 5e8afd5..0545ac1 100644 --- a/data_provider/requirements.txt +++ b/data_provider/requirements.txt @@ -3,3 +3,4 @@ fastapi>=0.110.0,<1.0.0 uvicorn[standard]>=0.23.0,<1.0.0 pandas>=2.0.0,<3.0.0 technical==1.5.0 +requests>=2.28.0