熔断器机制
概述
数据源配备熔断器(Circuit Breaker),连续失败后自动跳过请求,避免无效重试加重封禁。 代码中存在两套实现、三个实例,面向不同调用面,互不影响:
| 实现 | 位置 | 实例 | 参数 | 调用面 |
|---|---|---|---|---|
SourceCircuitBreaker | market_data_service/infra.py | em_breaker | 阈值 4 / 冷却 120s(EM_BREAKER_FAILURE_THRESHOLD / EM_BREAKER_COOLDOWN_SEC) | get_market_data() 东财降级 |
SourceCircuitBreaker | services/data_sources/stock_changes.py | push2ex_breaker | 阈值 2 / 冷却 30s(_create_breaker() 内联,不共用 EM_BREAKER_*) | 盘中异动监控(push2ex 域名) |
CircuitBreaker(core/resilience.py) | core/resilience.py,具名注册表 get_breaker(name) | eastmoney / tencent / tdx | 阈值 5 / reset_timeout 30s / 半开 1 次 | stock_service.py 的 _safe_call 路由 |
push2ex_breaker 的参数比主行情熔断更激进是有意为之:异动监控属体验型数据,
调用方自带 stale 兜底,东财不可达时每次轮询都白等 2~3s 完整失败链(实测 2.97s / 2.36s),
而主行情熔断阈值 4 需约 40s 才能触发。连续 2 次失败即熔断,轮询立刻返回 stale。
push2ex 与 push2 是不同域名、不同风控策略,因此使用独立实例,不与 em_breaker 混用。
状态机
| 状态 | 行为 |
|---|---|
| CLOSED | 正常放行所有请求 |
| OPEN | 跳过所有东财请求,直接降级 |
| HALF_OPEN | 允许一次探测请求,成功则恢复 |
实现代码
market_data_service/infra.py(SourceCircuitBreaker,两套实现中语义最简单的一版):
class SourceCircuitBreaker:
"""数据源熔断器:连续失败 N 次后进入 OPEN 状态,跳过请求;
冷却期后进入 HALF_OPEN 允许一次探测,成功则恢复。"""
def __init__(self, failure_threshold: int = 4, cooldown_sec: float = 120):
self._state: str = "CLOSED"
self._failure_count = 0
self._last_failure_time = 0.0
self._failure_threshold = failure_threshold
self._cooldown_sec = cooldown_sec
def can_request(self) -> bool:
if self._state == "CLOSED":
return True
if self._state == "OPEN":
if time.time() - self._last_failure_time >= self._cooldown_sec:
self._state = "HALF_OPEN"
return True # 允许一次探测
return False
return True
def record_success(self):
self._state = "CLOSED"
self._failure_count = 0
def record_failure(self):
self._last_failure_time = time.time()
if self._state == "HALF_OPEN":
self._state = "OPEN" # 探测失败,重新熔断
else:
self._failure_count += 1
if self._failure_count >= self._failure_threshold:
self._state = "OPEN"
logger.info(f"东财熔断器开启,{self._cooldown_sec}s 后恢复探测")
core/resilience.py::CircuitBreaker 是更完整的一版:带 threading.Lock(状态读取
线程安全)、half_open_requests 半开并发额度,并由 get_breaker(name) 按名注册复用
(_breakers 字典),状态经 get_all_breaker_stats() 由 GET /health/datasources
以 ok / degraded / down 三档对外暴露。stock_service.py 为三个数据源各注册一个实例,
并由 _safe_call 按label 关键字路由:
_breaker_eastmoney = get_breaker("eastmoney", failure_threshold=5, reset_timeout=30.0)
_breaker_tencent = get_breaker("tencent", failure_threshold=5, reset_timeout=30.0)
_breaker_tdx = get_breaker("tdx", failure_threshold=5, reset_timeout=30.0)
# 东财需要限流(高频会被封 IP),其他数据源不需要
_rl_eastmoney = get_rate_limiter("eastmoney", requests_per_second=1.5, max_burst=3)
| label 关键字 | 熔断器 | 限流器 |
|---|---|---|
eastmoney / em( / search / industry | _breaker_eastmoney | _rl_eastmoney(1.5 req/s,突发 3) |
tencent / stock( | _breaker_tencent | — |
kline / finance | _breaker_tdx | — |
配置参数
# market_data_service/infra.py — 主行情熔断器(可由环境变量覆盖)
em_breaker = SourceCircuitBreaker(
failure_threshold=s.EM_BREAKER_FAILURE_THRESHOLD, # 默认 4:连续 4 次失败触发熔断
cooldown_sec=s.EM_BREAKER_COOLDOWN_SEC, # 默认 120:熔断后冷却 120 秒
)
# services/data_sources/stock_changes.py — 盘中异动熔断器(更激进,写死不走配置)
push2ex_breaker = SourceCircuitBreaker(failure_threshold=2, cooldown_sec=30)
触发场景
| 场景 | 触发行为 |
|---|---|
| 东财返回 403 | 立即 record_failure()(风控信号,不重试) |
| 东财返回 429 | RateLimitError 分支 record_failure(),fund_flow 改走新浪备用源 |
| 网络超时 | DataSourceError 分支 record_failure(),瞬态错误才允许新浪兜底 |
| 东财正常返回 | record_success() 重置计数与状态 |
与降级链的联动
# fetchers.py::_orchestrate_live_sources 的东财降级段
if (not result or not is_valid_data(data_type, result)) and budget.can_attempt():
if em_breaker.can_request():
budget.consume()
try:
em_result = await _fetch_from_eastmoney(data_type, **kwargs)
if em_result and is_valid_data(data_type, em_result):
result = em_result
source_used = "eastmoney"
em_breaker.record_success()
except RateLimitError as e:
em_breaker.record_failure()
logger.warning(f"东财限流/风控 ({data_type}): {e}")
if data_type == "fund_flow" and budget.consume():
sina_result = await sina_fund_flow()
if sina_result:
result = build_fund_flow(sina_result)
source_used = "sina_backup"
except DataSourceError as e:
em_breaker.record_failure()
if should_fallback(e) and data_type == "fund_flow" and budget.consume():
... # 新浪备用源
else:
logger.debug(f"东财熔断器 OPEN,跳过 ({data_type})")
if data_type == "fund_flow" and budget.consume():
... # 新浪备用源
熔断器与重试预算(budget.can_attempt(),上限 RETRY_BUDGET_MAX=6)双重保护:
熔断器控制"是否还能打东财",预算控制"单次请求最多打几次外部源"。熔断 OPEN 时
不经东财、直接落到新浪备用源或 DB 历史兜底。
前端可观测
熔断器自身的状态变化记 INFO 日志。core/resilience.py::CircuitBreaker 的文案为
熔断器[{name}] {old} → {new_state}(如 熔断器[tdx] CLOSED → OPEN),
并同步落一条降级事件到环形缓冲(见上节)。
盘中异动端点(/realtime/stock-changes)在 push2ex_breaker OPEN 或全部类型失败时,
返回最近一次成功的陈旧数据并标注 _stale: true(api/realtime.py,陈旧上限 24h),
而不是直接报错。
前端「数据可信度中心」(/data-trust)消费上述降级事件与/health/datasources 的熔断状态:
- 熔断状态 → 源状态徽章(三态
fresh/stale/missing,非颜色单独承载信息) - 降级事件 → 降级时间线(熔断器名经
normalizeBreakerName()归一后显示中文源名)
降级事件时间线
CircuitBreaker 每次状态跃迁都会落一条事件到内存环形缓冲,供
GET /health/degradation-events 与前端「数据可信度中心」(/data-trust) 的降级时间线排查。
这一节描述的是 core/resilience.py::CircuitBreaker;上面 SourceCircuitBreaker
那两套实现不落事件。
# core/resilience.py
_degradation_log = DegradationLog() # deque(maxlen=100),进程内、不落库
record_degradation(source, from_state, to_state, reason="")
get_degradation_events(limit=50) -> list[dict] # 时间逆序,limit 上限 100
事件结构(reason 截断到 200 字):
| 字段 | 说明 |
|---|---|
timestamp | ISO 秒精度时间戳 |
source | 熔断器名。裸名如 tdx / eastmoney;降级链实例形如 fallback:tencent_kline(前缀 fallback: 由 fallback_registry._BREAKER_KEY_PREFIX 拼接) |
from_state / to_state | 跃迁前后状态(CLOSED / OPEN / HALF_OPEN) |
reason | 人话原因,可能为空串 |
写入方:唯一出口在 _transition_to()
_transition_to() 是熔断器唯一的状态跃迁出口,状态变更与事件落库在此一并完成:
def _transition_to(self, new_state: str, reason: str = ""):
old = self._state
if old == new_state:
return # 同态跃迁不重复落事件
self._state = new_state
# ...计数重置...
logger.info(f"熔断器[{self.name}] {old} → {new_state}")
record_degradation(self.name, old, new_state, reason or _default_reason(old, new_state))
此前 _transition_to() 只写 logger.info、从不调用 record_degradation(),
而 record_degradation 全仓零调用方 —— 写入路径根本不存在。结果是
/health/degradation-events 永远返回空列表:接口 200、列表为空、控制台无任何报错,
属静默故障,前端时间线因此长期空白。现在跃迁必然落事件。
原因文案的两个来源
一、调用方显式给出 —— record_success() 传「半开探测成功(n/m 次)」;
record_failure(reason) 传 _brief(exc)。注意 safe_call_with_resilience 本来就捕获了
err_msg,此前并未传给熔断器,修复后已接通:
def _brief(exc: BaseException, limit: int = 160) -> str:
msg = " ".join(str(exc).split())
return f"{type(exc).__name__}: {msg}"[:limit] if msg else type(exc).__name__
压成单行的原因:降级事件会直接渲染到前端页面,多行 traceback 片段(如 httpx 异常) 会把时间线撑成一堆竖条。
二、_default_reason(from_state, to_state) 兜底 —— 关键在于
OPEN → HALF_OPEN 没有任何外部调用方:它由 _check_transition() 的定时检查
在读取 state / can_request() / get_stats() 时自动触发。若无兜底,这条跃迁的
reason 就是空串,时间线上会出现一条无原因的"神秘降级"。
| 跃迁 | 兜底文案 |
|---|---|
OPEN → HALF_OPEN | 熔断窗口结束,进入半开探测 |
HALF_OPEN → CLOSED | 半开探测通过,恢复正常 |
HALF_OPEN → OPEN | 半开探测失败,重新熔断 |
OPEN → CLOSED | 手动或自动重置,恢复正常 |
reset() 刻意不落事件
reset() 改为直接改状态 + 仅记日志,不写降级事件。原因:手动重置是运维动作,
而 reset_all_breakers() 一次能重置全部熔断器 —— 若每次都落事件,100 条容量的
环形缓冲会被瞬间冲光,真正的自动降级事件被顶出去,时间线显示一堆「手动恢复」,
而刚才到底哪个源降级了反而查不到。
状态重置与计数清零是无条件执行的:即便原状态已是 CLOSED 也要清。
reset 的语义是「清空所有运行态」,不只是「确保关闭」—— 若加
if old == CLOSED: return 提前返回会跳过计数清零。
前端消费
/data-trust 的降级时间线直接渲染这些事件。熔断器名在前端经
normalizeBreakerName() 归一:fallback: 前缀先剥掉再取前缀部分映射
(fallback:tencent_kline → tencent → 「腾讯行情」),否则用户看到的是一串内部标识。
设计原则
- 403 立即记失败:东财 IP 级风控信号,一次 403 即记一次失败(不重试)并同步触发该主机 30s 共享冷却;达到各自阈值(主行情 4、异动监控 2)后才真正 OPEN
- 冷却后探测:冷却期满允许一次请求试探是否解封,成功即整体恢复
- 不缓存空数据:熔断期间不将空结果写入缓存
- 参数分档:主行情 4/120s、异动监控 2/30s,按"数据可等性"决定激进度
- 日志可观测:熔断器状态变化记 INFO,避免在降级循环里刷 WARNING
- 跃迁必留痕:自动降级与自动恢复都落事件;手动重置不落,避免冲刷缓冲