本文最后更新于 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 超时时,如果没有断路器,每次请求都会:
- 等待超时:默认 30 秒连接超时,大量请求堆积
- 盲目重试:同一个已经挂掉的 API 被反复调用,加剧对方服务器压力
- 连锁故障:线程池被阻塞,新请求无法处理,整个服务假死
- 无降级路径:要么拿到数据,要么报错,没有中间态
更麻烦的是,外部数据源的故障往往是”半死不活”的——不是完全不可用,而是时好时坏。你的重试逻辑每一次都心存侥幸,结果就是系统在”尝试-失败-重试-再失败”的循环中空耗资源。
我们需要的是一种机制,能够记住哪些数据源已经不可靠了,暂时不再去碰它,过一段时间再试探性看看恢复了没有。
断路器模式:三态状态机
断路器的核心思想来自电气工程:当电流过载时,断路器自动跳闸,切断电路,防止设备烧毁。在软件中,它是一个三态状态机:
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 FALLBACK_CACHE_KEY: str = ""
def execute(self, db, **kwargs) -> R: cached = enhanced_cache.get(cache_key) if cached: return self.build_response(CACHE_HIT, cached)
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
if data_list is None: fallback = persistent_cache.get(self.FALLBACK_CACHE_KEY) if fallback: return self.build_response(STALE_CACHE, fallback)
return self.build_response(UNAVAILABLE, default_data)
|
以 ReShare 的北向资金策略为例,数据源优先级链是:
1 2 3 4
| class NorthFlowFetchStrategy(DataFetchStrategy): SOURCE_PRIORITY = ["akshare"] CACHE_TTL = 3600 FALLBACK_CACHE_KEY = "north_flow_fallback"
|
对于股票历史数据,优先级链更长:
1
| BaoStock(自有服务器)→ AKShare 新浪 → AKShare 东财 → YFinance(全球兜底)
|
第三层:双层缓存
即使所有数据源都熔断了,我们也不希望用户看到一片空白。双层缓存提供了最后的兜底:
| 缓存层 |
实现 |
TTL |
用途 |
| 热缓存 |
enhanced_cache(内存 LRU) |
1 小时 |
正常加速,减少外部调用 |
| 冷缓存 |
persistent_cache(JSON 文件) |
无 TTL,永久 |
灾难兜底,标记为 stale |
冷缓存的核心价值:即使数据过期了,有数据总比没数据好。响应中会标记 stale=true 和 last_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, )
|
第四层:优雅降级响应
当一切兜底都失败时,系统不会崩溃,而是返回一个结构化的降级响应:
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
| for attempt in range(self.MAX_RETRIES): try: df = ak.stock_zh_a_daily(symbol=akshare_symbol, ...) if df.empty: 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 ) 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)
|
重试策略也值得一提:指数退避 + 随机抖动(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):
|
踩过的坑
坑 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) 的抖动后,重试自然分散开来。
小结
断路器模式的本质是用空间换时间——用一小段”放弃”的时间,换取系统整体的稳定性。它的核心价值不在于让故障消失(那是不可能的),而在于:
- 快速失败:不浪费时间等一个注定超时的请求
- 自动隔离:把问题源从调度链中暂时摘除
- 渐进恢复:通过半开探测避免二次故障
- 多层兜底:优先级降级 → 缓存兜底 → 优雅降级
在 AI 数据管道的场景下,外部源的不可靠是常态而非异常。把断路器作为基础设施的标准组件,让每个数据源调用都”有退路可走”,是构建高可用数据管道的第一步。
ReShare 的四层防御体系——断路器隔离、优先级降级、缓存兜底、优雅响应——并不复杂,但它把”外部源抽风 → 系统瘫痪”的直线关系,变成了”外部源抽风 → 自动降级 → 用户几乎无感”的弹性链路。这中间的差距,就是凌晨三点你被电话叫醒和早上七点自然醒来的区别。