跳到正文
SL Blog 技术探索 · 工程实践 · AI 时代思考
返回
📨 Kafka 3.9 教程 · 2 / 10 查看系列简介 →

Kafka 3.9 教程(一):快速入门

文章目录
  1. 目录
  2. 1.1 什么是事件流处理?
  3. 1.1.1 实时捕获
  4. 1.1.2 持久存储
  5. 1.1.3 实时处理与响应
  6. 1.1.4 路由分发
  7. 1.2 Apache Kafka 是什么?
  8. 1.2.1 平台概述
  9. 1.2.2 核心特性
  10. 1.2.3 部署方式
  11. 1.3 Kafka 的核心概念
  12. 1.3.1 服务器端组件
  13. 1.3.2 客户端组件
  14. 1.3.3 核心术语详解
  15. 1.4 快速开始步骤
  16. 环境要求
  17. 1.4.1 第一步:下载并解压
  18. 1.4.2 第二步:启动 Kafka 环境
  19. 1.4.3 第三步:创建主题存储事件
  20. 1.4.4 第四步:向主题写入事件
  21. 1.4.5 第五步:读取事件
  22. 1.4.6 第六步:使用 Kafka Connect 导入/导出数据
  23. 1.4.7 第七步:使用 Kafka Streams 处理数据
  24. 1.4.8 第八步:停止 Kafka 环境
  25. 1.5 Kafka 的典型使用场景
  26. 1.5.1 消息代理(Messaging)
  27. 1.5.2 网站活动追踪(Website Activity Tracking)
  28. 1.5.3 监控指标(Metrics)
  29. 1.5.4 日志聚合(Log Aggregation)
  30. 1.5.5 流处理(Stream Processing)
  31. 1.5.6 事件溯源(Event Sourcing)
  32. 1.5.7 提交日志(Commit Log)
  33. 附录:常用命令速查
  34. 主题管理
  35. 生产消费
  36. 消费者组
  37. 下一步学习

📌 本文是「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(分区)

主题的物理分片,实现并行处理和水平扩展。

分区的作用:

  1. 并行处理:多个分区可被多个消费者并行消费
  2. 负载均衡:生产者写入分散到多个分区
  3. 顺序保证:同一分区内的消息有序
  4. 扩展性:增加分区数可提高吞吐量

分区与消费者关系:

  • 一个消费者组中,每个分区只能被一个消费者消费
  • 消费者数不应超过分区数(多余的消费者会空闲)
  • 分区数决定了最大并行度

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

下一步学习


RAG 智能问答

针对本文继续提问:《Kafka 3.9 教程(一):快速入门》