Apache Kafka 是一个高吞吐、分布式的流处理平台与开源消息引擎,依托 追加日志(Commit Log)零拷贝(Zero-Copy)分区并发模型 实现海量事件的持久化存储与低延迟流转。

本文深入介绍 Spring Boot 与 Kafka 的工程集成方案,涵盖本地环境准备、核心依赖配置、生产者与消费者开发、JSON 对象序列化以及分布式事务消息等实践。

 

1. 什么是 Kafka#


Apache Kafka 是一个分布式事件流平台,核心基于发布与订阅(Pub/Sub)模式构建,具备高并发、高吞吐、可水平扩展与高容错等企业级特征。

1.1 核心特性#

  • 发布订阅与消息持久化:基于追加日志(Commit Log)组织数据,所有写入均为磁盘顺序 I/O,支持消息长期持久化与多组独立消费。
  • 实时流式数据处理:支持毫秒级流式数据接入与处理,可与 Flink、Spark Streaming 等大数据计算框架深度协同。
  • 高吞吐与水平扩展:通过分区(Partition)将主题数据分散在多个 Broker 节点,支持集群无损横向伸缩与并行高并发读写。
  • 高可用与多副本机制:提供 Leader 与 Follower 多副本冗余,配合 ISR(In-Sync Replicas)机制保障节点故障时数据不丢失。

1.2 核心应用场景#

在分布式与微服务架构中,Kafka 通常充当事件中枢或数据传输管道,广泛应用于以下业务场景:

  • 微服务间异步解耦:将非核心链路业务(如短信通知、权益发放、统计分析)解耦为异步事件驱动,降低服务间调用耦合与故障扩散风险。
  • 海量日志收集与监控:作为 ELK(Elasticsearch + Logstash + Kibana)或 ClickHouse 的前置缓冲区,汇聚全链路访问日志与系统监控指标。
  • 流数据管道与同步:捕获关系型数据库变更数据(CDC),实时同步并分发至搜索引擎、数据仓库或缓存集群。
  • 流量削峰填谷:在秒杀或大促等突发流量高峰期暂存高并发请求,下游根据自身承载能力匀速拉取消费,保护核心数据库。

 

2. 环境准备#


在开始集成开发之前,请确保本地或服务器已完成基础运行环境与服务的安装部署。

2.1 依赖环境要求#

  • JDK:JDK 1.8 或更高版本(推荐 JDK 17 / 21)。
  • Spring Boot:Spring Boot 2.7.x 或 3.x。
  • 构建工具:Maven 3.6+ 或 Gradle 7+。
  • Kafka 服务端:Kafka 2.8+(含内置 Zookeeper 或 KRaft 模式)。

2.2 启动本地 Kafka 与 Zookeeper#

(1) 访问 Kafka 官方下载页面 获取二进制解压包,解压后进入程序根目录。

(2) 启动 Zookeeper 协调服务进程:

1# 启动 Zookeeper 协调服务进程
2bin/zookeeper-server-start.sh config/zookeeper.properties

(3) 启动 Kafka Broker 服务节点:

1# 启动 Kafka 服务节点
2bin/kafka-server-start.sh config/server.properties

Note

默认配置下,Zookeeper 服务监听 2181 端口,Kafka Broker 服务监听 9092 端口。若需在 Linux 后台常驻运行,可在命令首部增加 nohup 并在命令末尾追加 &。

 

3. Kafka 集成 Spring Boot#


Spring Kafka(spring-kafka)对底层的 kafka-clients 进行了封装,提供了模板化发送工具 KafkaTemplate 以及声明式监听注解 @KafkaListener。

3.1 添加依赖#

配置文件默认路径:pom.xml

 1<dependencies>
 2    <!-- Spring Kafka 核心起步依赖 -->
 3    <dependency>
 4        <groupId>org.springframework.kafka</groupId>
 5        <artifactId>spring-kafka</artifactId>
 6    </dependency>
 7    <!-- Spring Boot Web 起步依赖(用于暴露 RESTful 调试接口) -->
 8    <dependency>
 9        <groupId>org.springframework.boot</groupId>
10        <artifactId>spring-boot-starter-web</artifactId>
11    </dependency>
12</dependencies>

3.2 配置 Kafka 连接#

配置文件默认路径:src/main/resources/application.yml

 1spring:
 2  kafka:
 3    bootstrap-servers: localhost:9092 # Kafka Broker 节点集群连接地址
 4    producer:
 5      key-serializer: org.apache.kafka.common.serialization.StringSerializer # 消息键序列化类
 6      value-serializer: org.apache.kafka.common.serialization.StringSerializer # 消息体序列化类
 7      acks: 1 # 生产者应答机制(0: 不等待, 1: Leader写入成功, all: ISR全副本写入)
 8      retries: 3 # 发送失败重试次数
 9      batch-size: 16384 # 批量发送阈值大小(字节)
10      properties:
11        linger.ms: 1 # 批量发送等待缓冲延迟(毫秒)
12    consumer:
13      group-id: test-group # 默认消费者组 ID
14      auto-offset-reset: earliest # 无初始位移或位移越界时的策略(earliest: 从头消费, latest: 从最新消费)
15      enable-auto-commit: true # 启用自动位移提交
16      auto-commit-interval: 1000 # 自动提交位移间隔周期(毫秒)
17      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer # 消息键反序列化类
18      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer # 消息体反序列化类
19    template:
20      default-topic: test-topic # 默认发送目标主题

 

4. 生产者开发#


通过 Spring Kafka 提供的 KafkaTemplate 组件,开发者可以便捷地向指定 Topic 投递同步或异步消息。

4.1 定义生产者服务#

创建生产者业务类,封装消息发送逻辑并挂载异步回调:

 1package com.example.kafka.service;
 2
 3import org.slf4j.Logger;
 4import org.slf4j.LoggerFactory;
 5import org.springframework.kafka.core.KafkaTemplate;
 6import org.springframework.kafka.support.SendResult;
 7import org.springframework.stereotype.Service;
 8
 9import java.util.concurrent.CompletableFuture;
10
11@Service
12public class KafkaProducerService {
13
14    private static final Logger log = LoggerFactory.getLogger(KafkaProducerService.class);
15
16    private final KafkaTemplate<String, String> kafkaTemplate;
17
18    // 依赖注入 KafkaTemplate
19    public KafkaProducerService(KafkaTemplate<String, String> kafkaTemplate) {
20        this.kafkaTemplate = kafkaTemplate;
21    }
22
23    /**
24     * 发送字符串消息并异步监听发送结果
25     *
26     * @param topic   目标主题名称
27     * @param message 待发送的消息内容
28     */
29    public void sendMessage(String topic, String message) {
30        CompletableFuture<SendResult<String, String>> future = kafkaTemplate.send(topic, message);
31        future.whenComplete((result, ex) -> {
32            if (ex == null) {
33                log.info("消息发送成功: topic={}, partition={}, offset={}",
34                        result.getRecordMetadata().topic(),
35                        result.getRecordMetadata().partition(),
36                        result.getRecordMetadata().offset());
37            } else {
38                log.error("消息发送失败: topic={}, error={}", topic, ex.getMessage(), ex);
39            }
40        });
41    }
42}

4.2 编写 API 接口与发送验证#

(1) 编写 REST Controller 控制器暴露触发入口:

 1package com.example.kafka.controller;
 2
 3import com.example.kafka.service.KafkaProducerService;
 4import org.springframework.web.bind.annotation.PostMapping;
 5import org.springframework.web.bind.annotation.RequestMapping;
 6import org.springframework.web.bind.annotation.RequestParam;
 7import org.springframework.web.bind.annotation.RestController;
 8
 9@RestController
10@RequestMapping("/kafka")
11public class KafkaProducerController {
12
13    private final KafkaProducerService kafkaProducerService;
14
15    public KafkaProducerController(KafkaProducerService kafkaProducerService) {
16        this.kafkaProducerService = kafkaProducerService;
17    }
18
19    /**
20     * HTTP POST 发送测试消息
21     */
22    @PostMapping("/send")
23    public String sendMessage(@RequestParam String topic, @RequestParam String message) {
24        kafkaProducerService.sendMessage(topic, message);
25        return "Message sent: " + message;
26    }
27}

(2) 使用 curl 发起 HTTP 请求测试消息发送:

1curl -X POST "http://localhost:8080/kafka/send?topic=test-topic&message=HelloKafka"

 

5. 消费者开发#


Spring Kafka 提供了声明式注解 @KafkaListener,后台自动维护消费线程池、分区拉取与位移提交。

5.1 定义消费者逻辑#

创建消费者服务类,使用 @KafkaListener 注解声明监听的目标主题和消费者组:

 1package com.example.kafka.consumer;
 2
 3import org.apache.kafka.clients.consumer.ConsumerRecord;
 4import org.slf4j.Logger;
 5import org.slf4j.LoggerFactory;
 6import org.springframework.kafka.annotation.KafkaListener;
 7import org.springframework.stereotype.Service;
 8
 9@Service
10public class KafkaConsumerService {
11
12    private static final Logger log = LoggerFactory.getLogger(KafkaConsumerService.class);
13
14    /**
15     * 监听指定 Topic 的消息
16     *
17     * @param record 消费记录对象,包含元数据信息
18     */
19    @KafkaListener(topics = "test-topic", groupId = "test-group")
20    public void listen(ConsumerRecord<String, String> record) {
21        log.info("接收到 Kafka 消息: topic={}, partition={}, offset={}, key={}, value={}",
22                record.topic(),
23                record.partition(),
24                record.offset(),
25                record.key(),
26                record.value());
27    }
28}

5.2 测试消费者#

(1) 使用 curl 再次向 test-topic 发送测试消息:

1curl -X POST "http://localhost:8080/kafka/send?topic=test-topic&message=HelloKafkaConsumer"

(2) 观察 Spring Boot 应用控制台输出,确认消息被正常接收与消费:

1接收到 Kafka 消息: topic=test-topic, partition=0, offset=1, key=null, value=HelloKafkaConsumer

 

6. Kafka 高级功能#


在生产实践中,除了简单的纯文本消息,系统往往需要传输复杂的 Java 实体对象,或在跨组件操作中保证消息发送与本地业务的原子性。

6.1 自定义消息对象的序列化#

默认情况下 Kafka 采用字符串序列化器。通过集成 Jackson 与 Spring Kafka 提供的 JsonSerializer / JsonDeserializer,可直接在生产者与消费者之间流转 Java DTO 实体。

(1) 引入 Jackson 与 Lombok 依赖(配置文件默认路径:pom.xml):

 1<!-- Jackson 核心依赖(用于 JSON 序列化与反序列化) -->
 2<dependency>
 3    <groupId>com.fasterxml.jackson.core</groupId>
 4    <artifactId>jackson-databind</artifactId>
 5</dependency>
 6<!-- Lombok 依赖(用于消除实体模型样板代码) -->
 7<dependency>
 8    <groupId>org.projectlombok</groupId>
 9    <artifactId>lombok</artifactId>
10    <optional>true</optional>
11</dependency>

(2) 定义消息实体传输模型 User:

 1package com.example.kafka.entity;
 2
 3import lombok.AllArgsConstructor;
 4import lombok.Data;
 5import lombok.NoArgsConstructor;
 6
 7import java.io.Serializable;
 8
 9@Data
10@NoArgsConstructor
11@AllArgsConstructor
12public class User implements Serializable {
13
14    private static final long serialVersionUID = 1L;
15
16    private String name;
17    private int age;
18}

(3) 配置 JSON 序列化与信任包规则(配置文件默认路径:src/main/resources/application.yml):

 1spring:
 2  kafka:
 3    producer:
 4      key-serializer: org.apache.kafka.common.serialization.StringSerializer # 键序列化器
 5      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer # 值序列化器:Spring 提供的 JSON 序列化工具
 6    consumer:
 7      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer # 键反序列化器
 8      value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer # 值反序列化器:Spring 提供的 JSON 反序列化工具
 9      properties:
10        spring.json.trusted.packages: "*" # 反序列化信任包路径,"*" 表示允许反序列化所有包路径下的对象

(4) 生产者发送实体对象:

 1package com.example.kafka.service;
 2
 3import com.example.kafka.entity.User;
 4import org.springframework.kafka.core.KafkaTemplate;
 5import org.springframework.stereotype.Service;
 6
 7@Service
 8public class UserProducerService {
 9
10    private final KafkaTemplate<String, Object> kafkaTemplate;
11
12    public UserProducerService(KafkaTemplate<String, Object> kafkaTemplate) {
13        this.kafkaTemplate = kafkaTemplate;
14    }
15
16    /**
17     * 发送结构化实体对象
18     */
19    public void sendUserMessage(String topic, User user) {
20        kafkaTemplate.send(topic, user);
21    }
22}

(5) 消费者接收实体对象:

 1package com.example.kafka.consumer;
 2
 3import com.example.kafka.entity.User;
 4import org.slf4j.Logger;
 5import org.slf4j.LoggerFactory;
 6import org.springframework.kafka.annotation.KafkaListener;
 7import org.springframework.stereotype.Service;
 8
 9@Service
10public class UserConsumerService {
11
12    private static final Logger log = LoggerFactory.getLogger(UserConsumerService.class);
13
14    /**
15     * 自动反序列化为 User 实体对象
16     */
17    @KafkaListener(topics = "user-topic", groupId = "test-group")
18    public void listenUser(User user) {
19        log.info("成功接收到结构化对象: {}", user);
20    }
21}

6.2 Kafka 事务支持#

Kafka 支持基于两阶段提交的跨分区事务消息,确保一组消息要么全部写入成功,要么全部放弃,常用于实现生产消费之间的原子流转。

(1) 开启生产者事务支持(配置文件默认路径:src/main/resources/application.yml):

1spring:
2  kafka:
3    producer:
4      transaction-id-prefix: trans- # 事务 ID 前缀,配置后自动开启 Kafka 事务管理器

(2) 在生产者业务中执行事务发送:

 1package com.example.kafka.service;
 2
 3import org.springframework.kafka.core.KafkaTemplate;
 4import org.springframework.stereotype.Service;
 5import org.springframework.transaction.annotation.Transactional;
 6
 7@Service
 8public class TransactionalProducerService {
 9
10    private final KafkaTemplate<String, String> kafkaTemplate;
11
12    public TransactionalProducerService(KafkaTemplate<String, String> kafkaTemplate) {
13        this.kafkaTemplate = kafkaTemplate;
14    }
15
16    /**
17     * 方式一:基于 executeInTransaction 显式事务回调
18     */
19    public void sendInTransaction(String topic, String message) {
20        kafkaTemplate.executeInTransaction(operations -> {
21            operations.send(topic, message);
22            return true;
23        });
24    }
25
26    /**
27     * 方式二:基于 Spring 的 @Transactional 声明式事务管理
28     */
29    @Transactional(rollbackFor = Exception.class)
30    public void sendWithAnnotation(String topic, String message) {
31        kafkaTemplate.send(topic, message);
32        // 若后续业务逻辑抛出未捕获异常,Kafka 事务将自动中止回滚
33    }
34}

Tip

开启事务后,消费者端的 isolation.level 默认是 read_uncommitted。如需过滤未提交或已中止的脏消息,需将消费者配置项 isolation.level 设置为 read_committed。

 

7. 总结#


Spring Kafka 极大地简化了与 Apache Kafka 的集成流程,提升了分布式流数据处理与消息通信的研发效率。

7.1 核心功能回顾#

(1) 生产者开发:使用 KafkaTemplate 发送同步与异步消息,配合回调函数监听消息分区、Offset 及发送状态。

(2) 消费者开发:使用 @KafkaListener 声明式监听目标主题与消费组,由后台线程池自动管理分区拉取与位移提交。

(3) 高级功能支持:集成 Jackson 序列化器实现复杂 Java 对象的透明序列化流转,并通过 transaction-id-prefix 开启分布式事务保障消息发送原子性。

7.2 生产落地建议#

  • 位移提交模式:对于金融支付与核心业务链路,建议将自动提交改为手动确认(AckMode.MANUAL_IMMEDIATE),彻底消除消息丢失隐患。
  • 死信队列与容错:针对业务消费异常与重试耗尽场景,引入死信队列(DLT / Dead Letter Queue)与告警机制,防止消费链路级联阻塞。
  • 集群高可用配置:在生产部署中将生产者 acks 设置为 all,并结合 Broker 端的 min.insync.replicas 策略,确保分区多副本强一致落盘。