Kafka 从入门到实战:Docker 部署与 Spring Boot 集成

在分布式系统中,消息中间件是实现服务解耦、异步通信的核心组件,而 Kafka 凭借其高吞吐、高可靠、可扩展的特性,成为日志收集、流式处理、消息队列等场景的首选方案。

一、认识 Kafka:是什么、从哪来、有啥本事

1.1 Kafka 是什么?

Kafka 是由 Apache 基金会维护的分布式流处理平台,最初由 LinkedIn 开发并于 2011 年开源。它本质上是一个高吞吐量的分布式消息队列,支持海量数据的实时采集、存储与分发,既能作为消息中间件实现服务间通信,也能作为流处理引擎处理实时数据(如用户行为分析、日志聚合)。

简单来说,Kafka 就像一个“分布式邮局”:生产者(发送消息的服务)是“寄信人”,消费者(接收消息的服务)是“收信人”,而 Kafka 则负责将“信件”(消息)高效、安全地投递到目标“邮箱”(Topic)。

1.2 Kafka 来源与归属

  • 起源:2010 年,LinkedIn 为解决内部日志收集与服务间通信的痛点,开发了 Kafka 原型,初衷是替代传统消息中间件(如 ActiveMQ),满足高吞吐场景需求。
  • 开源与归属:2011 年 Kafka 开源,2012 年加入 Apache 孵化器,2014 年成为 Apache 顶级项目,目前由 Apache 基金会维护,社区活跃,版本迭代稳定(最新稳定版为 3.7.x,本文使用 3.6.1 版本)。
  • 企业应用:除了 LinkedIn,Kafka 已被阿里、腾讯、字节、Netflix、Uber 等国内外大厂广泛采用,是大数据生态(Hadoop、Spark、Flink)中不可或缺的组件。

1.3 Kafka 核心特性

Kafka 能在众多消息中间件中脱颖而出,源于其四大核心特性:

  • 高并发:单机吞吐量可达 10 万级/秒,通过分区(Partition)机制实现并行读写,支持海量数据的实时传输(如双 11 订单日志处理)。
  • 容错性:基于副本(Replica)机制,每个分区可配置多个副本,当主副本(Leader)宕机时,从副本(Follower)自动切换,保证数据不丢失、服务不中断。
  • 可扩展性:支持横向扩展 Broker 节点,只需新增服务器并配置到集群,即可提升 Kafka 存储与处理能力,无需停机维护。
  • 消息顺序性:在单个分区内,消息严格按照发送顺序存储与消费(分区间无顺序保证),满足订单支付、日志时序等对顺序敏感的场景。

1.4 Kafka 典型应用场景

Kafka 的特性决定了其广泛的应用范围,主要包括四类场景:

  • 消息队列:实现服务解耦(如订单服务与库存服务通过 Kafka 通信,避免直接调用依赖)、异步通信(如用户注册后,异步发送短信/邮件通知)。
  • 日志收集:集中采集分布式系统中的日志(如应用日志、服务器日志),发送到 Kafka 后,再由 Flink/Spark 处理或存储到 Elasticsearch 供查询。
  • 运行指标监控:实时采集服务的运行指标(如 CPU 使用率、接口响应时间),通过 Kafka 传输到监控平台(如 Prometheus),实现异常告警。
  • 流式处理:作为流处理引擎(如 Kafka Streams、Flink)的数据源,实时处理数据(如实时计算用户订单金额、实时推荐)。

1.5 官方文档参考

二、Docker 搭建单机 Kafka 环境

Kafka 运行依赖 ZooKeeper(用于元数据管理,如 Topic 配置、Broker 节点信息),因此需同时部署 ZooKeeper 与 Kafka。本文将数据目录映射到宿主机 /usr/local 下,确保数据持久化,并明确指定版本避免兼容性问题。

2.1 环境准备

  • 宿主机:Linux 系统(如 CentOS 7、Ubuntu 20.04),已安装 Docker(建议 20.10+ 版本)。
  • 网络:确保宿主机 2181 端口(ZooKeeper)、9092 端口(Kafka)未被占用,且能正常访问(如需外部访问,需开放对应端口)。

2.2 步骤 1:创建宿主机数据目录

为避免容器销毁后数据丢失,需将 ZooKeeper 与 Kafka 的数据目录映射到宿主机 /usr/local 下,并配置权限(生产环境可根据实际需求调整权限,此处简化为 777):

# 创建 ZooKeeper 数据目录
mkdir -p /usr/local/zookeeper/data
# 创建 Kafka 数据目录
mkdir -p /usr/local/kafka/data
# 赋予目录读写权限(简化配置,生产环境需按需调整)
chmod -R 777 /usr/local/zookeeper /usr/local/kafka

2.3 步骤 2:启动 ZooKeeper 容器

使用官方 ZooKeeper 镜像(版本 3.8.4,与 Kafka 3.6.1 兼容),映射 2181 端口与数据目录:

docker run -d --name zookeeper \
  -p 2181:2181 \
  -v /usr/local/zookeeper/data:/data \  # 宿主机目录:容器内目录(容器内默认数据目录为 /data)
  zookeeper:3.8.4  # 明确版本,避免拉取 latest 版本导致兼容性问题
  • 参数说明:
    • -d:后台运行容器;
    • --name zookeeper:指定容器名为 zookeeper,方便后续操作;
    • -p 2181:2181:映射宿主机 2181 端口到容器 2181 端口(ZooKeeper 默认端口);
    • -v /usr/local/zookeeper/data:/data:持久化 ZooKeeper 数据到宿主机。

启动后,可通过 docker ps 查看容器是否正常运行:

docker ps | grep zookeeper  # 若输出容器信息,说明启动成功

2.4 步骤 3:启动 Kafka 容器

使用 bitnami 维护的 Kafka 镜像(配置更友好,版本 3.6.1),关联 ZooKeeper 并映射数据目录:

docker run -d --name kafka \
  -p 9092:9092 \
  -v /usr/local/kafka/data:/bitnami/kafka/data \  # 映射 Kafka 数据目录(bitnami 镜像默认数据目录为 /bitnami/kafka/data)
  -e KAFKA_BROKER_ID=1 \  # Broker 唯一标识(集群环境需不同,单机设为 1 即可)
  -e KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181 \  # 连接 ZooKeeper(容器名:端口,因 --link 可直接用容器名访问)
  -e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://<你的宿主机IP>:9092 \  # 关键!替换为宿主机实际IP(如 192.168.1.100)
  -e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092 \  # 允许容器内所有地址监听 9092 端口
  --link zookeeper \  # 关联 zookeeper 容器,保证网络互通
  bitnami/kafka:3.6.1  # 明确 Kafka 版本
  • 关键注意点:
    • KAFKA_ADVERTISED_LISTENERS:必须替换为宿主机的实际 IP(如内网 IP 或公网 IP),否则 Spring Boot 应用无法连接 Kafka(因为 Kafka 会返回容器内地址,外部服务无法访问);
    • --link zookeeper:在单机环境中,通过该参数让 Kafka 容器能访问 ZooKeeper 容器,集群环境需用 Docker Network 替代。

启动后,验证 Kafka 容器状态:

docker ps | grep kafka  # 若输出容器信息,说明启动成功

2.5 步骤 4:验证 Kafka 环境

为确保 Kafka 能正常生产/消费消息,需进入容器创建 Topic 并测试:

1. 进入 Kafka 容器
docker exec -it kafka /bin/bash  # 进入 Kafka 容器终端
2. 创建 Topic(如 default-topic
# 执行 Kafka 自带的创建 Topic 脚本
/opt/bitnami/kafka/bin/kafka-topics.sh \
  --create \
  --topic default-topic \  # Topic 名称(需与 Spring Boot 配置一致)
  --bootstrap-server localhost:9092 \  # 连接本地 Kafka
  --partitions 1 \  # 分区数(单机设为 1 即可)
  --replication-factor 1  # 副本数(单机设为 1,集群需 >=2)

执行成功后,会输出 Created topic default-topic.

3. 测试生产者发送消息
# 启动生产者终端
/opt/bitnami/kafka/bin/kafka-console-producer.sh \
  --topic default-topic \
  --bootstrap-server localhost:9092

此时可输入任意消息(如 test kafka message),按回车发送。

4. 测试消费者接收消息

打开新的终端,再次进入 Kafka 容器,启动消费者:

# 启动消费者终端(--from-beginning 表示从最早消息开始消费)
/opt/bitnami/kafka/bin/kafka-console-consumer.sh \
  --topic default-topic \
  --bootstrap-server localhost:9092 \
  --from-beginning

若能接收到之前生产者发送的 test kafka message,说明 Kafka 环境搭建成功。

三、Spring Boot 集成 Kafka:实现生产者与消费者

接下来,我们基于 Spring Boot 2.7.x 版本,实现 Kafka 生产者(通过接口触发消息发送)与消费者(接收并处理消息),完成接口调用日志的采集。

3.1 项目结构

先明确项目的核心结构(基于 Maven),便于理解代码组织:

com.buffer.kafkademo
├── controller  # 接口层:暴露 HTTP 接口触发消息发送
│   └── KafkaDemoController.java
├── service     # 服务层:生产者与消费者逻辑
│   ├── KafkaProducerService.java  # 生产者:发送消息到 Kafka
│   └── KafkaConsumerService.java  # 消费者:从 Kafka 接收消息
├── config      # 配置层(可选,本文用 yml 配置)
└── KafkaDemoApplication.java      # 启动类

3.2 步骤 1:添加 Maven 依赖

pom.xml 中引入 Spring Boot Web(暴露接口)、Spring Kafka(集成 Kafka)、Lombok(简化日志代码):

<dependencies>
    <!-- Spring Boot Web:提供 HTTP 接口 -->
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-web</artifactId>
    </dependency>

    <!-- Spring Kafka:集成 Kafka 客户端 -->
    <dependency>
        <groupId>org.springframework.kafka</groupId>
        <artifactId>spring-kafka</artifactId>
        <!-- 版本由 Spring Boot 父依赖管理,无需手动指定 -->
    </dependency>
    <!-- Spring Boot 测试(可选) -->
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-test</artifactId>
        <scope>test</scope>
    </dependency>
</dependencies>

3.3 步骤 2:配置 Kafka(application.yml)

src/main/resources/application.yml 中配置 Kafka 连接信息、生产者/消费者参数(注意替换 bootstrap-servers 为你的宿主机 IP):

spring:
  kafka:
    # Kafka 服务地址(替换为你的宿主机 IP,如 192.168.1.100:9092)
    bootstrap-servers: 19.168.10.23:9092

    # 生产者配置
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer  # Key 序列化方式
      value-serializer: org.apache.kafka.common.serialization.StringSerializer  # Value 序列化方式
      acks: 1  # 消息确认机制:1 表示 Leader 副本接收成功即返回(平衡可靠性与性能)
      retries: 3  # 发送失败重试次数:避免网络抖动导致的消息丢失
      batch-size: 16384  # 批量发送大小:16KB(达到阈值后批量发送,提升吞吐量)
      linger-ms: 1  # 延迟发送时间:1ms(即使未达批量阈值,1ms 后也发送,避免延迟过高)

    # 消费者配置
    consumer:
      group-id: kafka-demo-group  # 消费者组 ID(同一组内消费者共同消费 Topic,避免重复消费)
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer  # Key 反序列化方式
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer  # Value 反序列化方式
      auto-offset-reset: earliest  # 偏移量重置策略:earliest 表示从最早消息开始消费(首次启动时)
      enable-auto-commit: false  # 关闭自动提交偏移量(配合手动提交,保证消息不重复消费)

    # 消费者监听器配置
    listener:
      ack-mode: manual_immediate  # 手动提交偏移量(消费成功后立即提交,避免消息丢失)
      concurrency: 1  # 消费者并发数(单机设为 1,集群可根据分区数调整)
  • 配置说明:
    • acks: 1:生产环境中,若需更高可靠性,可设为 all(所有副本接收成功才返回),但会牺牲部分性能;
    • enable-auto-commit: false + ack-mode: manual_immediate:手动提交偏移量,确保消息被消费成功后再提交,避免消息重复消费或丢失;
    • group-id:同一 Topic 的不同消费者组可重复消费消息(如 Topic 同时供日志存储和实时分析使用)。

3.4 步骤 3:实现 Kafka 生产者(KafkaProducerService)

生产者负责将接口调用信息发送到 Kafka 的 default-topic,通过 KafkaTemplate 简化发送逻辑,并添加回调处理发送成功/失败的情况:

package com.buffer.kafkademo.service;

import lombok.extern.slf4j.Slf4j;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Service;
import org.springframework.util.concurrent.ListenableFuture;
import org.springframework.util.concurrent.ListenableFutureCallback;

/**
 * Kafka 生产者服务:发送接口调用日志到 Kafka
 */
@Service
@Slf4j  // Lombok 注解,简化日志对象创建(无需手动 new Logger)
public class KafkaProducerService {

    // 注入 Spring Kafka 提供的 KafkaTemplate(已自动配置,无需手动初始化)
    private final KafkaTemplate<String, String> kafkaTemplate;

    // 构造器注入(Spring 推荐,避免循环依赖)
    public KafkaProducerService(KafkaTemplate<String, String> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    /**
     * 发送消息到指定 Topic
     * @param topic 消息主题(需与消费者监听的 Topic 一致)
     * @param message 消息内容(如接口调用日志)
     */
 public void sendRecordMessage(String topic, String message) {
        try {
            ListenableFuture<?> future = kafkaTemplate.send(topic, message);
            future.addCallback(new ListenableFutureCallback<Object>() {
                @Override
                public void onSuccess(Object result) {
                    logger.info("Message sent successfully to topic: {}, message: {}", topic, message);
                }

                @Override
                public void onFailure(Throwable ex) {
                    logger.error("Failed to send message to topic: {}, message: {}", topic, message, ex);
                }
            });
        } catch (Exception e) {
            logger.error("Error while sending message to topic: {}, message: {}", topic, message, e);
        }
    }
}

3.5 步骤 4:实现 Kafka 消费者(KafkaConsumerService)

消费者通过 @KafkaListener 注解监听 default-topic,接收生产者发送的消息并处理(如打印日志、存储到数据库),同时手动提交偏移量确保消息不重复消费:

package com.buffer.kafkademo.service;

import lombok.extern.slf4j.Slf4j;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.stereotype.Service;

/**
 * Kafka 消费者服务:监听 Topic 并处理消息
 */
@Service
@Slf4j
public class KafkaConsumerService {

    /**
     * 监听指定 Topic 的消息
     * @param message 接收到的消息内容
     * @param acknowledgment 偏移量提交对象(手动提交需用到)
     */
    @KafkaListener(topics = "default-topic", groupId = "kafka-demo-group")
    public void consumeMessage(String message, Acknowledgment acknowledgment) {
        try {
            // 1. 业务逻辑处理:此处模拟处理接口调用日志(实际场景可存储到 ES/MySQL)
            log.info("Received message from Kafka | Topic: default-topic, Content: {}", message);
            
            // 示例:若消息是 JSON 格式,可解析为对象(需引入 Jackson 依赖)
            // ApiLog apiLog = new ObjectMapper().readValue(message, ApiLog.class);
            // log.info("Parsed API Log | Path: {}, Time: {}", apiLog.getApiPath(), apiLog.getInvokeTime());

            // 2. 手动提交偏移量:确保消息处理成功后再提交,避免消息丢失
            acknowledgment.acknowledge();
            log.info("Offset committed successfully | Message processed: {}", message);
        } catch (Exception e) {
            // 处理失败:不提交偏移量(Kafka 会重新推送消息),同时打印异常
            log.error("Failed to process message | Content: {}", message, e);
            // 实际场景可添加重试次数限制(如超过 3 次失败则存入死信队列)
        }
    }
}
Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐