📌 本文是「Kafka 3.9 教程」系列第 一 篇 · 系列目录
本文档基于 Apache Kafka 3.9 官方文档整理
目录
1.1 什么是事件流处理?
事件流处理是现代分布式系统的核心基础设施,它相当于人体的中枢神经系统。在当今「永远在线」的世界中,企业越来越依赖软件定义和自动化,事件流处理正是支撑这一趋势的技术基础。
从技术角度来说,事件流处理包括四个核心步骤:
1.1.1 实时捕获
从以下事件源实时捕获数据:
- 数据库(数据变更)
- 传感器(物联网设备数据)
- 移动设备(地理位置更新)
- 云服务(API 调用)
- 软件应用程序(日志、事件)
这些数据以事件流的形式被捕获和处理。
1.1.2 持久存储
事件流被持久化存储,以便:
- 后续检索和分析
- 支持时间回溯查询
- 保证数据不丢失
Kafka 提供了高持久性的存储机制,确保数据安全。
1.1.3 实时处理与响应
对事件流进行:
- 实时处理:毫秒级响应业务逻辑
- 追溯处理:对历史数据进行批量分析
- 复杂计算:聚合、转换、过滤、连接等操作
1.1.4 路由分发
根据业务规则,将事件流路由到:
- 不同的下游系统
- 不同的处理管道
- 不同的存储位置
1.2 Apache Kafka 是什么?
1.2.1 平台概述
Apache Kafka 是一个分布式事件流平台,将三个关键能力结合在一起:
| 能力 | 说明 |
|---|---|
| 发布和订阅 | 像消息队列或企业消息系统一样发布和订阅事件流 |
| 持久存储 | 按时间顺序持久化存储事件流,支持重放 |
| 实时处理 | 事件发生时立即进行处理,无需等待批量 |
1.2.2 核心特性
Kafka 以分布式、高度可扩展、弹性、容错和安全的方式提供上述功能:
- 分布式架构:支持跨数据中心部署
- 高可扩展性:支持数千个 Broker,百万级消息/秒
- 弹性设计:支持水平扩展,故障自动恢复
- 容错机制:多副本保证数据不丢失
- 安全可控:支持 SSL/TLS、SASL、ACL 等安全机制
1.2.3 部署方式
Kafka 可以部署在多种环境中:
| 部署方式 | 说明 |
|---|---|
| 裸机硬件 | 高性能、高吞吐量的传统部署方式 |
| 虚拟机 | 虚拟化环境,资源隔离 |
| 容器 | Docker、Kubernetes 编排部署 |
| 云环境 | AWS、Azure、GCP 等公有云 |
1.3 Kafka 的核心概念
1.3.1 服务器端组件
Broker(代理节点)
Kafka 集群中的服务器,负责:
- 接收和存储生产者发送的消息
- 处理消费者的消息读取请求
- 管理分区的 Leader/Follower 关系
特点:
- 每个 Broker 可以同时是多个分区的 Leader
- Broker 之间通过 Zookeeper 或 KRaft 进行协调
- Broker 故障时,自动进行 Leader 选举
Kafka Connect
运行在服务器上的数据集成框架:
- 持续导入数据:将外部系统数据导入 Kafka
- 持续导出数据:将 Kafka 数据导出到外部系统
1.3.2 客户端组件
Producer(生产者)
发布(写入)事件到 Kafka 主题的应用程序。
主要职责:
- 将业务数据封装为消息
- 选择目标分区
- 批量发送消息以提高效率
- 处理发送失败和重试
Consumer(消费者)
订阅(读取和处理)事件的应用程序。
主要职责:
- 从指定主题拉取消息
- 按分区顺序处理消息
- 管理消费位置(offset)
- 支持消费者组实现负载均衡
支持的语言
Kafka 提供多种语言的官方客户端:
- Java/Scala:完整功能支持,包含 Kafka Streams 库
- Go:高性能客户端
- Python:confluent-kafka-python 等
- C/C++:librdkafka
- REST API:通过代理访问
1.3.3 核心术语详解
Event(事件)
记录「发生了什么」的基本数据单元。
事件结构:
{
"key": "user-123",
"value": {
"action": "purchase",
"amount": 99.99,
"product": "book"
},
"timestamp": 1699900800000,
"headers": {
"content-type": "application/json"
}
}
| 字段 | 说明 | 可选 |
|---|---|---|
| key | 用于分区路由,相同 key 到相同分区 | 是 |
| value | 实际消息内容 | 是(但通常必须有) |
| timestamp | 事件时间或处理时间 | 否(自动添加) |
| headers | 元数据键值对 | 是 |
Topic(主题)
事件的分类容器,类似于文件系统中的文件夹。
特点:
- 多生产者:一个主题可以有多个生产者同时写入
- 多消费者:一个主题可以有多个消费者组订阅
- 持久存储:消息按配置保留,不会自动删除
- 按需读取:消费者可以多次读取历史消息
# 创建主题示例
kafka-topics.sh --create \
--topic user-events \
--bootstrap-server localhost:9092 \
--partitions 3 \
--replication-factor 2
Partition(分区)
主题的物理分片,实现并行处理和水平扩展。
分区的作用:
- 并行处理:多个分区可被多个消费者并行消费
- 负载均衡:生产者写入分散到多个分区
- 顺序保证:同一分区内的消息有序
- 扩展性:增加分区数可提高吞吐量
分区与消费者关系:
- 一个消费者组中,每个分区只能被一个消费者消费
- 消费者数不应超过分区数(多余的消费者会空闲)
- 分区数决定了最大并行度
Replication(副本)
分区的多份拷贝,实现容错和高可用。
核心概念:
- Leader:主副本,处理所有读写请求
- Follower:从副本,从 Leader 同步数据
- ISR(In-Sync Replicas):与 Leader 保持同步的副本集合
副本因子配置:
# Broker 默认配置
default.replication.factor=3
# 创建主题时指定
kafka-topics.sh --create \
--topic my-topic \
--replication-factor 3
1.4 快速开始步骤
环境要求
- Java 8+:必须安装 JDK 或 JRE
- 操作系统:Linux/macOS/Windows
1.4.1 第一步:下载并解压
# 下载 Kafka 3.9.1(推荐 Scala 2.13 版本)
wget https://archive.apache.org/dist/kafka/3.9.1/kafka_2.13-3.9.1.tgz
# 解压
tar -xzf kafka_2.13-3.9.1.tgz
cd kafka_2.13-3.9.1
1.4.2 第二步:启动 Kafka 环境
Kafka 支持两种启动模式:KRaft 模式(推荐)和 ZooKeeper 模式。
KRaft 模式(推荐,无需 ZooKeeper)
# 1. 生成集群 UUID
KAFKA_CLUSTER_ID="$(bin/kafka-storage.sh random-uuid)"
echo $KAFKA_CLUSTER_ID
# 2. 格式化日志目录
bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/kraft/server.properties
# 3. 启动 Kafka 服务器
bin/kafka-server-start.sh config/kraft/server.properties
Docker 方式(JVM 版本)
# 拉取镜像
docker pull apache/kafka:3.9.1
# 启动容器
docker run -p 9092:9092 apache/kafka:3.9.1
Docker 方式(Native 版本,实验性)
# 拉取原生镜像(基于 GraalVM,更轻量但仅用于测试)
docker pull apache/kafka-native:3.9.1
docker run -p 9092:9092 apache/kafka-native:3.9.1
ZooKeeper 模式(旧版本兼容)
# 1. 启动 ZooKeeper
bin/zookeeper-server-start.sh config/zookeeper.properties
# 2. 启动 Kafka Broker(新开终端)
bin/kafka-server-start.sh config/server.properties
1.4.3 第三步:创建主题存储事件
# 创建主题
kafka-topics.sh --create \
--topic quickstart-events \
--bootstrap-server localhost:9092
# 查看主题详情
kafka-topics.sh --describe \
--topic quickstart-events \
--bootstrap-server localhost:9092
输出示例:
Topic: quickstart-events TopicId: xxx PartitionCount: 1 ReplicationFactor: 1
Topic: quickstart-events Partition: 0 Leader: 0 Replicas: 0 Isr: 0
1.4.4 第四步:向主题写入事件
使用控制台生产者发送消息:
kafka-console-producer.sh \
--topic quickstart-events \
--bootstrap-server localhost:9092
# 输入消息(每行一个事件)
> 这是我的第一条消息
> 这是我的第二条消息
> {"user": "张三", "action": "登录"}
1.4.5 第五步:读取事件
使用控制台消费者读取消息:
# 从头开始消费(读取所有历史消息)
kafka-console-consumer.sh \
--topic quickstart-events \
--from-beginning \
--bootstrap-server localhost:9092
# 实时消费(只消费新消息)
kafka-console-consumer.sh \
--topic quickstart-events \
--bootstrap-server localhost:9092
输出:
这是我的第一条消息
这是我的第二条消息
{"user": "张三", "action": "登录"}
1.4.6 第六步:使用 Kafka Connect 导入/导出数据
配置 Kafka Connect
编辑 config/connect-standalone.properties,添加:
plugin.path=libs/connect-file-3.9.1.jar
创建测试数据文件
echo -e "foo\nbar" > test.txt
启动连接器(独立模式)
# 启动两个连接器:源连接器和目标连接器
bin/connect-standalone.sh \
config/connect-standalone.properties \
config/connect-file-source.properties \
config/connect-file-sink.properties
验证数据流动
# 查看目标文件内容
cat test.sink.txt
# 输出:foo, bar
# 或使用控制台消费者查看
kafka-console-consumer.sh \
--topic connect-test \
--from-beginning \
--bootstrap-server localhost:9092
1.4.7 第七步:使用 Kafka Streams 处理数据
Kafka Streams 是内置的流处理库,可以对事件进行实时处理。
WordCount 示例代码
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.kstream.KStream;
import org.apache.kafka.streams.kstream.KTable;
import org.apache.kafka.streams.kstream.Produced;
import java.util.Arrays;
import java.util.Properties;
public class WordCountDemo {
public static void main(String[] args) {
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount-demo");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName());
StreamsBuilder builder = new StreamsBuilder();
KStream<String, String> textLines = builder.stream("quickstart-events");
KTable<String, Long> wordCounts = textLines
.flatMapValues(line -> Arrays.asList(line.toLowerCase().split(" ")))
.groupBy((key, word) -> word)
.count();
wordCounts.toStream().to("wordcount-output", Produced.with(Serdes.String(), Serdes.Long()));
KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();
}
}
1.4.8 第八步:停止 Kafka 环境
# 停止所有服务
# Linux/macOS
rm -rf /tmp/kafka-logs /tmp/zookeeper /tmp/kraft-combined-logs
# Windows(手动删除目录)
# C:\tmp\kafka-logs
# C:\tmp\zookeeper
1.5 Kafka 的典型使用场景
1.5.1 消息代理(Messaging)
Kafka 可作为传统消息中间件(如 ActiveMQ、RabbitMQ)的替代品。
适用场景:
- 系统间解耦
- 异步任务处理
- 流量削峰填谷
- 任务队列
优势:
- 更高的吞吐量
- 内置分区和副本
- 更好的容错能力
1.5.2 网站活动追踪(Website Activity Tracking)
Kafka 最初就是为此场景设计的。
典型流程:
用户浏览页面 → 发送事件到 Kafka → 实时处理 → 存储到数据仓库
事件类型:
- 页面浏览(Page View)
- 搜索查询(Search)
- 点击行为(Click)
- 表单提交(Form Submit)
1.5.3 监控指标(Metrics)
收集和聚合分布式应用的运营数据。
应用场景:
- 性能监控
- 错误追踪
- 服务健康检查
- 业务指标统计
1.5.4 日志聚合(Log Aggregation)
集中收集和处理服务器日志。
对比传统方案:
| 特性 | 传统日志收集 | Kafka 日志聚合 |
|---|---|---|
| 延迟 | 较高 | 低延迟实时 |
| 可靠性 | 一般 | 高(副本机制) |
| 吞吐量 | 中等 | 高 |
| 扩展性 | 有限 | 水平扩展 |
1.5.5 流处理(Stream Processing)
多阶段数据处理管道。
典型架构:
原始数据 → Kafka → 处理阶段1 → 新主题 → 处理阶段2 → ... → 最终结果
处理框架:
- Kafka Streams:内置,轻量级
- Apache Flink:分布式流处理框架
- Apache Storm:实时计算系统
- Spark Streaming:微批处理
1.5.6 事件溯源(Event Sourcing)
将状态变更记录为时间有序的事件序列。
特点:
- 完整的事件历史
- 支持时间回溯
- 便于审计和调试
1.5.7 提交日志(Commit Log)
作为分布式系统的外部提交日志。
应用场景:
- 数据库变更捕获(CDC)
- 分布式系统协调
- 事件驱动的微服务架构
- 数据复制同步
附录:常用命令速查
主题管理
# 创建主题
kafka-topics.sh --create --topic my-topic --bootstrap-server localhost:9092
# 列出所有主题
kafka-topics.sh --list --bootstrap-server localhost:9092
# 查看主题详情
kafka-topics.sh --describe --topic my-topic --bootstrap-server localhost:9092
# 修改主题配置
kafka-configs.sh --alter --topic my-topic --add-config retention.ms=86400000 --bootstrap-server localhost:9092
# 删除主题
kafka-topics.sh --delete --topic my-topic --bootstrap-server localhost:9092
生产消费
# 控制台生产者
kafka-console-producer.sh --topic my-topic --bootstrap-server localhost:9092
# 控制台消费者
kafka-console-consumer.sh --topic my-topic --from-beginning --bootstrap-server localhost:9092
# 带密钥的消费
kafka-console-consumer.sh --topic my-topic --property print.key=true --bootstrap-server localhost:9092
消费者组
# 列出消费者组
kafka-consumer-groups.sh --list --bootstrap-server localhost:9092
# 查看消费者组详情
kafka-consumer-groups.sh --describe --group my-group --bootstrap-server localhost:9092
# 重置偏移量
kafka-consumer-groups.sh --reset-offsets --group my-group --topic my-topic --to-earliest --execute --bootstrap-server localhost:9092