使用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发送状态校验:在发送代码中添加发送状态判断,确认消息是否发送成功:
如果发送失败,检查Kafka集群连接配置,确保boolean isSent = streamBridge.send(topic, pageEvent); System.out.println("消息发送结果:" + isSent);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
相关产品推荐
相关产品推荐

