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

Spring 2.X转3.X:Kafka中@EnableBinding与@StreamListener替代方案咨询

Spring Cloud Stream 2.X到3.X迁移:Kafka消费者替代方案咨询

我正在进行Spring Cloud Stream从2.X到3.X的迁移,原消费者代码依赖@EnableBinding(Sink::class)和@StreamListener(Sink.INPUT)实现消息消费,但这两个注解已被废弃并移除。现在运行Kafka消费者测试时出现如下断言错误:

org.opentest4j.AssertionFailedError:
expected: 1L
but was: 0L

我尝试用@KafkaListener(topics=["kafka-test"], group-id="test-group")替换@StreamListener,但测试仍报相同错误。

我曾对消费方法做过如下改写:
原方法:

fun consume(message: Message<*>) {

改写为:

fun consume(): Consumer<Message<*>> = Consumer { message ->

但调用代码kafkaConsumer.consume(GenericMessage("""{"aV": 1, "dI": 3, "id": $alertId}""".toByteArray()))时,出现参数不匹配错误:

Too many arguments for public open fun consume(): Consumer<Message<*>>

之后我再次修改消费方法:

@KafkaListener(topics = ["test-topic"], groupId = "group-test")
fun consume(message: GenericMessage<ByteArray>) {

测试依旧出现expected:1L but was:0L的错误。

即使未迁移依赖,仅移除@EnableBinding和@StreamListener也会触发该错误。请问是否存在与这两个注解功能等效的Kafka消费者替代方案?

当前迁移后的依赖版本:

  • Spring-kafka 3.1.4
  • spring-cloud-stream 4.1.1
  • kafka-clients 6.2.1
  • spring-integration 6.1.6

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 02:33:15