断路器模式:让 AI 数据管道不怕外部源抽风

本文最后更新于 2026年7月22日 上午

当数据源开始”抽风”

凌晨三点,你的量化交易系统准时启动。一切看起来很正常——直到日志里开始刷屏:

1
2
3
4
5
WARNING - akshare 新浪渠道获取数据失败: ConnectionError
WARNING - 重试 1/3...
WARNING - 重试 2/3...
WARNING - 重试 3/3...
ERROR - 所有数据源不可用

北向资金数据没拉到,龙虎榜数据没拉到,融资融券数据也没拉到。等你在早上七点打开终端,发现整个情绪面分析模块已经瘫痪了四个小时。

这不是假设场景。在我们的 ReShare 量化数据中台项目中,外部数据源的不可靠性是每天都要面对的现实:新浪财经的 API 间歇性超时,东方财富的反爬策略随时升级,BaoStock 的自有服务器偶尔也会宕机。一个”抽风”的外部源,如果没有保护机制,可以像传染病一样拖垮整个数据管道。

这篇文章分享我们在 ReShare 项目中实践的断路器模式(Circuit Breaker Pattern)——如何让数据管道在外部源”抽风”时自动降级、自动隔离、自动恢复,而不是一路报错到底。

问题:没有断路器的世界

先看看没有保护机制时会发生什么。一个典型的数据获取调用链:

1
用户请求 → API 路由 → MarketDataService → AKShare → 新浪财经 API

当新浪财经 API 超时时,如果没有断路器,每次请求都会:

  1. 等待超时:默认 30 秒连接超时,大量请求堆积
  2. 盲目重试:同一个已经挂掉的 API 被反复调用,加剧对方服务器压力
  3. 连锁故障:线程池被阻塞,新请求无法处理,整个服务假死
  4. 无降级路径:要么拿到数据,要么报错,没有中间态

更麻烦的是,外部数据源的故障往往是”半死不活”的——不是完全不可用,而是时好时坏。你的重试逻辑每一次都心存侥幸,结果就是系统在”尝试-失败-重试-再失败”的循环中空耗资源。

我们需要的是一种机制,能够记住哪些数据源已经不可靠了,暂时不再去碰它,过一段时间再试探性看看恢复了没有

断路器模式:三态状态机

断路器的核心思想来自电气工程:当电流过载时,断路器自动跳闸,切断电路,防止设备烧毁。在软件中,它是一个三态状态机:

1
2
3
4
5
CLOSED(正常)─── 连续失败 N 次 ───→ OPEN(熔断)
↑ │
│ 等待 T
│ ↓
└── 连续成功 M 次 ── HALF_OPEN(半开)
状态 行为 迁移条件
CLOSED(关闭) 正常放行所有请求 连续失败达到阈值 → OPEN
OPEN(打开) 直接拒绝所有请求 超时时间到 → HALF_OPEN
HALF_OPEN(半开) 放行少量探测请求 成功达标 → CLOSED;任意失败 → OPEN

关键设计点:

  • CLOSED 到 OPEN:不是一次失败就熔断,而是连续失败达到阈值。偶发错误不会触发熔断。
  • OPEN 到 HALF_OPEN:熔断不是永久的。超时后进入半开状态,允许少量探测请求验证是否恢复。
  • HALF_OPEN 到 CLOSED:需要连续成功达到阈值才能完全恢复。这防止了”假恢复”——偶尔成功一次但整体仍不稳定。

ReShare 的四层防御体系

在实际项目中,单靠断路器是不够的。ReShare 构建了四层防御体系,从内到外逐层兜底:

1
2
3
4
5
6
7
8
9
10
11
Layer 1: 断路器(CircuitBreaker)
── 按数据源隔离,连续失败 5 次熔断,30 秒后半开探测
↓ 熔断器打开则跳过
Layer 2: 优先级降级(DataFetchStrategy)
── 按优先级遍历数据源,A 挂了试 B,B 挂了试 C
↓ 所有源都失败
Layer 3: 缓存兜底(EnhancedCache + PersistentCache)
── 热缓存 TTL 1 小时,冷缓存持久化到磁盘
↓ 连缓存都没有
Layer 4: 优雅降级响应
── 返回空数据 + status="unavailable" + last_success 时间戳

第一层:断路器

这是最内层的防线,实现在 data_source_exceptions.py 中:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
class CircuitState(Enum):
CLOSED = "closed" # 正常状态,允许请求
OPEN = "open" # 熔断状态,拒绝请求
HALF_OPEN = "half_open" # 半开状态,允许探测请求

@dataclass
class CircuitBreakerConfig:
failure_threshold: int = 5 # 触发熔断的失败次数
success_threshold: int = 3 # 恢复所需的连续成功次数
timeout_seconds: int = 30 # 熔断超时时间
half_open_max_calls: int = 3 # 半开状态最大探测次数

class CircuitBreaker:
def can_execute(self) -> bool:
with self._lock:
if self._stats.state == CircuitState.CLOSED:
return True
elif self._stats.state == CircuitState.OPEN:
elapsed = (now - self._stats.last_failure_time).total_seconds()
if elapsed >= self.config.timeout_seconds:
self._update_state(CircuitState.HALF_OPEN)
return True # 放行探测请求
return False # 熔断中,直接拒绝
elif self._stats.state == CircuitState.HALF_OPEN:
if self._stats.half_open_calls < self.config.half_open_max_calls:
self._stats.half_open_calls += 1
return True # 限量探测
return False

def record_failure(self, exception=None):
with self._lock:
self._stats.failure_count += 1
if self._stats.failure_count >= self.config.failure_threshold:
self._update_state(CircuitState.OPEN)

def record_success(self):
with self._lock:
self._stats.failure_count = 0
if self._stats.state == CircuitState.HALF_OPEN:
self._stats.success_count += 1
if self._stats.success_count >= self.config.success_threshold:
self._update_state(CircuitState.CLOSED)

注意几个工程细节:

  • 线程安全:所有状态变更都在 threading.Lock 保护下进行。断路器是多线程共享的,不加锁会导致状态混乱。
  • 原子状态迁移_update_state 方法统一处理状态变更和日志记录,避免散落在各处的状态修改。
  • 半开限流:HALF_OPEN 状态下不仅放行探测请求,还限制探测次数(half_open_max_calls),防止探测请求洪流。

每个数据源有自己的断路器实例,由 CircuitBreakerManager 统一管理:

1
2
3
4
5
6
7
8
9
10
class CircuitBreakerManager:
def __init__(self):
self._circuit_breakers: dict[str, CircuitBreaker] = {}
self._lock = threading.Lock()

def get_circuit_breaker(self, source_name: str) -> CircuitBreaker:
with self._lock:
if source_name not in self._circuit_breakers:
self._circuit_breakers[source_name] = CircuitBreaker(source_name)
return self._circuit_breakers[source_name]

第二层:策略模式 + 优先级降级

断路器解决的是”不要再调用已经挂掉的源”的问题,但请求还是要拿到数据的。策略模式负责在多个源之间按优先级降级:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
class DataFetchStrategy(ABC, Generic[T, R]):
SOURCE_PRIORITY: list[str] = [] # 数据源优先级列表
CACHE_TTL: int = 3600 # 缓存 TTL
FALLBACK_CACHE_KEY: str = "" # 持久化缓存键

def execute(self, db, **kwargs) -> R:
# 1. 检查热缓存
cached = enhanced_cache.get(cache_key)
if cached:
return self.build_response(CACHE_HIT, cached)

# 2. 按优先级遍历数据源
for source_name in self.SOURCE_PRIORITY:
breaker = self._circuit_breaker_manager.get_circuit_breaker(source_name)

if not breaker.can_execute():
logger.warning(f"数据源 {source_name} 熔断器打开,跳过")
continue

try:
raw_data = self.fetch_from_source(source_name, **kwargs)
breaker.record_success()
break
except Exception as e:
breaker.record_failure(e)
continue # 试下一个源

# 3. 持久化缓存兜底
if data_list is None:
fallback = persistent_cache.get(self.FALLBACK_CACHE_KEY)
if fallback:
return self.build_response(STALE_CACHE, fallback)

# 4. 优雅降级
return self.build_response(UNAVAILABLE, default_data)

以 ReShare 的北向资金策略为例,数据源优先级链是:

1
2
3
4
class NorthFlowFetchStrategy(DataFetchStrategy):
SOURCE_PRIORITY = ["akshare"] # 目前只有一个源
CACHE_TTL = 3600 # 1 小时热缓存
FALLBACK_CACHE_KEY = "north_flow_fallback" # 持久化兜底

对于股票历史数据,优先级链更长:

1
BaoStock(自有服务器)→ AKShare 新浪 → AKShare 东财 → YFinance(全球兜底)

第三层:双层缓存

即使所有数据源都熔断了,我们也不希望用户看到一片空白。双层缓存提供了最后的兜底:

缓存层 实现 TTL 用途
热缓存 enhanced_cache(内存 LRU) 1 小时 正常加速,减少外部调用
冷缓存 persistent_cache(JSON 文件) 无 TTL,永久 灾难兜底,标记为 stale

冷缓存的核心价值:即使数据过期了,有数据总比没数据好。响应中会标记 stale=truelast_success 时间戳,让调用方知道这是旧数据:

1
2
3
4
5
6
7
result = DataFetchResult(
data=cached_data_list,
status=DataFetchStatus.STALE_CACHE,
source="fallback_cache",
stale=True,
last_success=last_success, # "2026-07-21T15:30:00"
)

第四层:优雅降级响应

当一切兜底都失败时,系统不会崩溃,而是返回一个结构化的降级响应:

1
2
3
4
5
6
7
8
9
{
"data": null,
"status": "unavailable",
"source": "none",
"from_cache": false,
"stale": false,
"last_success": "2026-07-21T15:30:00",
"error_message": "所有数据源不可用"
}

调用方可以根据 status 字段决定如何处理:是展示旧数据、显示”数据暂时不可用”提示,还是触发告警。

数据源内部的双重保险

断路器和优先级降级是数据源之间的保护。在数据源内部,ReShare 还有另一层保险——以 AKShare 新浪渠道为例:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
# akshare_source.py
for attempt in range(self.MAX_RETRIES):
try:
df = ak.stock_zh_a_daily(symbol=akshare_symbol, ...)
if df.empty:
# AKShare 返回空 → 直接调用新浪 HTTP API 降级
df = self._sina_http_fallback(symbol, start_str, end_str)
return df
except Exception as e:
if attempt < self.MAX_RETRIES - 1:
delay = min(
self.RETRY_DELAY * (attempt + 1) + random.uniform(0, 1),
20 # 最大 20 秒
)
time.sleep(delay)
else:
# 最后一次重试失败 → 再次降级
df = self._sina_http_fallback(symbol, start_str, end_str)

_sina_http_fallback 绕过 AKShare 库,直接调用新浪财经的 HTTP API。这是一种”降级中的降级”——当库级别的调用失败时,退化到裸 HTTP 请求:

1
2
3
4
5
6
7
def _sina_http_fallback(self, symbol, start_str, end_str):
sina_url = (
f"http://money.finance.sina.com.cn/quotes_service/api/json_v2.php/"
f"CN_MarketData.getKLineData?symbol={prefix}{pure}&scale=240&datalen=1023"
)
resp = requests.get(sina_url, timeout=10)
# 解析 JSON → DataFrame

重试策略也值得一提:指数退避 + 随机抖动(jitter),避免多个客户端同步重试造成的”惊群效应”:

1
2
3
4
delay = min(
self.RETRY_DELAY * (attempt + 1) + random.uniform(0, 1),
20
)

可观测性:光有断路器不够

断路器做了自动决策,但运维人员需要知道发生了什么。ReShare 在 MarketDataService 中暴露了断路器状态查询和重置接口:

1
2
3
4
5
6
7
8
class MarketDataService:
def get_circuit_breaker_status(self) -> dict[str, dict]:
"""获取所有数据源的熔断器状态"""
return data_source_manager.get_circuit_breaker_stats()

def reset_circuit_breaker(self, source_name: str) -> None:
"""手动重置指定数据源的熔断器"""
data_source_manager.reset_circuit_breaker(source_name)

调用 get_circuit_breaker_status() 返回的状态示例:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
{
"akshare": {
"state": "open",
"failure_count": 5,
"total_calls": 142,
"total_failures": 7,
"last_failure_time": "2026-07-22T03:15:22",
"last_state_change_time": "2026-07-22T03:15:22"
},
"baostock": {
"state": "closed",
"failure_count": 0,
"total_calls": 98,
"total_failures": 1
}
}

同时,DataFetchStrategy 记录了每次数据获取的尝试历史(DataSourceAttempt),包括每个源的尝试结果、响应时间、错误信息,方便事后排查。

健康检测器(DataSourceHealthChecker)定期后台巡检,自动禁用连续失败的数据源,并在 5 分钟后自动重试——这和断路器的超时恢复机制形成了双保险:

1
2
3
4
class DataSourceHealthChecker:
def __init__(self, check_interval=300, max_consecutive_failures=3):
# 每 5 分钟定时检测不可用的数据源
# 连续失败 3 次标记为 UNAVAILABLE

踩过的坑

坑 1:阈值设太敏感

最初 failure_threshold=3,结果偶发的网络抖动就会触发熔断。高峰期数据源响应慢一点,连续三次超时就熔断了,但第四次其实就能成功。后来调到 5 次就好多了。

教训:阈值要根据数据源的稳定性来调。自有服务可以严格,公共 API 要宽松。

坑 2:半开状态不限制探测量

最初半开状态下没有限制探测请求数量,结果大量请求涌入刚恢复的数据源,又把它打挂了。加上 half_open_max_calls=3 后,最多放行 3 个探测请求,其他请求走降级路径。

坑 3:冷缓存没标记 stale

最初冷缓存返回的数据和正常数据格式一模一样,调用方分不清是实时数据还是旧数据。后来在响应中加了 stale 标志和 last_success 时间戳,让前端可以显示”数据更新于 2 小时前”。

坑 4:重试不加 jitter

多个定时任务同时启动,同时发现数据源失败,同时按固定间隔重试,形成了周期性的请求洪峰。加上 random.uniform(0, 1) 的抖动后,重试自然分散开来。

小结

断路器模式的本质是用空间换时间——用一小段”放弃”的时间,换取系统整体的稳定性。它的核心价值不在于让故障消失(那是不可能的),而在于:

  1. 快速失败:不浪费时间等一个注定超时的请求
  2. 自动隔离:把问题源从调度链中暂时摘除
  3. 渐进恢复:通过半开探测避免二次故障
  4. 多层兜底:优先级降级 → 缓存兜底 → 优雅降级

在 AI 数据管道的场景下,外部源的不可靠是常态而非异常。把断路器作为基础设施的标准组件,让每个数据源调用都”有退路可走”,是构建高可用数据管道的第一步。

ReShare 的四层防御体系——断路器隔离、优先级降级、缓存兜底、优雅响应——并不复杂,但它把”外部源抽风 → 系统瘫痪”的直线关系,变成了”外部源抽风 → 自动降级 → 用户几乎无感”的弹性链路。这中间的差距,就是凌晨三点你被电话叫醒和早上七点自然醒来的区别。


断路器模式:让 AI 数据管道不怕外部源抽风
https://normdist.com/2026/07/22/ND-20260722-001-circuit-breaker-for-data-pipeline/
作者
小瑞
发布于
2026年7月22日
许可协议