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

如何通过Spring Boot Web应用创建与Kafka的持久连接并收发消息

Spring Boot 对接Kafka持久连接实现指南

一、基础依赖与配置

首先引入Spring Kafka官方起步依赖,Maven配置如下,版本需要和你Kafka集群的大版本匹配,比如Kafka 2.8.x对应spring-kafka 2.8.x:

<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
</dependency>

在application.yml(或properties)中配置Kafka集群基础信息:

spring:
  kafka:
    bootstrap-servers: kafka-node1:9092,kafka-node2:9092,kafka-node3:9092 # 替换为你的集群地址
    # 生产者配置
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
      retries: 3
      acks: all
    # 消费者配置
    consumer:
      group-id: biz-consumer-group-01 # 替换为你的消费者组ID
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      auto-offset-reset: earliest
      enable-auto-commit: true

二、核心功能实现

Spring Kafka底层直接复用Kafka官方Java客户端,默认维护和集群的持久TCP连接,不需要你手动处理连接建立、保活、断连重连逻辑。

生产者消息发送

直接注入单例的KafkaTemplate即可,应用启动时会自动完成和Kafka集群的连接初始化,所有发送请求复用同一个长连接:

@Service
public class KafkaMsgProducer {
    @Resource
    private KafkaTemplate<String, String> kafkaTemplate;

    // 对外提供的发送消息接口,供上游调用
    public void send(String topic, String msg) {
        kafkaTemplate.send(topic, msg);
        // 如果需要监听发送结果,可以加回调
        // kafkaTemplate.send(topic, msg).addCallback(success -> {}, fail -> {});
    }
}

消费者消息监听

用@KafkaListener注解标注消费方法,应用启动后会自动创建消费者实例,和集群保持持久长连接监听指定Topic,有新消息时自动触发方法执行:

@Component
public class KafkaMsgConsumer {
    @KafkaListener(topics = "biz-topic-01", groupId = "biz-consumer-group-01")
    public void onMessageReceived(String msg) {
        // 此处写你的消息处理逻辑
        System.out.println("收到消息:" + msg);
    }
}

三、连接稳定性优化配置

如果你的网络环境有防火墙空闲掐断连接的规则,可以添加以下参数优化连接保活:

spring:
  kafka:
    properties:
      # 连接空闲3分钟就发心跳保活,避免被防火墙断开,默认是5分钟
      connections.max.idle.ms: 180000
    consumer:
      properties:
        heartbeat.interval.ms: 3000 # 消费者心跳间隔3秒
        session.timeout.ms: 15000 # 15秒没收到心跳就判定消费者下线

四、通信协议的影响

Kafka默认使用基于TCP的自定义二进制协议,这也是性能最高、最稳定的对接方式,完全兼容Spring Kafka的长连接管理逻辑,不需要你做额外适配。
只有特殊场景需要调整配置,不影响核心连接逻辑:

  • 集群开启SSL/TLS加密传输时,只需额外添加证书、加密算法的配置项即可
  • 集群开启SASL身份认证时,只需补充账号、认证方式等配置项即可
    不要使用HTTP协议对接Kafka,HTTP为短连接模式,无法实现持久连接,性能也会大幅下降。

内容的提问来源于stack exchange,提问作者Jsh0s

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 16:45:06