Reactor 错误处理与背压
2026/7/16大约 3 分钟
Reactor 错误处理与背压
一句话概括
响应式流中错误是终止事件——发生错误后流会停止发射。Reactor 提供了声明式的操作符来优雅容错、重试恢复,以及通过背压让下游控制上游的发射速率。
1. 错误处理策略
在响应式编程中,try-catch 不再适用。Reactor 用操作符代替了异常捕获:
onErrorReturn — 返回默认值
Mono.<String>error(new RuntimeException("oops"))
.onErrorReturn("fallback");
// → "fallback"
// 按异常类型区分
Flux.just("valid", "invalid")
.map(this::parse)
.onErrorReturn(IllegalArgumentException.class, "default");onErrorResume — 切换到备用流
Mono.<String>error(new RuntimeException("failed"))
.onErrorResume(e -> Mono.just("resumed: " + e.getMessage()));
// → "resumed: failed"
// 实际场景:调用远程 API 失败后走缓存
userService.findById(id)
.onErrorResume(e -> localCache.get(id));retry — 重试整个流
Flux.range(1, 3).map(i -> {
if (i == 2) throw new RuntimeException("retry me");
return i;
}).retry(2);
// 重试 2 次,总共最多执行 3 次重试的代价:retry() 会重新订阅整个源,包括之前已成功的元素。如果不想重发成功数据,考虑用 retryWhen + 状态管理。
timeout — 超时降级
Flux.interval(Duration.ofSeconds(10))
.timeout(Duration.ofMillis(100))
.onErrorResume(e -> Flux.just(-1L));
// 100ms 内没收到数据 → 异常 → 降级返回 -1实际场景:调用外部 API 设置超时阈值,超时后走缓存或默认值。
2. 错误处理组合模式
// 超时 + 重试 + 降级 = 三层防护
userService.findById(id)
.timeout(Duration.ofMillis(500))
.retry(2)
.onErrorResume(e -> localCache.get(id));这是一个典型的"快速失败 + 自动恢复 + 兜底降级"模式,在微服务调用中非常实用。
3. 背压(Backpressure)
背压是响应式编程区别于传统异步编程的核心特性——下游通过 request(n) 告诉上游"我处理不过来,慢点发"。
为什么需要背压?
数据源(每秒 1000 条)──→ 处理逻辑(每秒 10 条,但无背压)
↓
内存暴涨 → OOM没有背压的异步场景,生产速度远大于消费速度时,内存就会无限制堆积。
背压策略
| 策略 | 行为 | 适用场景 |
|---|---|---|
onBackpressureBuffer(n) | 缓冲最多 n 个元素,超出抛异常 | 下游偶尔慢,能接受等待 |
onBackpressureDrop(callback) | 丢弃来不及处理的数据 | 只关心最新值 |
onBackpressureLatest() | 保留最新元素,丢弃旧的 | 实时行情 |
// 策略一:缓冲
Flux.range(1, 1000)
.onBackpressureBuffer(100); // 最多缓冲 100 个
// 策略二:丢弃并记录
Flux.range(1, 1000)
.onBackpressureDrop(i -> log.warn("Dropped: {}", i));
// 策略三:保留最新
Flux.interval(Duration.ofMillis(10))
.onBackpressureLatest(); // 只保留最新值背压工作流程
Publisher Subscriber
│ │
│── onSubscribe(Subscription) ──▶│
│◀── request(n) ───────────│
│── onNext(item) ───────────▶│
│── onNext(item) ───────────▶│ (最多 n 个)
│◀── request(m) ───────────│ (处理完再要)订阅者通过 Subscription.request(n) 控制窗口大小,发布者严格按需发射。
4. 常见陷阱
❌ catch 内部异常而不传播
// 错误写法
Flux.just("1", "abc", "3")
.map(s -> {
try {
return Integer.parseInt(s);
} catch (NumberFormatException e) {
return null; // ❌ 返回 null 会触发空指针处理
}
});
// 正确做法:用操作符处理
Flux.just("1", "abc", "3")
.map(Integer::parseInt)
.onErrorContinue((err, item) -> log.warn("skip bad: {}", item));❌ 无界缓冲
Flux.interval(Duration.ofMillis(1))
.subscribe(this::slowConsumer);
// ❌ 生产 1ms/条,消费 100ms/条 → OOM
// ✅ 加上背压限制
Flux.interval(Duration.ofMillis(1))
.onBackpressureDrop()
.subscribe(this::slowConsumer);❌ retry 滥用
.retry(Long.MAX_VALUE) // ❌ 无限重试可能导致 DoS
.retry(3) // ✅ 限制重试次数关键总结
| 错误处理 | 适用场景 |
|---|---|
onErrorReturn | 快速降级,返回默认值 |
onErrorResume | 切换到备用流(缓存/备选 API) |
retry | 临时故障自动重试 |
timeout | 调用超时保护 |
| 背压策略 | 适用场景 |
|---|---|
buffer | 偶尔慢,接受延迟 |
drop | 只关心最新 |
latest | 实时行情,最新价 |
上一篇:Reactor 操作符详解
下一篇:Reactor 调度器与实战模式