错误发生后如何从Reactor Sink继续发送更多事件?
问题分析
你的核心问题在于:Reactor的Multicast Sink一旦调用emitError()进入错误状态后,就会永久终止,所有现有订阅者会收到错误,后续新订阅者也会直接触发该错误,无法继续接收新消息。加上Spring Redis消息监听器是多线程环境,单个消息处理的异常不应该影响整个发布流的可用性。
解决方案
1. 避免让Sink进入错误状态(核心修复)
不要在消息处理异常时调用sink.emitError(),而是捕获异常后仅记录日志,继续处理后续Redis消息。这样Sink会保持活跃状态,新老订阅者都能正常接收消息。
修改后的消息监听器代码:
@Override public void onMessage(Message message, byte[] pattern) { try { // 原有的消息处理逻辑 String body = new String(message.getBody()); // 处理消息发送失败的情况(比如背压) Sinks.EmitResult emitResult = sink.emitNext(body, Sinks.EmitFailureHandler.busyLooping(Duration.ofMillis(500))); if (emitResult.isFailure()) { // 记录发送失败的日志,比如背压导致的投递失败 log.warn("Failed to emit Redis message, result: {}", emitResult); } } catch (final Exception e) { // 仅记录异常,不终止Sink log.error("Error processing Redis message", e); } }
2. (可选)向订阅者传递错误但不终止流
如果需要让订阅者感知到消息处理错误,不要直接触发Sink的错误状态,而是将错误包装成普通事件发送。
步骤1:定义事件包装类
// 密封接口区分正常数据和错误事件 public sealed interface CacheUpdateEvent permits CacheUpdateData, CacheUpdateError { static CacheUpdateData data(String content) { return new CacheUpdateData(content); } static CacheUpdateError error(Throwable throwable) { return new CacheUpdateError(throwable); } } // 正常数据事件 record CacheUpdateData(String content) implements CacheUpdateEvent {} // 错误事件 record CacheUpdateError(Throwable throwable) implements CacheUpdateEvent {}
步骤2:修改Sink类型和消息处理逻辑
// 更新Sink的泛型为事件包装类 private final Sinks.Many<CacheUpdateEvent> sink = Sinks.many().multicast() .onBackpressureBuffer(Queues.SMALL_BUFFER_SIZE, false); @Override public void onMessage(Message message, byte[] pattern) { try { String body = new String(message.getBody()); Sinks.EmitResult result = sink.emitNext(CacheUpdateEvent.data(body), Sinks.EmitFailureHandler.busyLooping(Duration.ofMillis(500))); if (result.isFailure()) { log.warn("Failed to emit data event: {}", result); } } catch (final Exception e) { log.error("Error processing Redis message", e); // 发送错误事件,而非终止Sink Sinks.EmitResult errorResult = sink.emitNext(CacheUpdateEvent.error(e), Sinks.EmitFailureHandler.busyLooping(Duration.ofMillis(500))); if (errorResult.isFailure()) { log.warn("Failed to emit error event: {}", errorResult); } } }
步骤3:订阅者处理事件
redisPubSub.getPublisher() .bufferTimeout(MAX_BUFFER_SIZE, MAX_BUFFER_DURATION) .flatMap(events -> Flux.fromIterable(events) .handle((event, sink) -> { if (event instanceof CacheUpdateData data) { // 处理正常缓存更新数据 sink.next(data.content()); } else if (event instanceof CacheUpdateError error) { // 处理错误事件,比如记录日志,不终止订阅流 log.error("Received cache update error", error.throwable()); } }) ) // 订阅者自身逻辑出错时,保证流不终止 .onErrorContinue((e, obj) -> log.error("Error processing event", e)) .subscribe();
3. 优化订阅者发布逻辑
原getPublisher()方法中Flux.from(sink.asFlux())是冗余的,直接返回sink.asFlux()即可:
public Flux<CacheUpdateEvent> getPublisher() { return sink.asFlux().publishOn(Schedulers.single()); }
关键注意点
- Reactor的Multicast Sink是热流,一旦进入
error或complete状态就会永久终止,无法恢复。 - Spring Redis消息监听器是多线程环境,Sink的
emitNext()方法是线程安全的,无需额外同步。 - 背压处理:当前使用
onBackpressureBuffer(Queues.SMALL_BUFFER_SIZE, false),表示缓冲区满时新消息发送会失败而非丢弃旧消息,可根据业务需求调整参数。
内容的提问来源于stack exchange,提问作者sunalive
相关产品推荐
相关产品推荐

