You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于Reactor的数据流管道错误分流处理方案咨询

解决方案思路

一、Reactor原生错误分流(无需Either)

不用返回Either类型,直接在流中捕获错误并分流到单独处理流,让主业务流只处理正常的Message<MyPojo>:

  1. 定义错误处理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();
});
  1. 业务流中捕获错误并分流
    在每个处理阶段用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,复用其错误处理生态:

  1. 注入错误通道并发送错误消息
@Autowired
private MessageChannel errorChannel;

// 在业务流的onErrorResume中发送错误
.onErrorResume(e -> {
    Message<Throwable> errorMsg = MessageBuilder.withPayload(e)
            .copyHeaders(msg.getHeaders())
            .build();
    errorChannel.send(errorMsg);
    return Mono.empty();
})
  1. 配置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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.24 16:18:34