Skip to content

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,业务系统自己负责幂等。