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
相关产品推荐
相关产品推荐

