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

SpringBoot中如何在WebSocket的onText方法正确实现KafkaTemplate?

WebSocket实时数据流写入Kafka失败问题排查

我正在通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 06:45:01