基于Reactor的数据流管道错误分流处理方案咨询
解决方案思路
一、Reactor原生错误分流(无需Either)
不用返回Either类型,直接在流中捕获错误并分流到单独处理流,让主业务流只处理正常的Message<MyPojo>:
- 定义错误处理Sink
先创建一个Sink专门接收错误消息,集中处理日志记录、死信写入、偏移量提交逻辑:
Sink<Message<Throwable>> errorSink = Sinks.many().multicast().onBackpressureBuffer(); Flux<Message<Throwable>> errorFlux = errorSink.asFlux(); // 启动错误处理流 errorFlux.subscribe(errorMsg -> { // 记录错误日志 log.error("事件处理失败", errorMsg.getPayload()); // 从Header获取偏移量并提交 ReceiverOffset offset = errorMsg.getHeaders().get("receiverOffset", ReceiverOffset.class); if (offset != null) { offset.acknowledge(); } // 写入死信主题 kafkaSender.send(MessageBuilder.withPayload(errorMsg.getPayload()) .setHeader(KafkaHeaders.TOPIC, "dead-letter-topic") .build()) .subscribe(); });
- 业务流中捕获错误并分流
在每个处理阶段用onErrorResume捕获异常,将包含错误和原消息Header的Message发送到错误Sink,主流程返回空跳过错误消息:
Flux<Message<MyPojo>> businessFlow = kafkaReceiver.receive() .map(receiverRecord -> { // 解析消息并将偏移量存入Header MyPojo payload = parsePayload(receiverRecord.value()); return MessageBuilder.withPayload(payload) .setHeader("receiverOffset", receiverRecord.receiverOffset()) .build(); }) .flatMap(msg -> processStep1(msg) .onErrorResume(e -> { Message<Throwable> errorMsg = MessageBuilder.withPayload(e) .copyHeaders(msg.getHeaders()) .build(); errorSink.tryEmitNext(errorMsg); return Mono.empty(); })) .flatMap(msg -> processStep2(msg) .onErrorResume(e -> { Message<Throwable> errorMsg = MessageBuilder.withPayload(e) .copyHeaders(msg.getHeaders()) .build(); errorSink.tryEmitNext(errorMsg); return Mono.empty(); })) // 正常处理完成后提交偏移量 .doOnNext(msg -> { ReceiverOffset offset = msg.getHeaders().get("receiverOffset", ReceiverOffset.class); if (offset != null) { offset.acknowledge(); } }); // 启动业务流 businessFlow.subscribe();
这种方式下主业务流无需处理错误分支,所有错误自动分流到错误处理逻辑。
二、结合Spring Integration错误通道
如果已经使用Spring Integration,可以将Reactor流的错误消息发送到SI的errorChannel,复用其错误处理生态:
- 注入错误通道并发送错误消息
@Autowired private MessageChannel errorChannel; // 在业务流的onErrorResume中发送错误 .onErrorResume(e -> { Message<Throwable> errorMsg = MessageBuilder.withPayload(e) .copyHeaders(msg.getHeaders()) .build(); errorChannel.send(errorMsg); return Mono.empty(); })
- 配置SI错误处理逻辑
用@ServiceActivator注解实现错误通道的处理逻辑:
@ServiceActivator(inputChannel = "errorChannel") public void handleError(ErrorMessage errorMsg) { log.error("事件处理失败", errorMsg.getPayload()); // 获取原消息及偏移量 Message<?> originalMsg = errorMsg.getOriginalMessage(); if (originalMsg != null) { ReceiverOffset offset = originalMsg.getHeaders().get("receiverOffset", ReceiverOffset.class); if (offset != null) { offset.acknowledge(); } // 写入死信主题 kafkaSender.send(MessageBuilder.withPayload(originalMsg.getPayload()) .setHeader(KafkaHeaders.TOPIC, "dead-letter-topic") .build()) .subscribe(); } }
这种方式可以复用SI的重试、死信队列等成熟机制。
三、进阶:重试+错误分流组合
如果需要对可重试异常先进行重试,再分流最终失败的消息:
.flatMap(msg -> processStep1(msg) .retryWhen(Retry.backoff(3, Duration.ofSeconds(1)) .filter(e -> e instanceof RetryableException)) // 仅重试指定异常 .onErrorResume(e -> { Message<Throwable> errorMsg = MessageBuilder.withPayload(e) .copyHeaders(msg.getHeaders()) .build(); errorSink.tryEmitNext(errorMsg); return Mono.empty(); }))
关键注意事项
- 偏移量提交:无论处理成功或失败,都要确保
ReceiverOffset.acknowledge()被调用,结合乱序提交特性,单个偏移量提交不会影响其他消息的处理。 - Header保留:错误处理时必须复制原消息的Header,保证能获取到
ReceiverOffset。 - 背压处理:错误Sink选择
multicast().onBackpressureBuffer()类型,避免错误处理速度慢导致主业务流阻塞。
内容的提问来源于stack exchange,提问作者Ben
相关产品推荐
相关产品推荐

