欢迎光临
我们一直在努力

Spring Boot 集成 Kafka:从环境搭建到消息收发

Spring Boot 集成 Kafka:从环境搭建到消息收发

Kafka 是一款高吞吐、可持久化、支持水平扩展的分布式事件流平台,常用于异步解耦、日志采集、订单事件、数据同步和流式处理。

本文将通过一个“订单创建事件”示例,演示如何在 Spring Boot 中完成:

  • 使用 Docker 启动 Kafka
  • 创建 Kafka Topic
  • 使用 KafkaTemplate 发送 JSON 消息
  • 使用 @KafkaListener 消费消息
  • 查看 Topic 和消费结果
  • 了解消息可靠性与常见问题

一、环境要求

本文使用以下环境:

组件版本
JDK 21
Spring Boot 4.1.0
Apache Kafka 4.3.1
Maven 3.9+
Docker 20.10.4+

Spring Boot 会管理 spring-kafka 的依赖版本,因此不需要单独指定版本。

二、Kafka 核心概念

在开始编码之前,先了解几个常用概念。

1. Producer

Producer,即消息生产者,负责将消息发送到 Kafka。

2. Consumer

Consumer,即消息消费者,负责从 Kafka 中读取并处理消息。

3. Topic

Topic 是消息的逻辑分类。例如:

order-created
payment-success
user-registered

4. Partition

一个 Topic 可以划分为多个 Partition。

Partition 是 Kafka 实现并发处理和水平扩展的基础。同一个 Partition 内的消息有序,但不同 Partition 之间不保证全局顺序。

5. Consumer Group

多个消费者可以组成一个 Consumer Group。

同一个 Consumer Group 中,一个 Partition 同一时间只会分配给一个消费者。不同 Consumer Group 则可以分别消费同一条消息。

6. Offset

Offset 是消息在 Partition 中的位置。Kafka 使用 Offset 记录消费者已经读取到哪里。

三、使用 Docker 启动 Kafka

创建 docker-compose.yml:

services:
kafka:
image: apache/kafka:4.3.1
container_name: kafka
ports:
"9092:9092"

启动 Kafka:

docker compose up -d

查看容器状态:

docker compose ps

查看 Kafka 日志:

docker compose logs -f kafka

如果不想使用 Docker Compose,也可以直接运行:

docker run -d \\
–name kafka \\
-p 9092:9092 \\
apache/kafka:4.3.1

这里使用的是 Kafka 自带的 KRaft 模式,不再依赖 ZooKeeper。该单节点配置仅适合本地开发和学习。

四、创建 Spring Boot 项目

项目结构如下:

kafka-demo
├── pom.xml
└── src
└── main
├── java
│ └── com.example.kafkademo
│ ├── KafkaDemoApplication.java
│ ├── config
│ │ └── KafkaTopicConfig.java
│ ├── consumer
│ │ └── OrderEventConsumer.java
│ ├── controller
│ │ └── OrderController.java
│ ├── model
│ │ └── OrderEvent.java
│ └── producer
│ └── OrderEventProducer.java
└── resources
└── application.yml

五、添加 Maven 依赖

在 pom.xml 中添加 Spring Kafka:

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
https://maven.apache.org/xsd/maven-4.0.0.xsd"
>

<modelVersion>4.0.0</modelVersion>

<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>4.1.0</version>
<relativePath/>
</parent>

<groupId>com.example</groupId>
<artifactId>kafka-demo</artifactId>
<version>0.0.1-SNAPSHOT</version>
<name>kafka-demo</name>

<properties>
<java.version>21</java.version>
</properties>

<dependencies>
<!– Spring MVC –>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>

<!– Spring Kafka –>
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>

<!– 测试依赖 –>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>

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

<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
</plugin>
</plugins>
</build>

</project>

六、配置 Kafka

创建 src/main/resources/application.yml:

server:
port: 8080

spring:
application:
name: kafkademo

kafka:
bootstrap-servers: localhost:9092

producer:
# 等待所有同步副本确认消息
acks: all

# 发送失败时的重试次数
retries: 3

key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JacksonJsonSerializer

properties:
# 防止重试导致消息被重复写入 Kafka
"[enable.idempotence]": true

# 不在消息头中携带 Java 类名,降低服务间耦合
"[spring.json.add.type.headers]": false

consumer:
group-id: orderservice

# 当前消费组没有 Offset 时,从最早的消息开始读取
auto-offset-reset: earliest

# 由 Spring Kafka 在消息成功处理后提交 Offset
enable-auto-commit: false

key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.springframework.kafka.support.serializer.JacksonJsonDeserializer

properties:
# 指定没有类型消息头时的目标类型
"[spring.json.value.default.type]": "com.example.kafkademo.model.OrderEvent"

# 只信任当前项目的消息类型,不建议配置为 *
"[spring.json.trusted.packages]": "com.example.kafkademo.model"

# 不依赖生产者传递的 Java 类型消息头
"[spring.json.use.type.headers]": false

listener:
# 每条消息处理成功后提交对应 Offset
ack-mode: record

app:
kafka:
topic:
order-created: ordercreated

这里需要特别注意:

  • bootstrap-servers 是 Kafka 地址。
  • group-id 是消费者组名称。
  • acks: all 表示等待所有同步副本确认。
  • enable.idempotence: true 可以减少生产者重试造成的重复写入。
  • 关闭 Java 类型消息头,可以降低生产者与消费者之间的类名耦合。
  • 消费端仍然必须做好业务幂等,因为网络重试、消费失败等情况仍可能导致重复消费。

七、创建启动类

package com.example.kafkademo;

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;

@SpringBootApplication
public class KafkaDemoApplication {

public static void main(String[] args) {
SpringApplication.run(KafkaDemoApplication.class, args);
}
}

Spring Boot 会自动配置 KafkaTemplate 和 Kafka Listener Container,因此一般不需要额外添加 @EnableKafka。

八、定义订单事件

创建 OrderEvent.java:

package com.example.kafkademo.model;

import java.math.BigDecimal;
import java.time.LocalDateTime;

public record OrderEvent(
String orderId,
Long userId,
BigDecimal amount,
LocalDateTime createdAt
) {
}

Kafka 中实际保存的是序列化后的 JSON 数据,例如:

{
"orderId": "ORDER-20260801-001",
"userId": 10001,
"amount": 99.90,
"createdAt": "2026-08-01T15:30:00"
}

九、创建 Topic

创建 KafkaTopicConfig.java:

package com.example.kafkademo.config;

import org.apache.kafka.clients.admin.NewTopic;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.TopicBuilder;

@Configuration
public class KafkaTopicConfig {

@Bean
public NewTopic orderCreatedTopic(
@Value("${app.kafka.topic.order-created}") String topicName) {

return TopicBuilder.name(topicName)
.partitions(3)
.replicas(1)
.build();
}
}

应用启动时,Spring Boot 会通过 KafkaAdmin 检查 Topic:

  • Topic 不存在时自动创建。
  • Topic 已存在时不会重复创建。
  • 本地单节点 Kafka 的副本数只能配置为 1。

生产环境中通常建议通过运维平台、Terraform 或部署脚本统一管理 Topic,而不是由应用自动创建。

十、编写消息生产者

创建 OrderEventProducer.java:

package com.example.kafkademo.producer;

import com.example.kafkademo.model.OrderEvent;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.SendResult;
import org.springframework.stereotype.Service;

import java.util.concurrent.CompletableFuture;

@Service
public class OrderEventProducer {

private final KafkaTemplate<String, Object> kafkaTemplate;
private final String topicName;

public OrderEventProducer(
KafkaTemplate<String, Object> kafkaTemplate,
@Value("${app.kafka.topic.order-created}") String topicName) {
this.kafkaTemplate = kafkaTemplate;
this.topicName = topicName;
}

public CompletableFuture<SendResult<String, Object>> send(OrderEvent event) {
/*
* 使用 orderId 作为消息 Key。
* 相同 orderId 的消息会被发送到相同 Partition,
* 从而保证同一订单内的消息顺序。
*/

return kafkaTemplate.send(topicName, event.orderId(), event);
}
}

KafkaTemplate.send() 是异步方法,返回:

CompletableFuture<SendResult<K, V>>

通过 SendResult 可以获取:

  • Topic
  • Partition
  • Offset
  • 消息时间戳

不要在每次发送后调用 flush(),否则可能破坏 Kafka 的批量发送能力并降低吞吐量。

十一、提供消息发送接口

创建 OrderController.java:

package com.example.kafkademo.controller;

import com.example.kafkademo.model.OrderEvent;
import com.example.kafkademo.producer.OrderEventProducer;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.*;

import java.util.LinkedHashMap;
import java.util.Map;
import java.util.concurrent.CompletableFuture;

@RestController
@RequestMapping("/api/orders")
public class OrderController {

private final OrderEventProducer producer;

public OrderController(OrderEventProducer producer) {
this.producer = producer;
}

@PostMapping("/events")
public CompletableFuture<ResponseEntity<Map<String, Object>>> send(
@RequestBody OrderEvent event) {

return producer.send(event)
.thenApply(result -> {
var metadata = result.getRecordMetadata();

Map<String, Object> response = new LinkedHashMap<>();
response.put("message", "消息发送成功");
response.put("topic", metadata.topic());
response.put("partition", metadata.partition());
response.put("offset", metadata.offset());
response.put("orderId", event.orderId());

return ResponseEntity.ok(response);
});
}
}

这里没有直接返回“已发送”,而是等待 Kafka 返回确认结果后,再将 Topic、Partition 和 Offset 返回给客户端。

十二、编写消息消费者

创建 OrderEventConsumer.java:

package com.example.kafkademo.consumer;

import com.example.kafkademo.model.OrderEvent;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;

@Component
public class OrderEventConsumer {

private static final Logger log =
LoggerFactory.getLogger(OrderEventConsumer.class);

@KafkaListener(topics = "${app.kafka.topic.order-created}")
public void consume(ConsumerRecord<String, OrderEvent> record) {
OrderEvent event = record.value();

log.info(
"收到订单事件,key={}, orderId={}, userId={}, amount={}, " +
"topic={}, partition={}, offset={}",
record.key(),
event.orderId(),
event.userId(),
event.amount(),
record.topic(),
record.partition(),
record.offset()
);

// 在这里执行具体业务,例如:
// 1. 保存订单操作日志
// 2. 发送通知
// 3. 更新搜索索引
// 4. 触发下游业务
}
}

当方法正常执行结束后,Spring Kafka 会根据 ack-mode: record 提交该消息的 Offset。

如果方法抛出异常,Offset 不会按照成功路径提交,消息可能被重新投递。因此,消费逻辑必须具备幂等性。

十三、启动项目

首先确认 Kafka 正在运行:

docker compose ps

然后启动 Spring Boot:

mvn spring-boot:run

也可以先打包:

mvn clean package
java -jar target/kafka-demo-0.0.1-SNAPSHOT.jar

十四、发送测试消息

Linux 或 macOS

curl -X POST "http://localhost:8080/api/orders/events" \\
-H "Content-Type: application/json" \\
-d '{
"orderId": "ORDER-20260801-001",
"userId": 10001,
"amount": 99.90,
"createdAt": "2026-08-01T15:30:00"
}'

Windows PowerShell

$body = @{
orderId = "ORDER-20260801-001"
userId = 10001
amount = 99.90
createdAt = "2026-08-01T15:30:00"
} | ConvertTo-Json

Invoke-RestMethod `
Method Post `
Uri "http://localhost:8080/api/orders/events" `
ContentType "application/json" `
Body $body

接口响应示例:

{
"message": "消息发送成功",
"topic": "order-created",
"partition": 1,
"offset": 0,
"orderId": "ORDER-20260801-001"
}

消费者控制台输出:

收到订单事件,key=ORDER-20260801-001,
orderId=ORDER-20260801-001,
userId=10001,
amount=99.90,
topic=order-created,
partition=1,
offset=0

十五、使用 Kafka 命令查看 Topic

查看 Topic 列表:

docker compose exec kafka \\
/opt/kafka/bin/kafka-topics.sh \\
–bootstrap-server localhost:9092 \\
–list

查看 Topic 详情:

docker compose exec kafka \\
/opt/kafka/bin/kafka-topics.sh \\
–bootstrap-server localhost:9092 \\
–describe \\
–topic order-created

直接读取 Topic 中的消息:

docker compose exec kafka \\
/opt/kafka/bin/kafka-console-consumer.sh \\
–bootstrap-server localhost:9092 \\
–topic order-created \\
–from-beginning \\
–property print.key=true \\
–property key.separator=" => "

十六、为什么要使用消息 Key

发送消息时使用了订单号作为 Key:

kafkaTemplate.send(topicName, event.orderId(), event);

Kafka 默认会根据 Key 计算消息所属的 Partition。

这意味着:

ORDER-001 创建
ORDER-001 支付
ORDER-001 发货

只要这些消息使用相同 Key,就会进入相同 Partition,从而保持同一订单的消息顺序。

需要注意,Kafka 只能保证一个 Partition 内有序,不能保证多个 Partition 之间全局有序。

十七、如何理解消费者组

假设 order-created 有三个 Partition:

Partition 0
Partition 1
Partition 2

当一个消费组中有三个消费者实例时,通常每个消费者会分配到一个 Partition:

Consumer 1 -> Partition 0
Consumer 2 -> Partition 1
Consumer 3 -> Partition 2

如果启动四个消费者,由于只有三个 Partition,多出来的消费者不会分配到 Partition。

如果希望通知服务和积分服务都收到订单事件,需要使用不同的 Consumer Group:

notification-service
point-service

同一消费组用于负载均衡,不同消费组用于广播式消费。

十八、Kafka 能否保证消息绝不重复

不能简单地认为 Kafka 消费只会执行一次。

例如:

  • 消费者处理业务成功。
  • 数据库已经提交。
  • 消费者还没有提交 Offset。
  • 应用突然宕机。
  • 应用重启后重新消费这条消息。
  • 此时同一业务可能被执行两次。

    因此,消费端应使用业务唯一键实现幂等,例如:

    CREATE TABLE kafka_consume_record (
    id BIGINT PRIMARY KEY AUTO_INCREMENT,
    consumer_group VARCHAR(100) NOT NULL,
    message_key VARCHAR(200) NOT NULL,
    created_at DATETIME NOT NULL,
    UNIQUE KEY uk_group_message (consumer_group, message_key)
    );

    消费前先判断该业务消息是否已经处理,或者直接利用数据库唯一索引阻止重复写入。

    不要只使用 Kafka 的 Offset 作为业务幂等依据。更推荐使用订单号、支付流水号、事件 ID 等稳定的业务标识。

    十九、数据库事务与 Kafka 的一致性

    下面的代码存在一致性风险:

    orderRepository.save(order);
    kafkaTemplate.send("order-created", orderEvent);

    可能出现:

    • 数据库提交成功,但 Kafka 发送失败。
    • Kafka 发送成功,但数据库事务回滚。
    • 应用在两个操作之间宕机。

    生产环境中可以使用 Outbox Pattern:

  • 在同一个数据库事务中写入业务表和 Outbox 事件表。
  • 后台任务读取未发送的 Outbox 事件。
  • 将事件发送到 Kafka。
  • 发送成功后更新 Outbox 状态。
  • 消费端继续使用业务唯一键保证幂等。
  • 这比直接把“数据库事务”和“Kafka 发送”拼在一起更容易排查和补偿。

    二十、常见问题

    1. Connection to node -1 could not be established

    常见原因:

    • Kafka 没有启动。
    • bootstrap-servers 地址错误。
    • Docker 端口没有映射。
    • Kafka 的 advertised.listeners 配置不正确。

    首先检查:

    docker compose ps
    docker compose logs kafka

    2. Topic 不存在

    确认 Topic 是否创建:

    docker compose exec kafka \\
    /opt/kafka/bin/kafka-topics.sh \\
    –bootstrap-server localhost:9092 \\
    –list

    3. JSON 反序列化失败

    检查:

    • 生产者和消费者的字段结构是否兼容。
    • spring.json.value.default.type 是否为完整类名。
    • spring.json.trusted.packages 是否包含消息类所在包。
    • 消息是否确实为合法 JSON。

    不要为了快速解决问题直接配置:

    spring.json.trusted.packages: "*"

    这会扩大反序列化信任范围,存在安全风险。

    4. 消费者收不到历史消息

    auto-offset-reset: earliest 只在消费组没有已提交 Offset 时生效。

    如果当前消费组已经提交过 Offset,可以临时修改消费组名称:

    spring:
    kafka:
    consumer:
    group-id: orderservicetest2

    5. 启动多个消费者,但只有一个收到消息

    如果这些消费者属于同一个 Consumer Group,那么一条消息只会交给组内的一个消费者。

    如果希望每个服务都收到消息,需要配置不同的 Consumer Group。

    6. 副本数配置失败

    本地只有一个 Kafka Broker,因此只能配置:

    .replicas(1)

    如果配置为 3,会因为 Broker 数量不足导致 Topic 创建失败。

    二十一、生产环境建议

    本文示例适合学习和本地开发。生产环境还应考虑:

  • Kafka 集群至少部署三个 Broker。
  • Topic 副本数通常设置为三个。
  • 使用 SASL/SSL、ACL 和网络访问控制。
  • 为生产者配置合理的超时、重试和压缩方式。
  • 为消费者配置重试队列和死信 Topic。
  • 所有消费逻辑必须保证业务幂等。
  • 数据库与 Kafka 一致性场景优先考虑 Outbox Pattern。
  • 监控 Producer 发送失败、Consumer Lag 和消费异常。
  • 为消息增加 eventId、eventType、occurredAt 和版本号。
  • 不要在消息中传递密码、Token、银行卡号等敏感数据。
  • 不要无限重试永久性错误,避免问题消息阻塞整个 Partition。
  • Topic 建议由基础设施脚本统一创建和管理。
  • 一个更完整的事件结构可以设计为:

    {
    "eventId": "019ABCDEF123456",
    "eventType": "ORDER_CREATED",
    "eventVersion": 1,
    "occurredAt": "2026-08-01T15:30:00+08:00",
    "data": {
    "orderId": "ORDER-20260801-001",
    "userId": 10001,
    "amount": 99.90
    }
    }

    二十二、总结

    Spring Boot 通过 Spring Kafka 简化了 Kafka 的使用:

    • 使用 spring.kafka.* 完成客户端配置。
    • 使用 KafkaTemplate 异步发送消息。
    • 使用 @KafkaListener 消费消息。
    • 使用 NewTopic 在开发环境中自动创建 Topic。
    • 使用 JSON Serializer 和 Deserializer 传输对象。

    不过,完成消息收发只是 Kafka 应用的第一步。真正进入生产环境时,还需要重点处理:

    • 消息重复
    • 消息丢失
    • 消费失败
    • 消息积压
    • 顺序消费
    • 数据库一致性
    • 重试与死信队列
    • 监控和告警

    只有同时做好可靠性、幂等性和可观测性,Kafka 才能真正成为稳定的业务基础设施。

    参考资料

    • Spring Boot:Apache Kafka Support
    • Spring Kafka:Sending Messages
    • Spring Kafka:@KafkaListener
    • Apache Kafka 官方网站
    赞(0)
    未经允许不得转载:171主机测评 » Spring Boot 集成 Kafka:从环境搭建到消息收发
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址