Spring 3.1.0-M2兼容的reactor-kafka版本咨询及报错排查
问题:Spring 3.1.0-M2兼容的reactor-kafka版本咨询
场景描述
正在搭建基于Spring 3.1.0-M2的响应式应用,通过reactor-kafka消费Kafka消息并使用R2DBC写入MS SQL数据库。使用reactor-kafka 1.2.2.RELEASE时触发java.lang.NoSuchMethodError错误,现寻求与Spring 3.1.0-M2兼容的reactor-kafka版本。
报错信息
2023-04-03T08:28:19.934-06:00 ERROR 22144 --- [event-tracker-1] reactor.core.scheduler.Schedulers : KafkaScheduler worker in group main failed with an uncaught exception java.lang.NoSuchMethodError: 'void org.apache.kafka.clients.consumer.Consumer.close(long, java.util.concurrent.TimeUnit)' at reactor.kafka.receiver.internals.DefaultKafkaReceiver$CloseEvent.run(DefaultKafkaReceiver.java:685) ~[reactor-kafka-1.2.2.RELEASE.jar:1.2.2.RELEASE] at reactor.kafka.receiver.internals.DefaultKafkaReceiver.doEvent(DefaultKafkaReceiver.java:401) ~[reactor-kafka-1.2.2.RELEASE.jar:1.2.2.RELEASE] at reactor.kafka.receiver.internals.DefaultKafkaReceiver.lambda$start$14(DefaultKafkaReceiver.java:335) ~[reactor-kafka-1.2.2.RELEASE.jar:1.2.2.RELEASE] at reactor.core.publisher.LambdaSubscriber.onNext(LambdaSubscriber.java:160) ~[reactor-core-3.5.4.jar:3.5.4] at reactor.core.publisher.FluxPublishOn$PublishOnSubscriber.runAsync(FluxPublishOn.java:440) ~[reactor-core-3.5.4.jar:3.5.4] at reactor.core.publisher.FluxPublishOn$PublishOnSubscriber.run(FluxPublishOn.java:527) ~[reactor-core-3.5.4.jar:3.5.4] at reactor.kafka.receiver.internals.KafkaSchedulers$EventScheduler.lambda$decorate$1(KafkaSchedulers.java:100) ~[reactor-kafka-1.2.2.RELEASE.jar:1.2.2.RELEASE] at reactor.core.scheduler.WorkerTask.call(WorkerTask.java:84) ~[reactor-core-3.5.4.jar:3.5.4] at reactor.core.scheduler.WorkerTask.call(WorkerTask.java:37) ~[reactor-core-3.5.4.jar:3.5.4] at java.base/java.util.concurrent.FutureTask.run$$$capture(FutureTask.java:264) ~[na:na] at java.base/java.util.concurrent.FutureTask.run(FutureTask.java) ~[na:na] at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:304) ~[na:na] at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) ~[na:na] at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) ~[na:na] at java.base/java.lang.Thread.run(Thread.java:833) ~[na:na]
pom.xml配置
<parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>3.1.0-M2</version> <relativePath/> </parent> <properties> <java.version>17</java.version> </properties> <dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-r2dbc</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-webflux</artifactId> </dependency> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-streams</artifactId> </dependency> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> </dependency> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> </dependency> <dependency> <groupId>com.microsoft.sqlserver</groupId> <artifactId>mssql-jdbc</artifactId> <scope>runtime</scope> </dependency> <dependency> <groupId>io.r2dbc</groupId> <artifactId>r2dbc-mssql</artifactId> <version>1.0.0.RELEASE</version> <scope>runtime</scope> </dependency> <dependency> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> <optional>true</optional> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-test</artifactId> <scope>test</scope> </dependency> <dependency> <groupId>io.projectreactor</groupId> <artifactId>reactor-test</artifactId> <scope>test</scope> </dependency> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka-test</artifactId> <scope>test</scope> </dependency> <dependency> <groupId>io.projectreactor.kafka</groupId> <artifactId>reactor-kafka</artifactId> <version>1.2.2.RELEASE</version> </dependency> </dependencies>
消费者服务代码片段
@Transactional Flux<EventInterface> consumeEventDTO() { return reactiveKafkaConsumerTemplate .receiveAutoAck() .delayElements(Duration.ofSeconds(2L)) // BACKPRESSURE .doOnNext(consumerRecord -> log.info("received key={}, value={} from topic={}, offset={}", consumerRecord.key(), consumerRecord.value(), consumerRecord.topic(), consumerRecord.offset()) ) .map(ConsumerRecord::value) .flatMap(this::exists) .flatMap(this::save) .onErrorComplete(UncategorizedR2dbcException.class) .doOnNext(event -> { log.info("successfully consumed {}={}", EventInterface.class.getSimpleName(), event); }) .doOnError(throwable -> log.error("something bad happened while consuming : {}", throwable.getMessage())); }
解决方案
错误原因
NoSuchMethodError源于版本冲突:reactor-kafka 1.2.2.RELEASE适配的是kafka-clients 2.8.x版本,但Spring Boot 3.1.0-M2默认引入的kafka-clients版本为3.4.0,两者的Consumer.close(long, TimeUnit)方法签名存在兼容性问题。
兼容版本选择
与Spring 3.1.0-M2兼容的reactor-kafka版本为1.3.x系列(如1.3.8.RELEASE),该系列专门适配Spring Boot 3.x、Reactor Core 3.5.x(项目中已使用3.5.4)及kafka-clients 3.0+版本。
修改pom依赖
将reactor-kafka的版本替换为1.3.x稳定版:
<dependency> <groupId>io.projectreactor.kafka</groupId> <artifactId>reactor-kafka</artifactId> <version>1.3.8.RELEASE</version> </dependency>
额外说明
- 无需手动指定kafka-clients版本,Spring Boot父依赖会自动管理与reactor-kafka 1.3.x匹配的版本。
- reactor-kafka 1.3.x完全兼容Spring Boot 3.1.x的响应式生态,包括R2DBC、WebFlux等组件。
内容的提问来源于stack exchange,提问作者Greg Walters
相关产品推荐
相关产品推荐

