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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 01:57:07