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

Spring Data R2DBC自定义@Query方法传Flux参数报错求助

Spring Data R2DBC自定义仓库方法接收Flux参数报错问题

我在探索Spring Data R2DBC时,尝试在继承ReactiveCrudRepository的自定义仓库接口中实现一个接收Flux<Person>参数、返回Flux<Person>的方法,期望它具备和findAllById(Publisher idStream)相同的行为。但使用@Query注解编写自定义方法时,触发了“Value must not be null”异常,想咨询该用法是否被支持。

自定义仓库接口代码

public interface PersonRepository extends ReactiveCrudRepository<Person, Person> {

    @Query("SELECT id, name  FROM Person WHERE id = ?")
    Flux<Person> myMethod(Flux<Person> person);

}

实体类代码

@Getter
@Setter
@ToString
@Data
@NoArgsConstructor
@AllArgsConstructor
@Table(value = "Person", schema = "mySchema")
public class Person {
    @Id
    @Column("id")
    Long id;

    @Column("name")
    String name;

}

数据库配置类代码

@Configuration
@EnableR2dbcRepositories
@EnableR2dbcAuditing
@EnableTransactionManagement

public class DatabaseConfiguration extends AbstractR2dbcConfiguration {

......

}

主应用代码

@SpringBootApplication
@EnableR2dbcAuditing
@EnableConfigurationProperties({ApplicationProperties.class})
public class MyApplication {

    private static final Logger logger = LoggerFactory.getLogger(MyApplication.class);


    @Bean
    public CommandLineRunner consume(PersonRepository personRepository) {


        LongSupplier randomLong = () -> {
            return RandomUtils.nextLong(10L,20L);
        };

        Flux<Long> personIds = Flux.fromStream(LongStream.generate(randomLong).boxed());
        Flux<Person> persons = personIds.map( p -> {
            Person person = new Person();
            person.setId(p);
            return person;
        });


        personRepository.myMethod(persons.take(3)).subscribe( id -> logger.info("Processed Person Id: " + id));
        
        return null;
    }   

    .......
}

异常信息

reactor.core.Exceptions$ErrorCallbackNotImplemented: java.lang.IllegalArgumentException: Value must not be null
Caused by: java.lang.IllegalArgumentException: **Value must not be null**
    at org.springframework.util.Assert.notNull(Assert.java:201)
    Suppressed: reactor.core.publisher.FluxOnAssembly$OnAssemblyException: 
Assembly trace from producer [reactor.core.publisher.MonoMapFuseable] :
    reactor.core.publisher.Mono.map(Mono.java:3411)
    org.springframework.data.r2dbc.repository.query.StringBasedR2dbcQuery.createQuery(StringBasedR2dbcQuery.java:156)
Error has been observed at the following site(s):
    *__________Mono.map ⇢ at org.springframework.data.r2dbc.repository.query.StringBasedR2dbcQuery.createQuery(StringBasedR2dbcQuery.java:156)
    |_ Mono.flatMapMany ⇢ at org.springframework.data.r2dbc.repository.query.AbstractR2dbcQuery.execute(AbstractR2dbcQuery.java:88)
    *____Flux.usingWhen ⇢ at org.springframework.data.repository.core.support.RepositoryMethodInvoker$ReactiveInvocationListenerDecorator.decorate(RepositoryMethodInvoker.java:242)
Original Stack Trace:
        at org.springframework.util.Assert.notNull(Assert.java:201)
        at org.springframework.r2dbc.core.Parameter.from(Parameter.java:54)
        at org.springframework.data.r2dbc.query.QueryMapper.getBindValue(QueryMapper.java:426)
        at org.springframework.data.r2dbc.core.DefaultReactiveDataAccessStrategy.getBindValue(DefaultReactiveDataAccessStrategy.java:288)
        at org.springframework.data.r2dbc.repository.query.ExpressionEvaluatingParameterBinder.getBindValue(ExpressionEvaluatingParameterBinder.java:134)
        at org.springframework.data.r2dbc.repository.query.ExpressionEvaluatingParameterBinder.bindExpressions(ExpressionEvaluatingParameterBinder.java:82)
        at org.springframework.data.r2dbc.repository.query.ExpressionEvaluatingParameterBinder.bind(ExpressionEvaluatingParameterBinder.java:72)
        at org.springframework.data.r2dbc.repository.query.StringBasedR2dbcQuery$ExpandedQuery.<init>(StringBasedR2dbcQuery.java:197)
        at org.springframework.data.r2dbc.repository.query.StringBasedR2dbcQuery.lambda$createQuery$0(StringBasedR2dbcQuery.java:156)
        at reactor.core.publisher.FluxMapFuseable$MapFuseableSubscriber.onNext(FluxMapFuseable.java:113)
        at reactor.core.publisher.FluxDefaultIfEmpty$DefaultIfEmptySubscriber.onNext(FluxDefaultIfEmpty.java:101)
        at reactor.core.publisher.FluxMapFuseable$MapFuseableSubscriber.onNext(FluxMapFuseable.java:129)
        at reactor.core.publisher.FluxMapFuseable$MapFuseableSubscriber.onNext(FluxMapFuseable.java:129)
        at reactor.core.publisher.FluxPeekFuseable$PeekFuseableSubscriber.onNext(FluxPeekFuseable.java:210)
        at reactor.core.publisher.FluxMapFuseable$MapFuseableSubscriber.onNext(FluxMapFuseable.java:129)
        at reactor.core.publisher.Operators$MonoSubscriber.complete(Operators.java:1816)
        at reactor.core.publisher.MonoCollectList$MonoCollectListSubscriber.onComplete(MonoCollectList.java:128)
        at reactor.core.publisher.FluxConcatMap$ConcatMapImmediate.drain(FluxConcatMap.java:368)
        at reactor.core.publisher.FluxConcatMap$ConcatMapImmediate.onComplete(FluxConcatMap.java:276)
        at reactor.core.publisher.Operators.complete(Operators.java:137)
        at reactor.core.publisher.FluxIterable.subscribe(FluxIterable.java:148)
        at reactor.core.publisher.FluxIterable.subscribe(FluxIterable.java:87)
        at reactor.core.publisher.Flux.subscribe(Flux.java:8469)
        at reactor.core.publisher.FluxUsingWhen.subscribe(FluxUsingWhen.java:94)
        at reactor.core.publisher.Flux.subscribe(Flux.java:8469)
        at reactor.core.publisher.Flux.subscribeWith(Flux.java:8642)
        at reactor.core.publisher.Flux.subscribe(Flux.java:8439)
        at reactor.core.publisher.Flux.subscribe(Flux.java:8363)
        at reactor.core.publisher.Flux.subscribe(Flux.java:8306)

问题原因及解决方案

原因分析

你当前的写法不被Spring Data R2DBC支持:

  1. @Query中的占位符?无法直接绑定Flux<Person>类型参数,框架会尝试将整个Flux对象作为参数绑定,而非提取每个Person的id,最终触发空值断言失败。
  2. findAllById是框架内置方法,会自动处理Publisher类型的ID流,逐个提取ID执行查询,但自定义@Query方法不会自动做这个转换。

正确实现方式

方案一:直接复用findAllById

不需要自定义方法,将Flux<Person>转换为Flux<Long>(提取id)后调用内置方法即可:

// 替换主应用中的调用代码
personRepository.findAllById(persons.take(3).map(Person::getId))
               .subscribe(person -> logger.info("Processed Person Id: " + person.getId()));

方案二:修正自定义方法

如果一定要自定义方法,需要明确指定要绑定的实体属性,或者直接接收ID流:

  • 方式1:绑定实体的id属性
// 修改仓库接口方法
@Query("SELECT id, name FROM Person WHERE id = :person.id")
Flux<Person> myMethod(@Param("person") Flux<Person> personFlux);
  • 方式2:直接接收ID流
// 修改仓库接口方法
@Query("SELECT id, name FROM Person WHERE id = ?")
Flux<Person> myMethod(Flux<Long> personIds);

调用时传入提取后的ID流即可,和方案一的调用逻辑一致。

说明

Spring Data R2DBC的自定义@Query方法支持绑定Publisher类型参数,但必须明确指定要提取的属性(如:person.id),或者直接传入ID流。直接传入实体对象的Flux时,框架无法自动解析出要绑定的字段,从而引发空值异常。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 07:01:55