跳到主要内容

熔断器机制

概述​

数据源配备熔断器(Circuit Breaker),连续失败后自动跳过请求,避免无效重试加重封禁。 代码中存在两套实现、三个实例,面向不同调用面,互不影响:

实现位置实例参数调用面
SourceCircuitBreakermarket_data_service/infra.pyem_breaker阈值 4 / 冷却 120s(EM_BREAKER_FAILURE_THRESHOLD / EM_BREAKER_COOLDOWN_SEC)get_market_data() 东财降级
SourceCircuitBreakerservices/data_sources/stock_changes.pypush2ex_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()(风控信号,不重试)
东财返回 429RateLimitError 分支 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 字):

字段说明
timestampISO 秒精度时间戳
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))
2026-10-06 修复:事件缓冲恒为空

此前 _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
  • 跃迁必留痕:自动降级与自动恢复都落事件;手动重置不落,避免冲刷缓冲

相关文档​