Kafka生产者与消费者Demo
这一节从能跑起来的 Demo 开始,再解释生产者和消费者为什么这样配置。学习 Kafka 不建议一上来就背大量参数,先把“写入 Topic、拉取 Topic、提交 Offset”跑通。
本地启动 Kafka
官方 Quickstart 提供了 Docker 启动方式。下面用单节点 Kafka 作为学习环境。
bash
docker run -d --name kafka -p 9092:9092 apache/kafka:4.3.1创建 Topic:
bash
docker exec -it kafka /opt/kafka/bin/kafka-topics.sh \
--create \
--topic order.created \
--bootstrap-server localhost:9092查看 Topic:
bash
docker exec -it kafka /opt/kafka/bin/kafka-topics.sh \
--describe \
--topic order.created \
--bootstrap-server localhost:9092启动命令行生产者:
bash
docker exec -it kafka /opt/kafka/bin/kafka-console-producer.sh \
--topic order.created \
--bootstrap-server localhost:9092输入一行就是一条消息:
text
orderId=1001,userId=7,amount=99启动命令行消费者:
bash
docker exec -it kafka /opt/kafka/bin/kafka-console-consumer.sh \
--topic order.created \
--from-beginning \
--bootstrap-server localhost:9092如果能看到刚才输入的内容,说明最小链路已经跑通。
mermaid
flowchart TD
A["创建 Topic"] --> B["Console Producer 写消息"]
B --> C["Broker 持久化到 Partition"]
C --> D["Console Consumer 从头读取"]Spring Boot 依赖
使用 Spring Boot 项目时,一般引入 spring-kafka:
xml
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>application.yml
下面是学习版配置,重点是把 key 当字符串、value 当 JSON。
yaml
spring:
kafka:
bootstrap-servers: localhost:9092
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
acks: all
properties:
enable.idempotence: true
consumer:
group-id: order-demo-service
auto-offset-reset: earliest
enable-auto-commit: false
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
properties:
spring.json.trusted.packages: "*"
listener:
ack-mode: manual这里先解释几个容易踩坑的参数:
| 配置 | 为什么要配 | 不配可能怎样 |
|---|---|---|
acks: all | 等待足够副本确认 | Broker 故障时更容易丢消息 |
enable.idempotence: true | 开启生产者幂等 | 发送重试可能产生重复写入 |
enable-auto-commit: false | 禁止自动提交 offset | 业务没处理完就提交,宕机后可能漏消费 |
ack-mode: manual | 业务成功后手动确认 | 无法精确控制消费进度 |
auto-offset-reset: earliest | 没有旧 offset 时从最早开始读 | 学习环境中容易看不到历史消息 |
事件对象
事件对象建议表达“已经发生的事实”,不要表达“让别人做什么”的命令。
java
public record OrderCreatedEvent(
Long orderId,
Long userId,
Long amount,
Long createdAt
) {
}生产者代码
java
@Service
public class OrderKafkaProducer {
private final KafkaTemplate<String, OrderCreatedEvent> kafkaTemplate;
public OrderKafkaProducer(KafkaTemplate<String, OrderCreatedEvent> kafkaTemplate) {
this.kafkaTemplate = kafkaTemplate;
}
public CompletableFuture<SendResult<String, OrderCreatedEvent>> sendOrderCreated(OrderCreatedEvent event) {
String key = String.valueOf(event.orderId());
return kafkaTemplate.send("order.created", key, event);
}
}为什么 key 用 orderId?
- 同一个订单的事件会进入同一个分区。
- 同一个分区内 Kafka 保证顺序。
- 消费端处理订单状态流转时更安全。
如果不设置 key,同一个订单的“创建、支付、取消”可能进入不同分区,消费者并行处理时就可能先处理取消,再处理创建。
Controller 测试入口
java
@RestController
@RequestMapping("/orders")
public class OrderController {
private final OrderKafkaProducer producer;
public OrderController(OrderKafkaProducer producer) {
this.producer = producer;
}
@PostMapping("/{orderId}/events")
public String send(@PathVariable Long orderId) {
OrderCreatedEvent event = new OrderCreatedEvent(
orderId,
7L,
9900L,
System.currentTimeMillis()
);
producer.sendOrderCreated(event);
return "sent";
}
}调用:
bash
curl -X POST http://localhost:8080/orders/1001/events消费者代码
java
@Component
public class OrderKafkaConsumer {
@KafkaListener(topics = "order.created", groupId = "order-demo-service")
public void onMessage(OrderCreatedEvent event, Acknowledgment ack) {
try {
if (alreadyProcessed(event.orderId())) {
ack.acknowledge();
return;
}
handleBusiness(event);
markProcessed(event.orderId());
ack.acknowledge();
} catch (Exception ex) {
throw ex;
}
}
private boolean alreadyProcessed(Long orderId) {
return false;
}
private void handleBusiness(OrderCreatedEvent event) {
System.out.println("处理订单事件:" + event);
}
private void markProcessed(Long orderId) {
}
}这个消费者有两个重点:
- 先处理业务,再提交 offset。
- 处理前做幂等判断,处理后记录已处理状态。
生产消费流程
mermaid
flowchart TD
A["Controller 接收请求"] --> B["构造 OrderCreatedEvent"]
B --> C["KafkaTemplate 发送"]
C --> D["Kafka 写入 order.created"]
D --> E["消费者拉取消息"]
E --> F{"是否已经处理过?"}
F -->|"是"| G["直接提交 Offset"]
F -->|"否"| H["执行业务处理"]
H --> I["记录幂等标记"]
I --> J["手动提交 Offset"]为什么消费端需要幂等
Kafka 可能重复投递,常见原因包括:
- 消费者业务处理成功,但提交 offset 前宕机。
- 提交 offset 请求超时,消费者重试处理。
- Rebalance 后分区被分配给另一个消费者。
- 生产者重试或上游重复发送。
所以消费逻辑不能写成“收到就直接扣库存、加积分、发短信”。应该用业务唯一键防重复,例如订单 ID、事件 ID、流水号。
常见幂等实现:
sql
create table message_consume_log (
id bigint primary key auto_increment,
consumer_group varchar(128) not null,
message_key varchar(128) not null,
created_at datetime not null,
unique key uk_group_key (consumer_group, message_key)
);处理时先插入消费记录,唯一索引冲突就说明处理过:
java
@Transactional
public void consume(OrderCreatedEvent event) {
boolean inserted = tryInsertConsumeLog("order-demo-service", String.valueOf(event.orderId()));
if (!inserted) {
return;
}
handleBusiness(event);
}常见错误
| 错误 | 表面现象 | 根因 |
|---|---|---|
| Topic 不存在 | 发送失败或自动创建出错误 Topic | 没有统一管理 Topic |
| 不设置 key | 同一业务对象顺序混乱 | 消息分散到不同分区 |
| 自动提交 offset | 偶发漏消费 | offset 先提交,业务后失败 |
| 消费端不幂等 | 数据重复、短信重复、库存重复扣减 | 至少一次投递导致重复处理 |
| 所有服务用同一个 groupId | 只有一个服务收到消息 | 同组是竞争消费,不同组才是广播 |
小结
Kafka Demo 不难,难的是正确理解每一步的责任边界:Producer 负责可靠写入,Broker 负责保存和复制,Consumer 负责处理业务和提交 offset,业务系统自己负责幂等。
