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

使用Java 17的Spring Cloud Stream消费者无法接收Kafka消息

Spring Cloud Stream Kafka消费者收不到消息问题排查

问题描述

我正在开发一个基于Apache Kafka的Spring Cloud Stream消息传递应用,已搭建REST API通过StreamBridge向Kafka主题发布消息,但消费者始终无法接收消息。

发送端代码(REST Controller)

import lombok.AllArgsConstructor;
import ma.enset.KafkaSpringCloudStream.entity.PageEvent;
import org.springframework.cloud.stream.function.StreamBridge;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.RestController;
import java.util.Date;
import java.util.Random;

@RestController
@AllArgsConstructor
public class PageEventRestController {
    private StreamBridge streamBridge;

    @GetMapping("/publish/{topic}/{name}")
    private PageEvent publish(@PathVariable("topic") String topic,@PathVariable("name") String name){
        PageEvent pageEvent = PageEvent.builder()
                .name(name)
                .user(Math.random()>0.5 ? "User1" : "User2")
                .date(new Date((long) (Math.random()*System.currentTimeMillis())))
                .duration(new Random().nextInt(9001))
                .build();
        streamBridge.send(topic,pageEvent);
        return pageEvent;
    }
}

消费端代码(Service)

@Service
public class PageEventService {
    @Bean
    public Consumer<PageEvent> pageEventConsumer(){
        return (input) -> {
            System.out.println("********************************");
            System.out.println(input.toString());
            System.out.println("********************************");
        };
    }
}

配置文件(application.properties)

spring.application.name=app
spring.cloud.stream.bindings.pageEventConsumer-in-0.destination=testTopic

排查与解决步骤

  • 主题一致性检查:消费者绑定的是testTopic,但REST接口支持动态传入topic参数。如果调用接口时传入的不是testTopic(比如/publish/otherTopic/test),消息会发送到其他主题,消费者自然收不到。请确保调用时指定正确的testTopic,或者修改代码固定发送的主题。
  • 序列化/反序列化验证:确保PageEvent类具备JSON序列化能力——添加无参构造方法(可配合Lombok的@NoArgsConstructor)、Getter/Setter方法(或用@Data注解)。如果序列化失败,消息可能无法正常发送或被消费者丢弃。
  • StreamBridge发送状态校验:在发送代码中添加发送状态判断,确认消息是否发送成功:
    boolean isSent = streamBridge.send(topic, pageEvent);
    System.out.println("消息发送结果:" + isSent);
    
    如果发送失败,检查Kafka集群连接配置,确保spring.cloud.stream.kafka.binder.brokers指向正确的Kafka地址(默认是localhost:9092)。
  • Kafka主题存在性检查:确认Kafka已开启自动创建主题(默认auto.create.topics.enable=true),或手动创建testTopic主题。若主题不存在,消息无法存储,消费者也收不到。
  • 消费者配置校验:当前函数式消费者的绑定配置spring.cloud.stream.bindings.pageEventConsumer-in-0.destination=testTopic格式正确,若需要分组消费,可追加配置spring.cloud.stream.bindings.pageEventConsumer-in-0.group=consumer-group-1,避免重复消费。

内容的提问来源于stack exchange,提问作者Mohamed Sid Abdalla Rgagde

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 09:15:25