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支持:
@Query中的占位符?无法直接绑定Flux<Person>类型参数,框架会尝试将整个Flux对象作为参数绑定,而非提取每个Person的id,最终触发空值断言失败。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
相关产品推荐
相关产品推荐

