通用:Kafka 从入门到实战:Docker 部署与 Spring Boot 集成
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 官方文档参考
- Kafka 英文官方文档:Apache Kafka Documentation
- Docker 部署相关文档:Kafka Docker Guide
二、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 次失败则存入死信队列)
}
}
}
更多推荐



所有评论(0)