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

Axon Framework如何搭配reactive repository使用MessageDispatchInterceptor

问题核心原因

  • 你当前使用的MessageDispatchInterceptor是同步处理接口,内部构造的Mono如果没有被主动订阅,校验逻辑自然不会执行;而在Reactor原生非阻塞线程中调用block()会触发线程安全检查,直接抛出阻塞操作不支持的异常,这是响应式编程的基础约束。
  • Axon默认的命令派发响应式链路本身禁止任何阻塞操作,所以直接调用block()的方案完全不可行。

解决方案

方案1:使用Axon 4.6+ 提供的ReactiveMessageDispatchInterceptor(推荐)

这是Axon专门为响应式技术栈提供的拦截器组件,返回值直接为响应式发布器Publisher,Axon框架会自动完成订阅操作,不需要你手动处理阻塞逻辑。
实现代码示例:

class ReactiveSubnetCommandInterceptor : ReactiveMessageDispatchInterceptor<CommandMessage<*>> {

    @Autowired
    private lateinit var privateNetworkRepository: PrivateNetworkRepository

    override fun handle(message: CommandMessage<*>): Publisher<CommandMessage<*>> {
        if (message.payload !is CreateSubnetCommand) {
            return Mono.just(message)
        }
        val interceptCommand = message.payload as CreateSubnetCommand
        return privateNetworkRepository.findById(interceptCommand.privateNetworkId)
            // 替换为你自己的子网重叠校验逻辑
            .filter { network -> !network.isSubnetOverlap(interceptCommand.subnetCidr) }
            .switchIfEmpty(Mono.error(IllegalArgumentException("Requested subnet overlaps with an existing subnet.")))
            .thenReturn(message)
    }
}

之后将这个拦截器注册到响应式命令网关即可生效:

@Bean
fun reactiveCommandGateway(reactiveCommandBus: ReactiveCommandBus): ReactiveCommandGateway {
    return DefaultReactiveCommandGateway.builder()
        .commandBus(reactiveCommandBus)
        .dispatchInterceptors(ReactiveSubnetCommandInterceptor())
        .build()
}

方案2:Axon版本低于4.6的兼容降级方案

如果无法升级Axon版本,可以将阻塞操作放到专门的阻塞线程池执行,避免占用Reactor的非阻塞线程池:

class SubnetCommandInterceptor : MessageDispatchInterceptor<CommandMessage<*>> {

    @Autowired
    private lateinit var privateNetworkRepository: PrivateNetworkRepository

    // 定义专门处理阻塞操作的弹性线程池
    private val blockingScheduler = Schedulers.boundedElastic()

    override fun handle(messages: List<CommandMessage<*>?>): BiFunction<Int, CommandMessage<*>, CommandMessage<*>> {
        return BiFunction { _, command ->
            if (CreateSubnetCommand::class.simpleName == command.payloadType.simpleName){
                val interceptCommand = command.payload as CreateSubnetCommand
                privateNetworkRepository
                    .findById(interceptCommand.privateNetworkId)
                    .filter { network -> !network.isSubnetOverlap(interceptCommand.subnetCidr) }
                    .switchIfEmpty(Mono.error(IllegalArgumentException("Requested subnet overlaps with an existing subnet.")))
                    .subscribeOn(blockingScheduler)
                    .block()
            }
            command
        }
    }
}

注意:该方案是旧版本兼容方案,性能比原生响应式拦截器差,优先选择方案1。

额外提醒

这种前置的Set-based一致性校验只能降低冲突概率,无法100%保证并发场景下的一致性,最终必须在聚合内部再做一次校验兜底,避免并发写入导致的校验失效。

内容的提问来源于stack exchange,提问作者Patrick

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 07:24:05