SpringBoot中如何在WebSocket的onText方法正确实现KafkaTemplate?
我正在通过HTTP WebSocket接收实时数据流,想要把消息写入Kafka,但当前的消息发送实现存在问题,求排查原因。
项目包结构
- seismickafkaproducer - kafkaconfig - streaming
核心代码
streaming包下的WebSocketClient类
@EnableKafka public class WebSocketClient implements WebSocket.Listener { private final CountDownLatch latch; private final JsonParser parser; @Autowired KafkaTemplate<String, String> kafkaTemplate; public WebSocketClient(CountDownLatch latch, JsonParser parser) { this.latch = latch; this.parser = parser; } public void onOpen(WebSocket webSocket) { WebSocket.Listener.super.onOpen(webSocket); } public CompletionStage<?> onText(WebSocket webSocket, CharSequence data, boolean last, KafkaTemplate<String, String> kafkaTemplate) { String message = data.toString(); System.out.println("Message received, writing to kafka"); kafkaTemplate.send("seismic", message); return WebSocket.Listener.super.onText(webSocket, data, last); } public void onError(WebSocket websocket, Throwable error) { System.out.println(websocket.toString()); WebSocket.Listener.super.onError(websocket, error); } }
kafkaconfig包下的KafkaProducerConfig类
@Configuration public class KafkaProducerConfig { @Value("${spring.kafka.bootstrap-servers}") private String bootstrapServer; public Map<String, Object> producerConfig() { Map<String, Object> properties = new HashMap<>(); properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServer); properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class); return properties; } @Bean public ProducerFactory<String, String> producerFactory() { return new DefaultKafkaProducerFactory<>(producerConfig()); } @Bean public KafkaTemplate<String, String> kafkaTemplate( ProducerFactory<String, String> producerFactory ) { return new KafkaTemplate<>(producerFactory); } }
kafkaconfig包下的KafkaTopicConfig类
@Configuration public class KafkaTopicConfig { @Bean public NewTopic seismicTopic() { return TopicBuilder.name("seismic").build(); } }
问题排查方向
onText方法参数覆盖依赖注入实例:onText方法中声明了
KafkaTemplate<String, String> kafkaTemplate参数,该参数不会被Spring自动注入,调用时若未传值会导致空指针异常。应删除此方法参数,直接使用类中通过@Autowired注入的kafkaTemplate实例。序列化配置不匹配:ProducerConfig中值序列化器使用
JsonSerializer,但KafkaTemplate泛型为<String, String>,JsonSerializer会将String序列化为带引号的JSON格式,不仅不符合预期,还可能引发序列化异常。需将VALUE_SERIALIZER_CLASS_CONFIG改为StringSerializer.class。@EnableKafka注解位置错误:
@EnableKafka应标注在配置类或主启动类上,而非WebSocketClient类,否则Kafka相关Bean无法被正确加载。WebSocketClient实例化方式错误:若WebSocketClient是手动new出来的,而非由Spring容器管理,
@Autowired注入的kafkaTemplate会为null。需确保WebSocketClient通过Spring容器创建,比如添加@Component注解,或通过构造器注入KafkaTemplate。缺少发送结果回调:kafkaTemplate.send是异步操作,当前代码未处理发送结果,无法得知是否发送失败。建议添加回调捕获异常:
kafkaTemplate.send("seismic", message) .addCallback( success -> System.out.println("消息发送成功:" + success.getRecordMetadata()), failure -> System.err.println("消息发送失败:" + failure.getMessage()) );
内容的提问来源于stack exchange,提问作者paul445

