首页 / Spring Boot 入门教程 / 消息队列:Kafka

Spring Boot 入门教程

消息队列:Kafka

本教程共 48 篇 · 第 48 篇 · 更新于 2026-08-13 · 约 5 分钟阅读

Kafka消息队列KafkaTemplate@KafkaListener生产者消费者JSON 序列化Spring Kafka

本节目标:理解 Kafka 的基本概念,学会在 Spring Boot 里集成 Kafka,写生产者发消息、消费者收消息,并处理 JSON 消息。

Kafka 是什么

Kafka 是一个分布式消息系统,核心是「发布-订阅」。可以把 Kafka 想成邮局:发件人(生产者)把信件投进邮局,收件人(消费者)从邮局取信。两边互不认识,也不用同时在线。

几个必须懂的概念:

  • Broker:Kafka 服务器,负责存消息。
  • Topic:消息的分类,类似信件上的主题标签。
  • Partition:Topic 拆成多个分区,消息按分区存储,支持并行读写。
  • Offset:消息在分区里的编号,消费者靠它记录读到哪了。
  • Consumer Group:一组消费者共同消费一个 Topic。一条消息只会发给组里的一个成员;不同组各拿一份。

Kafka 和 RabbitMQ 的定位不同。Kafka 吞吐量极高,擅长海量日志、事件流、数据管道;RabbitMQ 更轻,适合普通业务消息。选型看场景,没有谁全面胜出。

添加依赖与配置

Spring Boot 通过 spring-kafka 项目提供自动配置,官方文档里加的就是这个依赖:

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

Kafka 没有专门的 spring-boot-starter-kafkaspring-kafka 本身就是 starter 的角色,版本由 Spring Boot 统一管理。

配置集中在 spring.kafka.*

spring:
  kafka:
    bootstrap-servers: "localhost:9092"   # Broker 地址
    consumer:
      group-id: "myGroup"                 # 消费者组
Warning

网上很多老教程写 localhost:2181,那是 ZooKeeper 的端口,早就过时了。Kafka 3.x 起用 KRaft 模式,不再需要 ZooKeeper,直连 9092 即可。

生产者:KafkaTemplate

KafkaTemplate 由自动配置创建,直接注入使用。它的 send 方法把消息发到指定 Topic:

package com.example.demo.service;

import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Service;

@Service
public class MessageProducer {

    private final KafkaTemplate<String, String> kafkaTemplate;

    public MessageProducer(KafkaTemplate<String, String> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    public void send(String topic, String message) {
        // 第一个参数是 Topic,第二个是消息内容
        kafkaTemplate.send(topic, message);
    }

    public void sendWithKey(String topic, String key, String message) {
        // 带 key 发送:相同 key 的消息会进同一个分区,保证顺序
        kafkaTemplate.send(topic, key, message);
    }
}

send 返回 CompletableFuture,需要确认发送结果时可以异步等待。

Topic 不会自动创建。想在启动时建好,声明一个 NewTopic Bean:

package com.example.demo.config;

import org.apache.kafka.clients.admin.NewTopic;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration(proxyBeanMethods = false)
public class KafkaTopicConfig {

    @Bean
    public NewTopic orderTopic() {
        // 名字、分区数、副本数
        return new NewTopic("orders", 3, (short) 1);
    }
}

Topic 已存在时这个 Bean 会被忽略,重复声明不会报错。

消费者:@KafkaListener

消费消息比发送更简单。在任意 Bean 的方法上标 @KafkaListener,指定 Topic 即可:

package com.example.demo.listener;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;

@Component
public class OrderListener {

    // 直接收消息内容
    @KafkaListener(topics = "orders")
    public void onMessage(String message) {
        System.out.println("收到订单消息:" + message);
    }

    // 也可以收完整的 ConsumerRecord,拿 key、分区、偏移量
    @KafkaListener(topics = "orders")
    public void onRecord(ConsumerRecord<String, String> record) {
        System.out.println("key=" + record.key()
                + ",分区=" + record.partition()
                + ",offset=" + record.offset()
                + ",内容=" + record.value());
    }
}

消费组从配置里读:spring.kafka.consumer.group-id。同一个组内,一条消息只被消费一次;想要多个系统各拿一份,用不同的组名。

监听容器(负责拉取消息、调用监听方法)由自动配置创建,参数在 spring.kafka.listener.* 下调整,比如并发数和提交方式:

spring:
  kafka:
    listener:
      concurrency: 3   # 并发消费线程数,别超过分区数
Tip

消费者默认从最新消息开始读(auto.offset.reset=latest)。想消费历史消息,改成 earliest;想手动指定位置,在监听方法里用 @Header(KafkaHeaders.OFFSET) 拿偏移量。

默认情况下,消息处理成功后偏移量自动提交,这叫自动确认。业务上要求「处理成功才提交」时,改成手动确认:

spring:
  kafka:
    listener:
      ack-mode: manual   # 手动提交偏移量
@Component
public class OrderListener {

    @KafkaListener(topics = "orders")
    public void onOrder(Order order, Acknowledgment ack) {
        try {
            System.out.println("处理订单:" + order.id());
            ack.acknowledge(); // 处理成功后才提交偏移量
        } catch (Exception e) {
            // 不确认,消息会重新投递
        }
    }
}

手动确认的代价是代码变多,换来的是「不丢消息」的承诺:处理中途崩溃,偏移量没提交,重启后消息还会再消费一遍。

JSON 消息

业务消息很少是纯字符串,更多是对象。spring-kafka 内置了 Jackson 序列化器,配置一下就能收发 JSON:

spring:
  kafka:
    producer:
      value-serializer: org.springframework.kafka.support.serializer.JacksonJsonSerializer
    consumer:
      value-deserializer: org.springframework.kafka.support.serializer.JacksonJsonDeserializer
      properties:
        # 反序列化成哪个类型
        "[spring.json.value.default.type]": com.example.demo.Order
        # 允许反序列化的包,防止反序列化攻击
        "[spring.json.trusted.packages]": com.example.demo

定义消息对象:

package com.example.demo;

import java.math.BigDecimal;

public record Order(Long id, String item, BigDecimal amount) {
}

生产者直接发对象:

@Service
public class OrderProducer {

    private final KafkaTemplate<String, Order> kafkaTemplate;

    public OrderProducer(KafkaTemplate<String, Order> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    public void createOrder(Order order) {
        kafkaTemplate.send("orders", order);
    }
}

消费者方法参数直接写对象类型:

@Component
public class OrderListener {

    @KafkaListener(topics = "orders")
    public void onOrder(Order order) {
        System.out.println("处理订单:" + order.id() + "," + order.item());
    }
}
Warning

spring.json.trusted.packages 一定要配。Jackson 反序列化不认识的包默认拒绝,这是防止恶意消息的安全机制,别图省事配成 *

幂等与顺序

Kafka 的投递语义是「至少一次」:消息可能被重复消费。比如消费者处理完还没来得及提交偏移量就崩溃,重启后同一条消息会再来一遍。所以消费端代码要设计成幂等的——同一个订单处理两次,结果和一次一样。常见做法是用业务主键去重:处理前查一下记录,存在就跳过。

顺序性也有讲究:同一个 key 的消息保证进同一个分区,分区内严格有序。想保证「同一个用户的订单消息按序处理」,发送时用用户 ID 做 key;消费端并发数不能超过分区数,否则一个分区被多个线程抢着消费,顺序就乱了。

事务与高级话题

消息发送和数据库操作要一起成功、一起失败时,用 Kafka 事务。配置事务前缀后,Spring Boot 自动创建 KafkaTransactionManager

spring:
  kafka:
    producer:
      transaction-id-prefix: "tx-"   # 声明使用事务

配合 @Transactional 注解,方法内数据库写和 Kafka 发送就绑进同一个事务。

Kafka Streams 是另一座大山:用 @EnableKafkaStreams 开启,通过 StreamsBuilder 对消息做流式处理(过滤、聚合、Join)。它处理的是「持续流动的数据」,适合实时统计场景。

本地跑起来

本地开发需要先启动一个 Kafka。用 Docker Compose 最省事:

services:
  kafka:
    image: apache/kafka:latest
    ports:
      - "9092:9092"
    environment:
      KAFKA_NODE_ID: 1
      KAFKA_PROCESS_ROLES: "broker,controller"
      KAFKA_LISTENERS: "PLAINTEXT://:9092,CONTROLLER://:9093"
      KAFKA_ADVERTISED_LISTENERS: "PLAINTEXT://localhost:9092"
      KAFKA_CONTROLLER_LISTENER_NAMES: "CONTROLLER"
      KAFKA_CONTROLLER_QUORUM_VOTERS: "1@localhost:9093"

启动后运行应用,用生产者接口发消息,观察消费者日志,一套本地消息链路就通了。

写测试时可以用 @EmbeddedKafka 注解,它会在测试里启动一个内嵌 Kafka,不依赖外部环境:

package com.example.demo;

import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.kafka.test.context.EmbeddedKafka;

@SpringBootTest
@EmbeddedKafka(topics = "orders", bootstrapServersProperty = "spring.kafka.bootstrap-servers")
class OrderListenerTest {
    // 测试代码
}

bootstrapServersProperty 把内嵌 Broker 的地址映射进 Spring Boot 配置,自动配置无需改动。

小结

Kafka 集成是 Spring Boot 里「配置最少、威力最大」的模块之一:加 spring-kafka 依赖,写两行配置,注入 KafkaTemplate 发消息、加 @KafkaListener 收消息。进阶方向记住三条:JSON 序列化要配 trusted packages、事务靠 transaction-id-prefix、本地调试用 Docker Compose 或 @EmbeddedKafka

上一篇
任务调度与异步
下一篇
已经是最后一篇啦