第一章:Kafka 基础概念与架构
1.1 什么是消息队列与 Kafka 的定位
| 概念名称 | 说明 | 注意事项 |
|---|---|---|
| 消息队列(Message Queue) | 一种应用程序间通信的中间件,通过发送和接收消息实现异步通信和解耦。 | 消息队列通常用于解耦系统、削峰填谷、异步处理等场景。 |
| 同步通信 | 调用方等待被调用方返回结果后才继续执行,如 HTTP 请求。 | 容易造成服务阻塞,系统耦合度高。 |
| 异步通信 | 发送方发送消息后无需等待接收方处理结果,提高系统响应速度和吞吐量。 | 需要处理消息丢失、顺序、幂等性等问题。 |
| Kafka 定位 | 分布式流处理平台,具备高吞吐、低延迟、可扩展、持久化存储等特点,适用于大数据场景。 | Kafka 不仅是消息队列,还可用于日志聚合、流处理、事件溯源等。 |
| 主要应用场景 | 日志收集、监控数据聚合、业务解耦、流式处理、消息广播等。 | 根据场景选择合适的消息中间件,Kafka 适合高吞吐、大数据量场景,而非低延迟微服务通信。 |
1.2 Kafka 核心组件介绍(Broker、Topic、Partition、Producer、Consumer、Consumer Group)
| 组件名称 | 说明 | 注意事项 |
|---|---|---|
| Broker | Kafka 集群中的一个服务器节点,负责存储和转发消息。 | 每个 Broker 有唯一 ID,多个 Broker 组成集群,提升可用性和吞吐量。 |
| Topic | 消息的主题分类,生产者将消息发送到特定 Topic,消费者订阅 Topic 接收消息。 | Topic 是逻辑概念,可划分为多个 Partition 以实现并行处理。 |
| Partition | Topic 的分区,每个 Partition 是一个有序、不可变的消息序列,分布在不同 Broker 上。 | 分区是 Kafka 实现水平扩展和并发处理的基础,每个 Partition 只能由一个 Consumer 消费。 |
| Producer | 生产者,向 Kafka 的 Topic 发送消息的应用程序。 | 生产者可指定消息发送到哪个 Partition,或由分区策略自动选择。 |
| Consumer | 消费者,从 Kafka 订阅 Topic 并拉取消息的应用程序。 | 消费者通过拉取(pull)方式获取消息,控制消费速度。 |
| Consumer Group | 消费者组,一组消费者的逻辑集合,共同消费一个或多个 Topic,实现负载均衡。 | 同一 Consumer Group 内的消费者共享位移(offset),每个 Partition 只能被组内一个消费者消费。 |
1.3 Kafka 架构原理与数据流向
| 概念名称 | 说明 | 注意事项 |
|---|---|---|
| 数据流向 | Producer → Topic (Partition) → Consumer (Consumer Group) | 消息按 Partition 有序存储,消费者组内消费者并行消费不同 Partition。 |
| Leader Partition | 每个 Partition 的主副本,负责处理所有读写请求。 | Leader 由 Controller Broker 选举产生,客户端直接与 Leader 交互。 |
| Follower Partition | Partition 的副本,从 Leader 同步数据,用于容灾。 | Follower 不处理客户端请求,仅被动复制数据。 |
| ISR(In-Sync Replicas) | 与 Leader 保持同步的副本集合,包含 Leader 和同步良好的 Follower。 | ISR 列表由 Broker 维护,若 Follower 落后过多将被移出 ISR。 |
| Controller Broker | 集群中的控制器,负责 Partition Leader 选举、Broker 上下线处理等元数据管理。 | 每个集群有且仅有一个 Controller,由 ZooKeeper 或 KRaft 协议选举产生。 |
| Replication Factor | 每个 Partition 的副本数量,决定数据冗余度和可用性。 | 建议设置为 3 以保证高可用,但会增加存储开销。 |
1.4 ZooKeeper 在 Kafka 中的作用(旧版)与 KRaft 模式(新版)
| 概念名称 | 说明 | 注意事项 |
|---|---|---|
| ZooKeeper(旧版) | Kafka 依赖 ZooKeeper 管理集群元数据,如 Broker 注册、Topic 配置、Leader 选举等。 | Kafka 3.0 前必须依赖 ZooKeeper,增加了系统复杂性和运维成本。 |
| KRaft 模式 | Kafka Raft Metadata Mode,Kafka 自研的共识协议,替代 ZooKeeper 进行元数据管理。 | 从 Kafka 2.8 开始支持,3.0+ 推荐使用,简化架构,提升性能和可维护性。 |
| Controller 选举 | 在 KRaft 模式下,Controller 由集群内部通过 Raft 协议选举产生,无需外部 ZooKeeper。 | KRaft 模式下集群可独立运行,减少外部依赖。 |
| 元数据存储 | KRaft 将元数据存储在内部 Topic(如 __cluster_metadata)中,由 Controller 管理。 | 数据一致性由 Raft 协议保证,避免 ZooKeeper 成为瓶颈。 |
| 混合模式 | 过渡阶段支持 ZooKeeper 和 KRaft 共存,便于升级。 | 建议新集群直接使用 KRaft,老集群逐步迁移。 |
| KRaft 节点角色 | 支持 Controller 和 Broker 角色,通过配置决定节点类型。 | 可部署专用 Controller 节点或混合部署,根据集群规模选择。 |
第二章:开发环境搭建与快速入门
2.1 搭建本地 Kafka 开发环境(单机/集群)
| 操作步骤 | 说明 | 注意事项 |
|---|---|---|
| 下载 Kafka 发行版 | 从 Apache 官网下载 Kafka 二进制包,如 kafka_2.13-3.8.0.tgz。 | 选择与 Scala 版本兼容的包,开发环境建议使用最新稳定版。 |
| 启动 ZooKeeper(旧版) | 执行 bin/zookeeper-server-start.sh config/zookeeper.properties | 仅在使用 ZooKeeper 模式时需要,KRaft 模式可跳过。 |
| 启动 Kafka Broker | 执行 bin/kafka-server-start.sh config/server.properties | 确保配置文件中 broker.id、listeners、log.dirs 等正确设置。 |
| 配置 KRaft 模式(新版) | 使用 bin/kafka-storage.sh 格式化存储目录,并启动 server.properties 启用 KRaft。 | 需设置 process.roles=broker,controller,controller.quorum.voters 等参数。 |
| 单机环境验证 | 使用命令行创建 Topic、发送和接收消息测试。 | 确保防火墙开放 9092 端口,主机名解析正常。 |
| 伪集群搭建 | 修改多个 server.properties 文件,配置不同 broker.id、端口和日志目录。 | 用于本地测试多 Broker 场景,模拟集群行为。 |
| 环境变量与脚本封装 | 设置 KAFKA_HOME,编写启动/停止脚本简化操作。 | 提高开发效率,避免重复命令输入。 |
2.2 创建第一个 Kafka 生产者程序
| 方法名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
KafkaProducer(Properties) | new KafkaProducer<>(props) | 构造生产者实例,传入配置属性。 | 必须配置 bootstrap.servers、key.serializer、value.serializer。 |
send(ProducerRecord) | producer.send(record) | 发送消息到指定 Topic。 | 方法返回 Future 对象,可调用 get() 实现同步发送。 |
send(ProducerRecord, Callback) | producer.send(record, callback) | 异步发送消息,并指定回调函数处理结果。 | 回调在生产者线程中执行,避免阻塞。 |
flush() | producer.flush() | 强制将所有待发送消息立即发送出去。 | 通常在批量发送后调用,确保消息发出。 |
close() | producer.close() | 关闭生产者,释放资源。 | 必须在程序结束前调用,否则可能导致消息丢失或资源泄漏。 |
close(Duration) | producer.close(timeout) | 带超时时间的关闭方法,等待在途请求完成。 | 避免无限等待,合理设置超时时间。 |
代码示例:构造生产者实例
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
代码示例:异步发送(带回调)
ProducerRecord<String, String> record = new ProducerRecord<>("test-topic", "key1", "value1");
producer.send(record, (metadata, exception) -> {
if (exception != null) {
System.err.println("发送失败: " + exception);
} else {
System.out.println("发送成功到: " + metadata.topic() + "-" + metadata.partition());
}
});
2.3 创建第一个 Kafka 消费者程序
| 方法名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
KafkaConsumer(Properties) | new KafkaConsumer<>(props) | 构造消费者实例,传入配置属性。 | 必须配置 bootstrap.servers、group.id、key.deserializer、value.deserializer。 |
subscribe(Collection) | consumer.subscribe(Arrays.asList("topic1", "topic2")) | 订阅一个或多个 Topic。 | 订阅后会自动触发再平衡(Rebalance)。 |
subscribe(Pattern) | consumer.subscribe(Pattern.compile("test.*")) | 订阅符合正则表达式的所有 Topic。 | 适用于动态创建的 Topic。 |
poll(Duration) | consumer.poll(timeout) | 从服务器拉取消息,返回 ConsumerRecords。 | 必须在循环中持续调用,timeout 设置拉取阻塞时间。 |
commitSync() | consumer.commitSync() | 同步提交当前消费位移(offset),阻塞直到提交成功。 | 确保消息处理成功后再提交,防止重复消费。 |
commitAsync() | consumer.commitAsync() | 异步提交位移,不阻塞主线程。 | 提交失败不会重试,可在关闭前调用 commitSync() 补救。 |
close() | consumer.close() | 关闭消费者,释放资源。 | 必须调用,通常在 finally 块中执行。 |
代码示例:构造消费者实例
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "test-group");
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("auto.commit.interval.ms", "1000");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
代码示例:消费消息并同步提交位移
consumer.subscribe(Arrays.asList("test-topic"));
try {
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
System.out.printf("topic=%s, partition=%d, offset=%d, key=%s, value=%s%n",
record.topic(), record.partition(), record.offset(), record.key(), record.value());
}
consumer.commitSync();
}
} finally {
consumer.close();
}
2.4 使用命令行工具测试消息收发
| 命令名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
kafka-topics.sh --create | bin/kafka-topics.sh --create --topic <name> --partitions <n> --replication-factor <n> --bootstrap-server <host> | 创建 Topic | 单机测试 replication-factor 可设为 1,生产环境建议为 3。 |
kafka-topics.sh --list | bin/kafka-topics.sh --list --bootstrap-server <host> | 列出所有 Topic | 验证 Topic 是否创建成功。 |
kafka-topics.sh --describe | bin/kafka-topics.sh --describe --topic <name> --bootstrap-server <host> | 查看 Topic 详细信息(分区、副本、Leader 等) | 用于排查分区分布和副本状态。 |
kafka-console-producer.sh | bin/kafka-console-producer.sh --bootstrap-server <host> --topic <name> | 启动控制台生产者,手动输入消息发送 | 每行输入即发送一条消息,Ctrl+C 退出。 |
kafka-console-consumer.sh | bin/kafka-console-consumer.sh --bootstrap-server <host> --topic <name> --from-beginning | 启动控制台消费者,消费指定 Topic 消息 | --from-beginning 从头开始消费,否则从最新消息开始。 |
kafka-consumer-groups.sh --list | bin/kafka-consumer-groups.sh --list --bootstrap-server <host> | 列出所有消费者组 | 查看当前活跃的消费者组。 |
kafka-consumer-groups.sh --describe | bin/kafka-consumer-groups.sh --describe --group <name> --bootstrap-server <host> | 查看消费者组的消费进度(位移、LAG 等) | LAG 表示未消费消息数,用于监控消费延迟。 |
命令行使用示例:
# 创建 Topic
bin/kafka-topics.sh --create --topic test-topic --partitions 3 --replication-factor 1 --bootstrap-server localhost:9092
# 列出所有 Topic
bin/kafka-topics.sh --list --bootstrap-server localhost:9092
# 查看 Topic 详情
bin/kafka-topics.sh --describe --topic test-topic --bootstrap-server localhost:9092
# 发送消息
bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic test-topic
# 消费消息
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic test-topic --from-beginning
# 查看消费者组
bin/kafka-consumer-groups.sh --list --bootstrap-server localhost:9092
# 查看消费进度
bin/kafka-consumer-groups.sh --describe --group test-group --bootstrap-server localhost:9092
第三章:Kafka 生产者(Producer)详解
3.1 生产者核心配置参数
| 配置参数 | 语法示例 | 用途 | 注意事项 |
|---|---|---|---|
bootstrap.servers | props.put("bootstrap.servers", "host1:9092,host2:9092") | 指定 Kafka 集群的初始连接地址列表,生产者通过它发现集群元数据。 | 至少配置两个 Broker 以防止单点故障,地址用逗号分隔。 |
key.serializer | props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer") | 指定消息键(Key)的序列化类,必须实现 Serializer 接口。 | 常见值:StringSerializer、IntegerSerializer、ByteArraySerializer 等。 |
value.serializer | props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer") | 指定消息值(Value)的序列化类。 | 与 key.serializer 类似,必须正确配置,否则发送失败。 |
acks | props.put("acks", "all") | 控制消息写入副本的确认机制。0=不等待确认,1=Leader 确认,all=所有 ISR 副本确认。 | acks=all 提供最高持久性,但延迟较高;acks=1 平衡性能与可靠性。 |
retries | props.put("retries", 3) | 设置生产者在遇到可重试异常(如网络抖动、Leader 选举)时的重试次数。 | 设为 Integer.MAX_VALUE 可实现无限重试,需配合 retry.backoff.ms 使用。 |
batch.size | props.put("batch.size", 16384) | 每个批次(Batch)的字节数,生产者会累积消息直到达到此大小再发送。 | 较大值可提高吞吐量,但增加延迟;默认 16KB。 |
linger.ms | props.put("linger.ms", 5) | 批次等待更多消息的时间(毫秒),用于增加批次大小。 | 设为 >0 可减少小消息发送次数,提升吞吐;但会增加延迟。 |
buffer.memory | props.put("buffer.memory", 33554432) | 生产者缓冲区总大小,用于暂存待发送消息。 | 默认 32MB,若消息产生速度过快可能导致缓冲区满,抛出 BufferExhaustedException。 |
compression.type | props.put("compression.type", "snappy") | 消息压缩类型,可选 none、gzip、snappy、lz4、zstd。 | 压缩可减少网络传输和存储开销,但增加 CPU 消耗。 |
enable.idempotence | props.put("enable.idempotence", true) | 启用幂等性生产者,确保消息在单分区不重复、不丢失、有序。 | 需配合 acks=all、retries=Integer.MAX_VALUE、max.in.flight.requests.per.connection=5(≤5)使用。 |
max.in.flight.requests.per.connection | props.put("max.in.flight.requests.per.connection", 5) | 每个连接最多允许的未确认请求个数。 | 幂等性模式下必须 ≤5,设为 1 可保证严格顺序,但降低吞吐。 |
3.2 发送消息的三种方式(同步、异步、带回调)
| 发送方式 | 语法与代码示例 | 用途 | 注意事项 |
|---|---|---|---|
| 同步发送 | RecordMetadata metadata = producer.send(record).get(); | 调用 send 后立即阻塞,等待 Broker 确认,返回消息元数据。 | 简单直观,但吞吐低,适用于关键消息或调试。 |
| 异步发送(无回调) | producer.send(record); | 发送后立即返回,不关心结果,吞吐最高。 | 无法得知发送是否成功,可能丢失消息,仅用于测试或容忍丢失的场景。 |
| 异步发送(带回调) | producer.send(record, new Callback() { ... }); | 发送后立即返回,通过回调函数处理成功或失败。 | 推荐方式,兼顾性能与错误处理,回调在生产者 I/O 线程中执行,避免阻塞。 |
代码示例:
// 同步发送
ProducerRecord<String, String> record = new ProducerRecord<>("topic", "key", "value");
RecordMetadata metadata = producer.send(record).get();
System.out.println("发送到: " + metadata.topic() + "-" + metadata.partition() + " offset=" + metadata.offset());
// 异步发送(无回调)
producer.send(record);
// 异步发送(带回调)
producer.send(record, new Callback() {
@Override
public void onCompletion(RecordMetadata metadata, Exception exception) {
if (exception != null) {
System.err.println("发送失败: " + exception.getMessage());
} else {
System.out.println("成功发送到 " + metadata.topic() + "-" + metadata.partition());
}
}
});
3.3 消息序列化与自定义 Serializer
| 方法/概念 | 说明 | 注意事项 |
|---|---|---|
| Serializer 接口 | 实现 org.apache.kafka.common.serialization.Serializer 接口,定义自定义序列化逻辑,将 Java 对象转为字节数组。 | 必须实现 serialize 方法,注意处理 null 值。 |
| 配置自定义 Serializer | props.put("value.serializer", "com.example.CustomSerializer"),指定自定义序列化类的全限定名。 | 类必须在 classpath 中,且无参构造函数可用。 |
close() 方法 | 在生产者关闭时调用,用于清理资源(如关闭连接、释放缓冲区)。 | 通常不需要显式资源管理,但复杂序列化器(如涉及连接池)应实现。 |
configure() 方法 | 在序列化器初始化时调用,可获取生产者配置参数。 | isKey 为 true 表示用于序列化 Key,false 表示 Value。 |
代码示例:
public class CustomSerializer implements Serializer<MyObject> {
@Override
public byte[] serialize(String topic, MyObject data) {
// 自定义序列化逻辑,如 JSON、Protobuf 等
return JSON.toJSONString(data).getBytes(StandardCharsets.UTF_8);
}
@Override
public void configure(Map<String, ?> configs, boolean isKey) {
this.encoding = (String) configs.get("serializer.encoding");
}
@Override
public void close() {
// 如关闭 Protobuf 池
}
}
3.4 分区策略与自定义 Partitioner
| 方法/概念 | 说明 | 用途 | 注意事项 |
|---|---|---|---|
| 默认分区策略 | 若消息有 Key,则对 Key 的 hash 取模;若无 Key,则轮询分配。 | 实现负载均衡,相同 Key 的消息总在同一分区,保证顺序。 | Key 为 null 时,生产者使用粘性分区(sticky partitioning)优化批次发送。 |
| 自定义 Partitioner | 实现 org.apache.kafka.clients.producer.Partitioner 接口 | 定义自己的分区分配逻辑,如按用户 ID、地理位置等路由。 | 必须实现 partition 方法,返回 0 到 numPartitions-1 之间的整数。 |
| 配置自定义 Partitioner | props.put("partitioner.class", "com.example.CustomPartitioner") | 在生产者配置中指定自定义分区器类名。 | 类必须在 classpath 中,且无参构造函数可用。 |
close() 和 configure() | 分别用于初始化和关闭时执行逻辑。 | 可在 configure 中读取配置参数,在 close 中释放资源。 | 通常可留空,除非分区器持有状态或资源。 |
代码示例:
public class CustomPartitioner implements Partitioner {
@Override
public int partition(String topic, Object key, byte[] keyBytes,
Object value, byte[] valueBytes, Cluster cluster) {
List<PartitionInfo> partitions = cluster.partitionsForTopic(topic);
int numPartitions = partitions.size();
// 自定义逻辑,如按用户 ID 哈希
return (int) (Math.abs(((String) key).hashCode()) % numPartitions);
}
@Override
public void configure(Map<String, ?> configs) { }
@Override
public void close() { }
}
3.5 消息确认机制(acks)与可靠性保证
| 配置参数 | 说明 | 用途 | 注意事项 |
|---|---|---|---|
acks=0 | 生产者不等待任何确认。 | 消息可能丢失,吞吐最高,可靠性最低,适用于可容忍丢失的场景(如日志)。 | 不推荐用于关键业务。 |
acks=1 | 只需 Leader 副本写入成功即确认,Follower 可能未同步。 | 平衡性能与可靠性,但若 Leader 故障且 Follower 未同步,消息会丢失。 | 常用折中方案。 |
acks=all(或 -1) | 需所有 ISR(In-Sync Replicas)副本都写入成功才确认。 | 提供最高持久性,即使 Leader 故障,消息也不会丢失。 | 需配合 min.insync.replicas 使用。 |
min.insync.replicas | Broker 端参数,定义 ISR 最小副本数。当 ISR 数量 < 此值时,Producer 写入会失败。 | 与 acks=all 配合使用,确保数据冗余。例如 replication.factor=3, min.insync.replicas=2,允许一个副本故障。 | Broker 配置:min.insync.replicas=2 |
3.6 消息重试机制与幂等性生产者
| 配置/概念 | 语法示例 | 用途 | 注意事项 |
|---|---|---|---|
retries | props.put("retries", Integer.MAX_VALUE) | 设置重试次数,应对临时性故障(如网络抖动、Leader 选举)。 | 默认 0,不重试。设为 Integer.MAX_VALUE 可无限重试。 |
retry.backoff.ms | props.put("retry.backoff.ms", 100) | 重试前等待的时间(毫秒),避免密集重试。 | 默认 100ms,可根据网络状况调整。 |
enable.idempotence | props.put("enable.idempotence", true) | 启用幂等性,确保每条消息在单个分区只被写入一次(不重复、不丢失、有序)。 | 自动设置 retries=Integer.MAX_VALUE, acks=all, max.in.flight.requests.per.connection=5。 |
| 幂等性原理 | 生产者为每个消息分配 PID(Producer ID)和 Sequence Number | Broker 根据 PID 和 Sequence Number 去重,防止重复写入。 | 仅保证单分区内的幂等性,跨分区或生产者重启后不保证。 |
| 事务性生产者(Transactional) | props.put("transactional.id", "txn-01"); producer.initTransactions(); | 提供跨分区、跨会话的原子性写入,确保一组消息全部成功或全部失败。 | 需启用幂等性,使用 producer.beginTransaction(), send(), commitTransaction() 等方法。 |
3.7 生产者拦截器(Interceptor)
| 方法名称 | 说明 | 用途 | 注意事项 |
|---|---|---|---|
onSend(ProducerRecord) | ProducerRecord<K, V> onSend(ProducerRecord<K, V> record),在消息发送前拦截。 | 可修改消息内容或记录日志。返回值为实际发送的消息,可包装或修改。 | 返回值为实际发送的消息。 |
onAcknowledgement(RecordMetadata, Exception) | void onAcknowledgement(RecordMetadata metadata, Exception exception),在收到 Broker 确认后拦截。 | 可记录发送结果或监控。无论成功或失败都会调用。 | 可用于监控和告警。 |
close() | 拦截器关闭时调用。 | 释放资源,通常用于清理工作。 | 如关闭监控客户端等。 |
configure(Map) | 拦截器初始化时调用。 | 获取配置参数,可读取生产者配置中的自定义参数。 | 如 this.metricPrefix = (String) configs.get("metric.prefix")。 |
| 配置拦截器 | props.put("interceptor.classes", "com.example.ProducerInterceptor") | 指定拦截器类名,可配置多个,按顺序执行。 | 拦截器链中任一环节抛出异常会中断流程,后续拦截器不执行。 |
代码示例:
public class LoggingInterceptor implements ProducerInterceptor<String, String> {
@Override
public ProducerRecord<String, String> onSend(ProducerRecord<String, String> record) {
System.out.println("发送消息: " + record.value());
return record; // 可返回修改后的 record
}
@Override
public void onAcknowledgement(RecordMetadata metadata, Exception exception) {
if (exception != null) {
System.err.println("发送失败: " + exception);
} else {
System.out.println("发送成功: " + metadata.offset());
}
}
@Override
public void configure(Map<String, ?> configs) {
this.metricPrefix = (String) configs.get("metric.prefix");
}
@Override
public void close() {
// 关闭监控客户端等
}
}
// 配置多个拦截器
props.put("interceptor.classes",
"com.example.LoggingInterceptor,com.example.MonitoringInterceptor");
第四章:Kafka 主题(Topic)与分区管理
4.1 Topic 的创建、查看与删除
| 操作方式 | 语法/命令 | 用途 | 注意事项 |
|---|---|---|---|
| 命令行创建 | bin/kafka-topics.sh --create --topic <name> --partitions <n> --replication-factor <n> --bootstrap-server <host> | 创建指定分区数和副本因子的 Topic | 副本因子不能大于 Broker 数量,否则创建失败。 |
| 命令行查看列表 | bin/kafka-topics.sh --list --bootstrap-server <host> | 列出集群中所有 Topic | 可用于验证 Topic 是否创建成功。 |
| 命令行查看详情 | bin/kafka-topics.sh --describe --topic <name> --bootstrap-server <host> | 查看 Topic 的分区分布、Leader、副本、ISR 等信息 | 输出包含 Partition、Leader、Replicas、Isr 字段,用于排查问题。 |
| 命令行删除 | bin/kafka-topics.sh --delete --topic <name> --bootstrap-server <host> | 删除指定 Topic(需 server.properties 中 delete.topic.enable=true) | 删除后数据不会立即清除,需等待日志清理线程执行。 |
| Java API 创建 | Admin.createTopics(new NewTopic(topic, partitions, replication)) | 使用 Admin 客户端在代码中创建 Topic | 需引入 kafka-clients 依赖,适合自动化部署场景。 |
| Java API 删除 | Admin.deleteTopics(Collection) | 使用 Admin 客户端删除 Topic | 异步操作,返回 DeleteTopicsResult 可监听结果。 |
| Java API 查看 | Admin.describeTopics(Collection) | 获取 Topic 配置和元数据信息 | 返回 TopicDescription 对象,包含分区和副本信息。 |
命令行示例:
# 创建 Topic
bin/kafka-topics.sh --create --topic user-logs --partitions 3 --replication-factor 2 --bootstrap-server localhost:9092
# 查看详情
bin/kafka-topics.sh --describe --topic user-logs --bootstrap-server localhost:9092
# 删除 Topic
bin/kafka-topics.sh --delete --topic user-logs --bootstrap-server localhost:9092
Java API 示例:
Admin admin = Admin.create(config);
// 创建 Topic
NewTopic newTopic = new NewTopic("orders", 6, (short) 3);
CreateTopicsResult result = admin.createTopics(Collections.singleton(newTopic));
// 查看 Topic
DescribeTopicsResult describeResult = admin.describeTopics(Arrays.asList("user-logs"));
KafkaFuture<Map<String, TopicDescription>> future = describeResult.values();
// 删除 Topic
DeleteTopicsResult deleteResult = admin.deleteTopics(Arrays.asList("temp-data"));
4.2 分区与副本机制(Leader、Follower)
| 概念名称 | 说明 | 注意事项 |
|---|---|---|
| Partition | Topic 的分区,每个 Partition 是一个有序、不可变的消息序列,独立存储和处理。 | 分区数决定并行度,增加分区可提升吞吐,但过多分区会增加管理开销。 |
| Replication Factor | 每个 Partition 的副本数量,包括 1 个 Leader 和多个 Follower。 | 副本用于容灾,建议生产环境设置为 3,确保高可用。 |
| Leader | 每个 Partition 的主副本,负责处理所有读写请求,客户端直接与 Leader 交互。 | Leader 选举由 Controller 负责,若 Leader 宕机,Controller 会从 ISR 中选举新 Leader。 |
| Follower | 副本的从节点,被动从 Leader 拉取消息,保持数据同步。 | Follower 不对外提供服务,仅用于数据冗余和故障转移。 |
| 数据分布 | 分区均匀分布在多个 Broker 上,实现负载均衡。 | 避免所有 Leader 集中在少数 Broker 上,影响性能。 |
| 写流程 | 生产者发送消息到 Leader → Leader 写入本地日志 → Follower 从 Leader 拉取数据同步。 | 消息写入 Leader 成功即返回(根据 acks 配置),Follower 异步复制。 |
| 读流程 | 消费者从 Leader 拉取消息,Follower 不参与读取。 | 保证数据一致性,避免从不同副本读取导致顺序错乱。 |
4.3 副本同步与 ISR(In-Sync Replicas)机制
| 概念名称 | 说明 | 注意事项 |
|---|---|---|
| ISR(In-Sync Replicas) | 与 Leader 保持同步的副本集合,包括 Leader 自身和同步进度未落后的 Follower。 | ISR 是动态变化的,Broker 会定期检查 Follower 的同步延迟。 |
| 同步条件 | Follower 在 replica.lag.time.max.ms 时间内从 Leader 获取最新消息且未断开连接。 | 默认值为 30 秒,超过此时间 Follower 将被移出 ISR。 |
| OSR(Out-of-Sync Replicas) | 落后过多或断开连接的 Follower 副本,不在 ISR 中。 | OSR 中的副本无法参与 Leader 选举,存在数据丢失风险。 |
| Leader 选举 | 当 Leader 宕机时,Controller 从 ISR 中选举新 Leader,确保数据不丢失。 | 仅从 ISR 中选举,避免选择落后副本导致数据不一致。 |
| Unclean Leader Election | 若 ISR 为空,是否允许从 OSR 中选举 Leader(由 unclean.leader.election.enable 控制)。 | 设为 true 可能导致数据丢失,生产环境建议设为 false。 |
| 高水位(HW) | 表示已提交(committed)的消息的偏移量,所有 ISR 副本都已写入该位置之前的消息。 | 消费者只能读取到 HW 之前的消息,保证一致性。 |
| LEO(Log End Offset) | 日志末端偏移量,表示下一条待写入消息的位置。 | Leader 和 Follower 都有自己的 LEO,用于计算同步进度。 |
4.4 分区再平衡与 Reassignment
| 操作方式 | 语法/命令 | 用途 | 注意事项 |
|---|---|---|---|
| 自动再平衡 | 消费者组内成员变化时自动触发 | 消费者加入或退出时重新分配 Partition | 再平衡期间消费暂停,应尽量减少触发频率。 |
| 手动 Reassignment | bin/kafka-reassign-partitions.sh --generate --topics-to-move-json-file <file> --broker-list | 生成分区迁移计划 | 用于 Broker 扩容或缩容时重新分布分区。 |
| 执行 Reassignment | bin/kafka-reassign-partitions.sh --execute --reassignment-json-file <file> --bootstrap-server | 执行生成的分区迁移计划 | 迁移过程不影响服务,但会增加网络和磁盘 I/O。 |
| 验证 Reassignment | bin/kafka-reassign-partitions.sh --verify --reassignment-json-file <file> --bootstrap-server | 验证分区迁移是否完成 | 迁移未完成前不要删除临时文件。 |
| 取消 Reassignment | bin/kafka-reassign-partitions.sh --cancel --bootstrap-server | 取消正在进行的分区迁移 | 取消后可能留下部分迁移数据,需手动清理。 |
| 增加分区数 | bin/kafka-topics.sh --alter --topic <name> --partitions <n> | 增加 Topic 的分区数量(只能增加不能减少) | 分区数增加后,消费者组会触发再平衡,注意 key 到分区的映射可能变化。 |
命令行示例:
# 生成分区迁移计划
echo '{"topics":[{"topic":"user-logs"}],"version":1}' > topics.json
bin/kafka-reassign-partitions.sh --generate --topics-to-move-json-file topics.json --broker-list "0,1,2"
# 执行迁移
bin/kafka-reassign-partitions.sh --execute --reassignment-json-file reassignment.json --bootstrap-server localhost:9092
# 验证迁移
bin/kafka-reassign-partitions.sh --verify --reassignment-json-file reassignment.json --bootstrap-server localhost:9092
# 取消迁移
bin/kafka-reassign-partitions.sh --cancel --bootstrap-server localhost:9092
# 增加分区数
bin/kafka-topics.sh --alter --topic user-logs --partitions 6 --bootstrap-server localhost:9092
4.5 Topic 配置参数调优
| 参数名称 | 默认值 | 用途 | 推荐设置/调优建议 | 注意事项 |
|---|---|---|---|---|
retention.ms | 604800000(7 天) | 控制消息保留时间,超过此时间的消息将被清理。 | 根据业务需求设置,如 1 天、3 天或更长。 | 设置过短可能导致消费者来不及消费,过长则占用存储。 |
retention.bytes | -1(不限制) | 控制单个 Partition 的最大日志大小。 | 根据磁盘容量和数据量设置,如 1073741824(1GB)。 | 与 retention.ms 共同作用,任一条件满足即触发清理。 |
segment.bytes | 1073741824(1GB) | 日志段文件大小,达到该值后滚动到新文件。 | 可根据写入频率调整,频繁写入可设小些(如 512MB)。 | 较小的段便于清理,但增加文件句柄开销。 |
segment.ms | 604800000(7 天) | 日志段最大存活时间,超过后强制滚动。 | 可设置为 24 小时,确保每天一个段文件便于管理。 | 与 segment.bytes 共同控制日志滚动。 |
cleanup.policy | delete | 清理策略:delete(按时间/大小删除)或 compact(压缩,保留 key 最新值)。 | 普通日志用 delete,需要保留状态的用 compact。 | compact 模式适用于 key-value 状态存储场景。 |
max.message.bytes | 1048588(1MB) | 单条消息最大大小。 | 若需发送大消息,可调大(如 10MB),需同步调整生产者和消费者配置。 | 过大消息影响性能和延迟,建议拆分或使用外部存储。 |
min.insync.replicas | 1 | ISR 中最小副本数,配合 acks=all 使用,确保消息写入多数副本。 | 建议设为 2(副本数 ≥3 时),提高数据可靠性。 | 若 ISR 数不足,生产者将收到 NotEnoughReplicasException。 |
unclean.leader.election.enable | false | 是否允许从非 ISR 副本中选举 Leader。 | 生产环境建议 false,避免数据丢失;对可用性要求极高可设 true。 | 设为 true 可能导致数据不一致。 |
compression.type | producer | 消息压缩类型:none、gzip、snappy、lz4、zstd。 | 建议 lz4 或 zstd,压缩比高且性能好。 | 压缩减少网络和磁盘开销,但增加 CPU 使用。 |
第五章:Kafka 消费者(Consumer)详解
5.1 消费者核心配置参数
| 配置参数名称 | 语法示例 | 用途 | 注意事项 |
|---|---|---|---|
bootstrap.servers | props.put("bootstrap.servers", "host1:9092,host2:9092") | 指定 Kafka 集群的初始连接地址列表。 | 客户端会自动发现所有 Broker,建议配置 2-3 个以提高可用性。 |
group.id | props.put("group.id", "consumer-group-1") | 指定消费者所属的消费者组 ID。 | 同一组内的消费者共享 Topic 的消费负载。 |
key.deserializer | props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer") | 指定 key 的反序列化类。 | 必须实现 Deserializer 接口,常用有 String、Integer、ByteArray 等。 |
value.deserializer | props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer") | 指定 value 的反序列化类。 | 同上,需与生产者序列化方式匹配。 |
enable.auto.commit | props.put("enable.auto.commit", "true") | 是否开启自动提交位移。 | 设为 true 时由 auto.commit.interval.ms 控制提交频率。 |
auto.commit.interval.ms | props.put("auto.commit.interval.ms", "5000") | 自动提交位移的间隔时间(毫秒)。 | 仅在 enable.auto.commit=true 时生效,过短增加 Broker 压力,过长增加重复消费风险。 |
auto.offset.reset | props.put("auto.offset.reset", "earliest") | 当位移不存在或无效时的重置策略。 | 可选值:earliest(从头开始)、latest(从最新消息开始)、none(抛异常)。 |
session.timeout.ms | props.put("session.timeout.ms", "45000") | 消费者与 Broker 的心跳超时时间。 | 超时后 Broker 认为消费者失效并触发 Rebalance,建议设置为心跳间隔的 3 倍以上。 |
heartbeat.interval.ms | props.put("heartbeat.interval.ms", "3000") | 消费者向 Broker 发送心跳的间隔。 | 必须小于 session.timeout.ms,通常为其 1/3。 |
max.poll.records | props.put("max.poll.records", "500") | 每次 poll() 调用返回的最大记录数。 | 控制单次处理消息量,避免处理超时导致 Rebalance。 |
max.poll.interval.ms | props.put("max.poll.interval.ms", "300000") | 两次 poll() 调用之间的最大间隔。 | 若处理消息时间过长超过此值,消费者会被踢出组。 |
fetch.min.bytes | props.put("fetch.min.bytes", "1") | 每次拉取请求返回的最小数据量。 | 设为较大值可减少网络请求,但增加延迟。 |
fetch.max.wait.ms | props.put("fetch.max.wait.ms", "500") | 拉取请求在 Broker 端等待数据积累的最大时间。 | 配合 fetch.min.bytes 使用,实现批量拉取。 |
5.2 订阅主题与消息拉取机制
| 方法名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
subscribe(Collection) | consumer.subscribe(Arrays.asList("t1")) | 订阅指定的一个或多个 Topic。 | 订阅后自动参与再平衡,不可与 assign 混用。 |
subscribe(Pattern) | consumer.subscribe(Pattern.compile("t.*")) | 订阅符合正则表达式的 Topic。 | 动态匹配新创建的 Topic,适合监控类应用。 |
assign(Collection) | consumer.assign(Arrays.asList(tp)) | 手动为消费者分配特定的分区(TopicPartition)。 | 跳过消费者组管理,适用于特定分区处理或外部位移管理。 |
poll(Duration) | consumer.poll(Duration.ofMillis(100)) | 拉取消息,阻塞指定时间或直到有数据可返回。 | 必须在循环中调用,是消费者主循环的核心。 |
wakeup() | consumer.wakeup() | 唤醒阻塞的 poll() 或 commit 操作,用于优雅关闭。 | 唯一可安全跨线程调用的方法,用于中断阻塞操作。 |
position(TopicPartition) | consumer.position(tp) | 获取指定分区下一次拉取的起始位移。 | 仅在分配了分区后调用有效。 |
endOffsets(Collection) | consumer.endOffsets(partitions) | 获取指定分区当前的最后一条消息的位移(即日志末尾)。 | 用于计算 LAG 或初始化位移。 |
代码示例:
// 订阅 Topic
consumer.subscribe(Arrays.asList("order-topic", "log-topic"));
// 正则订阅
consumer.subscribe(Pattern.compile("metrics-.*"));
// 手动分配分区
TopicPartition tp = new TopicPartition("test-topic", 0);
consumer.assign(Arrays.asList(tp));
// 获取位移
long nextOffset = consumer.position(new TopicPartition("t", 0));
Set<TopicPartition> partitions = consumer.assignment();
Map<TopicPartition, Long> ends = consumer.endOffsets(partitions);
5.3 反序列化与自定义 Deserializer
| 方法名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
configure(Map, boolean) | public void configure(Map<String, ?> configs, boolean isKey) | 初始化 Deserializer,获取配置参数。 | isKey 为 true 表示反序列化 key,false 表示 value。 |
deserialize(String, byte[]) | public T deserialize(String topic, byte[] data) | 执行反序列化逻辑,将字节数组转换为 Java 对象。 | 可根据 topic 或 data 内容进行差异化处理。 |
close() | public void close() | 关闭 Deserializer,释放资源。 | 通常用于关闭缓存或连接,但大多数 Deserializer 无状态,无需实现。 |
注: 自定义 Deserializer 需实现
org.apache.kafka.common.serialization.Deserializer<T>接口。
代码示例:
public class CustomDeserializer implements Deserializer<String> {
@Override
public void configure(Map<String, ?> configs, boolean isKey) {
this.isKey = isKey;
}
@Override
public String deserialize(String topic, byte[] data) {
return data == null ? null : new String(data, StandardCharsets.UTF_8);
}
@Override
public void close() {
// 清理资源
}
}
5.4 消费者组(Consumer Group)与再平衡(Rebalance)
| 概念名称 | 说明 | 注意事项 |
|---|---|---|
| 消费者组(Consumer Group) | 一组消费者的逻辑集合,共同消费一个或多个 Topic,实现消息的广播或负载均衡。 | 同一 Group 内每个 Partition 只能被一个 Consumer 消费,实现队列模式。 |
| 再平衡(Rebalance) | 当消费者组成员变化(加入/退出)或订阅 Topic 分区数变化时,Broker 重新分配 Partition 的过程。 | 过程中所有消费者暂停消费,影响吞吐,应尽量减少不必要的 Rebalance。 |
| Group Coordinator | 负责管理消费者组的 Broker 节点,处理成员加入、同步、位移提交等。 | 每个 Group 有唯一的 Coordinator,由 Group ID 哈希决定。 |
| JoinGroup Request | 消费者请求加入组的消息。 | 首次订阅或 Rebalance 时发送。 |
| SyncGroup Request | 在 Leader Consumer 选出后,向成员分发分区分配方案。 | 分配方案由 Group Leader 制定(默认 StickyAssignor)。 |
| Heartbeat Request | 消费者定期发送心跳以保持活跃状态。 | 由 heartbeat.interval.ms 控制频率,超时则被踢出组。 |
| Rebalance 避免策略 | 减少处理时间、合理设置 max.poll.interval.ms、避免长时间 GC。 | 处理消息逻辑不应阻塞 poll() 循环。 |
5.5 位移(Offset)管理:自动提交与手动提交
| 方法名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
commitSync() | consumer.commitSync() | 同步提交当前所有分区的位移,阻塞直到提交成功。 | 确保消息处理成功后再提交,防止数据丢失。 |
commitSync(Map) | consumer.commitSync(offsets) | 同步提交指定分区的位移。 | 精确控制每个分区的提交位移,常用于精确一次语义。 |
commitAsync() | consumer.commitAsync() | 异步提交位移,不阻塞线程。 | 性能更高,但失败不会重试。 |
commitAsync(Callback) | consumer.commitAsync(callback) | 异步提交并指定回调处理结果。 | 回调执行在消费者线程,避免复杂操作。 |
committed(TopicPartition) | consumer.committed(tp) | 查询指定分区已提交的位移。 | 返回 null 表示该分区无提交记录。 |
seek(TopicPartition, long) | consumer.seek(tp, offset) | 手动将消费者位移定位到指定位置。 | 用于重新处理消息或跳过数据,需在 poll() 前调用。 |
seekToBeginning(Collection) | consumer.seekToBeginning(partitions) | 将指定分区的位移重置到起始位置。 | 常用于测试或重新消费。 |
seekToEnd(Collection) | consumer.seekToEnd(partitions) | 将指定分区的位移定位到末尾(最新消息之后)。 | 实现从最新消息开始消费,等效于 auto.offset.reset=latest。 |
代码示例:
// 同步提交
try {
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
// 处理消息
}
consumer.commitSync();
}
} finally {
consumer.close();
}
// 指定分区提交
Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>();
offsets.put(tp, new OffsetAndMetadata(offset + 1));
consumer.commitSync(offsets);
// 异步提交
consumer.commitAsync((offsets, exception) -> {
if (exception != null) {
log.error("Commit failed", exception);
}
});
// 查询已提交位移
OffsetAndMetadata committedOffset = consumer.committed(tp);
// 定位位移
consumer.seek(tp, 100L);
consumer.seekToBeginning(consumer.assignment());
consumer.seekToEnd(consumer.assignment());
5.6 再平衡监听器(RebalanceListener)
| 方法名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
onPartitionsRevoked(Collection) | public void onPartitionsRevoked(Collection<TopicPartition> partitions) | 在分区被撤销前调用,用于提交位移或清理资源。 | 必须在此方法中完成位移提交,否则可能丢失。 |
onPartitionsAssigned(Collection) | public void onPartitionsAssigned(Collection<TopicPartition> partitions) | 在分区被分配后调用,可用于重置本地状态或 seek 位移。 | 可在此方法中进行初始化操作。 |
使用方式: 在调用 subscribe 时传入监听器实例:
consumer.subscribe(topics, new ConsumerRebalanceListener() {
@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
consumer.commitSync();
}
@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
for (TopicPartition tp : partitions) {
consumer.seek(tp, 0);
}
}
});
5.7 消费者拦截器(Interceptor)
| 方法名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
onConsume(ConsumerRecord) | public ConsumerRecord onConsume(ConsumerRecord record) | 在消息返回给用户前拦截并处理。 | 可修改消息内容或头信息,或记录日志。 |
onCommit(Map) | public void onCommit(Map<TopicPartition, OffsetAndMetadata> offsets) | 在位移提交完成后调用。 | 仅通知作用,不能修改提交行为。 |
configure(Map) | public void configure(Map<String, ?> configs) | 初始化拦截器。 | 获取消费者配置参数。 |
close() | public void close() | 关闭拦截器,释放资源。 | 在消费者关闭时调用。 |
配置方式:
props.put("interceptor.classes", "com.example.MyConsumerInterceptor");
代码示例:
public class LoggingConsumerInterceptor implements ConsumerInterceptor<String, String> {
@Override
public ConsumerRecords<String, String> onConsume(ConsumerRecords<String, String> records) {
for (ConsumerRecord<String, String> record : records) {
System.out.println("消费: " + record.topic() + " " + record.value());
}
return records;
}
@Override
public void onCommit(Map<TopicPartition, OffsetAndMetadata> offsets) {
System.out.println("提交位移: " + offsets);
}
@Override
public void configure(Map<String, ?> configs) {
this.configs = configs;
}
@Override
public void close() {
// 释放资源
}
}
第六章:Kafka 高级特性
6.1 Kafka Connect:数据导入导出
| 概念/组件名称 | 说明 | 注意事项 |
|---|---|---|
| Kafka Connect | 分布式、可扩展的数据集成工具,用于在 Kafka 与其他系统间可靠地导入导出数据。 | 支持批量和流式数据同步,无需编写代码。 |
| Source Connector | 源连接器,从外部系统(如数据库、文件)读取数据并写入 Kafka Topic。 | 常见实现:JDBC Source、File Source、Debezium(CDC)。 |
| Sink Connector | 目标连接器,从 Kafka 读取数据并写入外部系统(如数据库、HDFS、Elasticsearch)。 | 常见实现:JDBC Sink、HDFS Sink、Elasticsearch Sink。 |
| Connect Worker | 运行 Connector 和 Task 的进程,支持独立模式(单节点)和分布式模式(多节点)。 | 分布式模式支持自动故障转移和水平扩展。 |
| Task | Connector 的实际执行单元,一个 Connector 可拆分为多个 Task 并行运行。 | Task 数量通常与 Topic 分区数匹配,提升吞吐量。 |
| Converter | 负责在 Connect 内部数据格式与外部格式(如 JSON、Avro)之间转换。 | 常用:JsonConverter、AvroConverter(需 Schema Registry)。 |
| Transformations | 在数据流动过程中进行轻量级转换,如字段过滤、重命名、路由等。 | 使用 SMT(Single Message Transform)实现,避免额外处理服务。 |
配置示例(JDBC Source):
{
"name": "jdbc-source",
"config": {
"connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
"tasks.max": "1",
"connection.url": "jdbc:mysql://localhost:3306/db",
"table.whitelist": "users",
"mode": "incrementing",
"incrementing.column.name": "id",
"topic.prefix": "db-"
}
}
注意: 需确保数据库驱动在 CLASSPATH 中。
6.2 Kafka Streams:流处理入门
| 方法/组件名称 | 语法/用途 | 注意事项 |
|---|---|---|
StreamsBuilder | 构建流处理拓扑的构建器。 | 用于定义 KStream 和 KTable。 |
KStream | 表示一个持续的数据流,支持 map、filter、flatMap 等操作。 | 适用于事件流处理。 |
KTable | 表示一个 changelog 流,代表某个数据的最新状态。 | 适用于聚合和状态查询。 |
stream(String topic) | 从指定 Topic 创建 KStream。 | 消费 Topic 所有消息。 |
process(Processor) | 对每条记录执行自定义处理逻辑。 | 可访问记录元数据和状态存储。 |
map(KeyValueMapper) | 转换每条记录的 key 和 value。 | 无状态转换。 |
filter(Predicate) | 过滤满足条件的记录。 | 保留返回 true 的记录。 |
groupByKey() | 按 key 分组,用于后续聚合。 | 要求 key 已存在或使用 map 操作设置。 |
count() | 对分组数据进行计数。 | 结果写入状态存储,支持容错。 |
to(String topic) | 将流写入指定 Topic。 | 输出 Topic 需预先创建或自动创建。 |
KafkaStreams.start() | 启动流处理应用。 | 应用为长期运行服务。 |
KafkaStreams.close() | 关闭流处理应用。 | 释放资源,持久化状态。 |
代码示例:
StreamsBuilder builder = new StreamsBuilder();
// 从 Topic 创建流
KStream<String, String> source = builder.stream("orders");
KStream<String, String> stream = builder.stream("input-topic");
// 转换
stream.map((key, value) -> new KeyValue<>(key.toUpperCase(), value.length()));
// 过滤
stream.filter((key, value) -> value.contains("error"));
// 分组与聚合
KTable<String, Long> counts = stream.groupByKey().count(Materialized.as("counts-store"));
// 输出到 Topic
stream.to("output-topic");
// 启动
KafkaStreams streams = new KafkaStreams(builder.build(), config);
streams.start();
// 关闭
streams.close();
6.3 KSQL:流式 SQL 查询
| 命令/语句类型 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
CREATE STREAM | CREATE STREAM stream_name (column_defs) WITH (kafka_topic='topic', value_format='FORMAT'); | 创建流,映射到 Kafka Topic。 | value_format 支持 JSON、AVRO、DELIMITED 等。 |
CREATE TABLE | CREATE TABLE table_name (column_defs) WITH (kafka_topic='topic', value_format='FORMAT', key='key_column'); | 创建表,表示实体的当前状态。 | 需指定主键(key)。 |
SELECT | SELECT * FROM stream_name WHERE condition; | 查询流或表中的数据。 | 支持 WHERE、GROUP BY、JOIN 等标准 SQL 操作。 |
INSERT INTO | INSERT INTO stream_name SELECT ...; | 将查询结果插入到另一个流。 | 用于数据路由和转换。 |
CREATE STREAM AS SELECT | CREATE STREAM output_stream AS SELECT ... FROM input_stream ...; | 创建持续查询,将结果写入新 Topic。 | 查询持续运行,输出流自动创建。 |
JOIN | SELECT s.id, t.name FROM stream s LEFT JOIN table t ON s.id = t.id; | 流与表或流与流之间的连接操作。 | 支持 INNER、LEFT JOIN,流-流 JOIN 需窗口定义。 |
WINDOW | WINDOW TUMBLING (SIZE 5 MINUTES) | 定义时间窗口,用于聚合操作。 | 支持 Tumbling、Hopping、Session 窗口。 |
TERMINATE | TERMINATE query_id; | 停止持续查询。 | 停止后不再产生数据,但 Topic 保留。 |
示例:
-- 创建流
CREATE STREAM page_views (view_time BIGINT, user_id VARCHAR)
WITH (kafka_topic='page_views', value_format='JSON');
-- 创建表
CREATE TABLE users (user_id VARCHAR PRIMARY KEY, email VARCHAR)
WITH (kafka_topic='users', value_format='JSON', key='user_id');
-- 查询
SELECT user_id, view_time FROM page_views WHERE view_time > 1600000000;
-- 数据路由
INSERT INTO high_value_users SELECT * FROM users WHERE premium = true;
-- 持续查询
CREATE STREAM errors AS SELECT * FROM logs WHERE level = 'ERROR';
-- 窗口聚合
SELECT user_id, COUNT(*) FROM page_views
WINDOW TUMBLING (SIZE 1 MINUTE) GROUP BY user_id;
-- 停止查询
TERMINATE CSAS_ERRORS_1;
6.4 消息压缩(Compression)机制
| 压缩类型 | 特点 | 适用场景 | 注意事项 |
|---|---|---|---|
none | 无压缩,消息原文传输。 | 网络和磁盘资源充足,或消息本身已压缩。 | 默认值,性能最好,但占用资源最多。 |
gzip | 高压缩比,CPU 消耗较高。 | 对存储和带宽要求高,可接受较高 CPU 开销。 | 适合日志归档等场景。 |
snappy | 压缩比适中,速度快,由 Google 开发。 | 平衡压缩比和性能,通用推荐。 | Kafka 内置支持,广泛使用。 |
lz4 | 压缩速度极快,压缩比接近 snappy。 | 高吞吐场景,追求低延迟。 | Kafka 默认压缩类型(从 0.10+)。 |
zstd | 高压缩比,可调压缩级别,性能优秀。 | 需要高压缩比且 CPU 资源充足,现代推荐。 | Kafka 2.1+ 支持,压缩比优于 gzip,速度优于 snappy。 |
| 批处理压缩 | 压缩发生在消息批次(batch)级别,而非单条消息。 | 提升整体压缩效率。 | batch.size 和 linger.ms 影响压缩效果。 |
| 消费者解压 | 消费者自动解压,无需配置。 | 透明处理,开发者无感知。 | Broker 存储压缩后的消息,消费者拉取后解压。 |
生产者配置:
props.put("compression.type", "lz4"); // 推荐
props.put("compression.type", "zstd"); // 高压缩比
6.5 消息头(Headers)的使用
| 方法名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
Headers.add(String, byte[]) | record.headers().add("trace-id", bytes) | 添加消息头键值对。 | 值为 byte[],需自行序列化。 |
Headers.lastHeader(String) | Header header = record.headers().lastHeader("trace-id"); | 获取指定名称的最后一个头(允许重复键)。 | 返回 Header 对象,value() 返回 byte[]。 |
Headers.toArray() | Header[] headers = record.headers().toArray(); | 获取所有头的数组。 | 可用于遍历所有头信息。 |
Headers.remove(String) | record.headers().remove("temp-header"); | 移除指定名称的头。 | 移除后不再传递。 |
| 生产者设置头 | 在 ProducerRecord 创建后添加。 | 传递上下文信息,如认证、追踪、路由指令。 | 头信息不参与分区计算。 |
| 消费者读取头 | 在 ConsumerRecord 中通过 headers() 方法获取。 | 在处理消息前读取元数据。 | 头在消息传递中保持不变。 |
| 序列化注意事项 | 头值需手动序列化为 byte[],如 String 需 getBytes,对象可 JSON 或二进制序列化。 | 确保生产者和消费者使用相同编码。 | 使用 StandardCharsets.UTF_8 确保一致性。 |
代码示例:
// 生产者添加头
ProducerRecord<String, String> record = new ProducerRecord<>("topic", "key", "value");
record.headers().add("trace-id", "abc123".getBytes(StandardCharsets.UTF_8));
record.headers().add("auth-token", tokenBytes);
// 消费者读取头
ConsumerRecord<String, String> cr = ...;
Header header = cr.headers().lastHeader("trace-id");
if (header != null) {
String traceId = new String(header.value(), StandardCharsets.UTF_8);
}
// 遍历所有头
for (Header h : cr.headers()) {
System.out.println(h.key() + ": " + Arrays.toString(h.value()));
}
// 移除头
record.headers().remove("internal-flag");
6.6 事务消息(Transactions)支持
| 方法名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
initTransactions() | producer.initTransactions(); | 初始化事务,必须在启用事务的生产者上调用。 | 每个 transactional.id 在集群内唯一,用于恢复事务状态。 |
beginTransaction() | producer.beginTransaction(); | 开始新事务。 | 必须在 initTransactions() 之后调用。 |
send(ProducerRecord) | producer.send(record); | 在事务中发送消息,消息对消费者不可见。 | 所有发送在事务提交前对消费者不可见。 |
commitTransaction() | producer.commitTransaction(); | 提交事务,消息变为可见。 | 仅当所有消息成功写入且日志已同步。 |
abortTransaction() | producer.abortTransaction(); | 中止事务,丢弃所有已发送消息。 | 发生异常时调用,确保原子性。 |
sendOffsetsToTransaction | producer.sendOffsetsToTransaction(offsets, consumerGroupId); | 在事务中提交消费者位移,实现精确一次(exactly-once)处理。 | 需配合 Consumer 使用,确保消费-处理-生产的原子性。 |
| EOS(Exactly-Once Semantics) | Kafka 0.11+ 支持 | 保证消息在流处理中处理一次,不丢失不重复。 | 结合事务和幂等性实现,需设置 processing.guarantee="exactly_once_v2"(Streams)。 |
代码示例:
// 事务性生产者配置
Properties props = new Properties();
props.put("enable.idempotence", "true");
props.put("transactional.id", "tx-producer-1"); // 全局唯一
// ... 其他配置
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.initTransactions();
// 事务发送
try {
producer.beginTransaction();
producer.send(new ProducerRecord<>("t1", "k1", "v1"));
producer.send(new ProducerRecord<>("t2", "k2", "v2"));
producer.commitTransaction();
} catch (Exception e) {
producer.abortTransaction();
}
// 精确一次处理(EOS)
Map<TopicPartition, OffsetAndMetadata> offsets = ...;
producer.sendOffsetsToTransaction(offsets, "group-1");
注意: 同一
transactional.id只能有一个活跃生产者,否则会踢出旧实例。
第七章:性能调优与监控
7.1 生产者性能优化策略
| 优化策略 | 配置参数/方法 | 说明 | 推荐值/建议 | 注意事项 |
|---|---|---|---|---|
| 批量发送 | batch.size | 控制批次大小,达到后发送。 | 16KB ~ 128KB,根据消息大小调整。 | 过大增加延迟,过小降低吞吐。 |
| 延迟等待 | linger.ms | 批次未满时等待更多消息的时间。 | 5 ~ 100ms,平衡延迟与吞吐。 | 与 batch.size 协同作用,提升批次效率。 |
| 压缩 | compression.type | 启用压缩减少网络和磁盘开销。 | lz4 或 zstd,兼顾压缩比与性能。 | 压缩在 Producer 端进行,增加 CPU 使用。 |
| 异步发送 | enable.idempotence=true | 启用幂等性,确保单分区不重复。 | 必须设置,配合 max.in.flight.requests.per.connection=1 或 5(Kafka 2.1+)。 | transactional.id 用于事务,保证原子性。 |
| 并行度 | 多线程生产或增加分区 | 提升整体吞吐量。 | 分区数应大于等于生产者线程数。 | 分区是并行度的基础。 |
| 序列化器优化 | 自定义高效序列化(如 Avro、Protobuf) | 减少消息体积和序列化开销。 | 避免使用 Java 原生序列化。 | 结合 Schema Registry 管理模式。 |
| ACK 机制 | acks=1 或 acks=all | acks=1:Leader 写入即返回,低延迟;acks=all:ISR 全同步,高可靠。 | 高可靠性场景用 all,配合 min.insync.replicas>=2。 | acks=0 不可靠,仅用于可丢失场景。 |
| 连接池与缓冲 | buffer.memory | 生产者缓存消息的总内存大小。 | 根据吞吐需求调整,如 32MB ~ 128MB。 | 溢出时 send() 会阻塞或抛出异常。 |
7.2 消费者性能优化策略
| 优化策略 | 配置参数/方法 | 说明 | 推荐值/建议 | 注意事项 |
|---|---|---|---|---|
| 批量拉取 | fetch.min.bytes | 每次拉取请求返回的最小数据量。 | 1KB ~ 10KB,减少小批次拉取。 | 过大增加延迟,需与 fetch.max.wait.ms 平衡。 |
| 拉取等待 | fetch.max.wait.ms | 若数据不足 fetch.min.bytes,Broker 等待时间。 | 100 ~ 500ms。 | 避免频繁空拉取。 |
| 拉取大小 | fetch.max.bytes | 单次拉取请求的最大数据量。 | 50MB ~ 100MB(需 Broker message.max.bytes 支持)。 | 防止单次拉取过大阻塞。 |
| 会话超时 | session.timeout.ms | Consumer 与 Group Coordinator 的心跳超时。 | 10s ~ 30s。 | 过短导致误判宕机,过长故障恢复慢。 |
| 心跳间隔 | heartbeat.interval.ms | Consumer 向 Coordinator 发送心跳的间隔。 | 应小于 session.timeout.ms 的 1/3,如 3s。 | 确保心跳及时,避免不必要的再平衡。 |
| 最大轮询间隔 | max.poll.interval.ms | 两次 poll() 调用的最大间隔,超过则触发再平衡。 | 根据处理逻辑耗时设置,如 30s ~ 5min。 | 长时间处理需拆分或启用手动提交。 |
| 位移提交 | enable.auto.commit=false | 禁用自动提交,手动控制位移提交时机。 | 配合 commitSync() 或 commitAsync() 实现精确控制。 | 避免重复消费或数据丢失。 |
| 反序列化器优化 | 高效反序列化(如 Avro、Protobuf) | 减少 CPU 开销。 | 与生产者匹配。 | 避免反序列化成为瓶颈。 |
| 多线程消费 | 单 Consumer 多线程处理消息 | 解耦拉取与处理,提升 CPU 利用率。 | 注意位移提交的线程安全。 | 避免在多线程中直接提交位移。 |
7.3 Broker 配置调优
| 参数名称 | 默认值 | 说明 | 调优建议 | 注意事项 |
|---|---|---|---|---|
num.network.threads | 3 | 处理网络请求(接收)的线程数。 | 6 ~ 12,根据网络流量调整。 | 接收客户端请求。 |
num.io.threads | 8 | 执行磁盘 I/O(读写日志)的线程数。 | 16 ~ 32,通常为磁盘数的倍数。 | Kafka 依赖顺序 I/O,线程数影响吞吐。 |
socket.send.buffer.bytes | 100KB | Server 端 SO_SNDBUF 缓冲区大小。 | 1MB ~ 4MB,提升网络吞吐。 | 与客户端 socket.request.max.bytes 匹配。 |
socket.receive.buffer.bytes | 100KB | Server 端 SO_RCVBUF 缓冲区大小。 | 1MB ~ 4MB。 | 接收客户端数据。 |
queued.max.requests | 500 | 网络线程可发送给 I/O 线程的最大请求数。 | 1000 ~ 2000,缓冲更多请求。 | 防止请求堆积。 |
log.flush.interval.messages | - | 强制刷盘的消息条数间隔(不推荐使用)。 | 通常依赖 log.flush.interval.ms 或操作系统刷盘。 | 启用后影响性能,一般关闭。 |
log.flush.interval.ms | - | 强制刷盘的时间间隔(不推荐)。 | 通常不设置,依赖 log.flush.scheduler.interval.ms。 | 影响持久性和性能。 |
log.retention.hours | 168(7 天) | 消息保留时间。 | 根据业务需求设置,如 24、72 小时。 | 与 log.retention.bytes 共同作用。 |
log.segment.bytes | 1GB | 日志段文件大小。 | 512MB ~ 2GB,根据清理频率调整。 | 较小文件便于删除。 |
log.retention.check.interval.ms | 300000(5 分钟) | 检查是否需要清理日志的间隔。 | 300000 ~ 3600000(1 小时)。 | 频繁检查增加开销。 |
zookeeper.connection.timeout.ms | 6000 | Broker 连接 ZooKeeper 超时时间。 | 60000(1 分钟),网络不稳定时增大。 | 连接 ZooKeeper 失败会导致 Broker 不可用。 |
delete.topic.enable | true | 是否允许通过 API 删除 Topic。 | 生产环境建议 true,便于管理。 | 删除后数据不会立即清除。 |
7.4 Kafka 监控指标与常用工具(JMX、Prometheus、Grafana)
| 监控维度 | 关键指标(JMX MBean) | 说明 | 工具集成 |
|---|---|---|---|
| Broker | kafka.server:type=BrokerTopicMetrics,name=MessagesInPerSec | 每秒入站消息数,衡量生产负载。 | Prometheus 通过 JMX Exporter 抓取,Grafana 展示。 |
kafka.server:type=BrokerTopicMetrics,name=BytesInPerSec | 每秒入站字节数。 | ||
kafka.server:type=BrokerTopicMetrics,name=BytesOutPerSec | 每秒出站字节数,衡量消费负载。 | ||
kafka.server:type=ReplicaManager,name=UnderReplicatedPartitions | 非同步副本分区数,应为 0。 | 告警关键指标。 | |
kafka.server:type=KafkaRequestHandlerPool,name=RequestHandlerAvgIdlePercent | 请求处理线程空闲率,过低表示处理不过来。 | ||
| Producer | kafka.producer:type=producer-metrics,client-id="..." | 包括 record-send-rate、request-latency-avg、outgoing-byte-rate 等。 | 客户端监控,定位生产瓶颈。 |
| Consumer | kafka.consumer:type=consumer-fetch-manager-metrics,client-id="..." | 包括 records-lag-max(最大滞后)、fetch-rate、fetch-latency-avg 等。 | records-lag-max 是核心指标,监控消费延迟。 |
| Topic/Partition | kafka.log:type=Log,name=LogEndOffset,topic=...,partition=... | 分区 LEO,结合消费者位移计算滞后量。 | 需结合消费者组位移计算 consumer_lag。 |
| ZooKeeper | kafka.server:type=SessionExpireListener,name=ZooKeeperAuthFailures | ZooKeeper 认证失败次数。 | 监控 ZK 健康状态。 |
kafka.server:type=SessionExpireListener,name=ZooKeeperDisconnects | ZooKeeper 断开连接次数。 |
工具链:
| 工具 | 说明 |
|---|---|
| JMX Exporter + Prometheus + Grafana | 标准监控栈。部署 JMX Exporter 暴露指标,Prometheus 抓取,Grafana 可视化。 |
| Kafka Manager / Cruise Control / Confluent Control Center | 专用管理工具,提供集群视图、再平衡、监控等。Cruise Control 支持自动负载均衡。 |
告警设置:
UnderReplicatedPartitions > 0:及时发现集群异常- ISR 变化频繁
Consumer Lag > 阈值- Request Latency 骤升
使用 Prometheus Alertmanager 或工具内置告警。
7.5 日志清理策略(Log Retention)
| 策略类型 | 相关配置参数 | 说明 | 触发条件 | 注意事项 |
|---|---|---|---|---|
| 基于时间删除 | log.retention.hours(或 minutes/days)、log.roll.hours | 按消息时间戳删除过期数据。 | 消息时间戳超过 retention.ms。 | 时间基于消息的 timestamp,需确保生产者时间准确。 |
| 基于大小删除 | log.retention.bytes、log.segment.bytes | 按分区日志总大小删除最旧段。 | 分区日志总大小超过 retention.bytes。 | 与 log.retention.hours 任一满足即触发清理。 |
| 清理检查周期 | log.retention.check.interval.ms(默认 300000ms=5 分钟) | 清理服务检查日志是否可删除的频率。 | 每 check.interval.ms 执行一次清理检查。 | 频繁检查增加 CPU 开销。 |
| 日志段滚动 | log.segment.bytes、log.segment.ms | 控制何时创建新日志段文件。 | 当前段大小超过 segment.bytes 或存活时间超过 segment.ms。 | 滚动是删除的前提,小段文件便于精确清理。 |
| 删除流程 | - | 异步删除,不影响服务。 | 确定可删除段 → 标记为 .deleted → 后台线程删除文件。 | 删除不可逆。 |
| 压缩(Compact) | cleanup.policy=compact、min.compaction.lag.ms、max.compaction.lag.ms | 保留每个 key 的最新值,适用于状态存储。 | 满足 min.compaction.lag.ms 后,后台线程合并日志,保留最新 key。 | 不删除消息,仅合并;delete 和 compact 可同时启用。 |
| 混合策略 | cleanup.policy=delete,compact | 先按时间/大小删除过期数据,再对剩余数据压缩。 | 结合两种策略优点。 | 适用于需要长期保留但又需节省空间的场景。 |
| 特殊保留 | retention.ms=-1 | 永久保留,永不删除。 | 用于审计日志等。 | 需手动管理磁盘空间。 |
第八章:Java Kafka 实战案例
8.1 日志收集系统设计与实现
场景描述: 构建一个分布式日志收集系统,将多台服务器的应用日志(如 Nginx、应用日志)统一收集到 Kafka,供后续分析、监控或存储。
架构设计:
[应用服务器] → (Filebeat/Fluentd) → [Kafka Cluster] → (Kafka Consumer) → [Elasticsearch/Logstash/S3]
日志采集端:
- 使用 Filebeat 或 Fluentd 作为轻量级日志采集器,监控日志文件变化。
- 配置 Filebeat 输出到 Kafka:
output.kafka:
hosts: ["kafka-broker1:9092", "kafka-broker2:9092"]
topic: app-logs-topic
partition.round_robin:
reachable_only: false
required_acks: 1
compression: gzip
max_message_bytes: 1000000
Kafka Producer(Java,可选): 若需自定义处理(如过滤、格式化),可用 Java 编写 Producer。
Properties props = new Properties();
props.put("bootstrap.servers", "kafka1:9092,kafka2:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("acks", "1"); // 平衡吞吐与可靠性
props.put("retries", 3);
props.put("batch.size", 16384); // 16KB
props.put("linger.ms", 10); // 等待 10ms 凑批
props.put("compression.type", "gzip"); // 减少网络传输
props.put("buffer.memory", 33554432); // 32MB
Topic 设计:
- 分区数:根据日志源数量和吞吐量预估,如 10-50 个。
- 副本数:
replication.factor=3,保证高可用。 - 保留策略:
retention.ms=604800000(7 天),或按大小retention.bytes。
Consumer(Java):
- 消费日志写入 Elasticsearch 或 HDFS。
- 幂等写入:确保日志不重复(如使用 Logstash 或自定义去重逻辑)。
- 批量处理:提升写入下游存储的效率。
8.2 订单系统解耦:异步处理订单
场景描述: 电商系统中,用户下单后需执行库存扣减、积分计算、短信通知、物流准备等多个操作。使用 Kafka 解耦,提升系统响应速度和可靠性。
架构设计:
[订单服务] → (Kafka Producer) → [order-topic] → [库存服务] (Consumer)
├──→ [积分服务] (Consumer)
├──→ [通知服务] (Consumer)
└──→ [物流服务] (Consumer)
订单服务(Producer):
@Service
public class OrderService {
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
public void createOrder(Order order) {
// 1. 保存订单到数据库(本地事务)
orderRepository.save(order);
// 2. 发送订单消息到 Kafka
try {
String orderJson = objectMapper.writeValueAsString(order);
kafkaTemplate.send("order-topic", order.getOrderId().toString(), orderJson);
log.info("Order message sent: {}", order.getOrderId());
} catch (Exception e) {
log.error("Failed to send order message", e);
// 处理发送失败:重试、记录日志、或降级
}
}
}
下游服务(Consumer - 以库存服务为例):
@Component
public class InventoryConsumer {
@KafkaListener(topics = "order-topic", groupId = "inventory-group")
public void consumeOrder(String key, String orderJson,
ConsumerRecord<String, String> record) {
try {
Order order = objectMapper.readValue(orderJson, Order.class);
// 执行扣减库存逻辑
inventoryService.deductStock(order.getItems());
log.info("Inventory deducted for order: {}", order.getOrderId());
} catch (Exception e) {
log.error("Failed to process order for inventory", e);
// 关键:不要提交位移!让 Kafka 重试
// 可以记录错误日志,或发送到死信队列(DLQ)
throw e;
}
}
}
配置要点:
- Producer:
acks=all,retries=Integer.MAX_VALUE,enable.idempotence=true,确保消息不丢失。 - Consumer:
enable.auto.commit=false,手动提交位移,确保处理成功后再提交。 - 死信队列(DLQ):处理失败的消息可发送到
order-failed-topic,供人工或后台程序处理。
8.3 实时数据统计:使用 Kafka Streams
场景描述: 实时统计每分钟的订单金额总额、订单数,或用户点击流的 PV/UV。
public class RealTimeAnalytics {
public static void main(String[] args) {
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "realtime-analytics");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka1:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
StreamsBuilder builder = new StreamsBuilder();
// 1. 从订单 Topic 读取数据
KStream<String, String> orderStream = builder.stream("order-topic");
// 2. 转换为 <key, Order> 流
KStream<String, Order> orderObjectStream = orderStream.mapValues(json -> {
try {
return objectMapper.readValue(json, Order.class);
} catch (Exception e) {
throw new RuntimeException("Deserialization failed", e);
}
});
// 3. 按时间窗口聚合(例如:每分钟)
KTable<Windowed<String>, Long> orderCountTable = orderObjectStream
.groupByKey()
.windowedBy(TimeWindows.of(Duration.ofMinutes(1)))
.count();
KTable<Windowed<String>, Double> orderAmountTable = orderObjectStream
.map((key, order) -> new KeyValue<>(key, order.getAmount()))
.groupByKey()
.windowedBy(TimeWindows.of(Duration.ofMinutes(1)))
.reduce(Double::sum);
// 4. 将结果写回 Kafka Topic
orderCountTable.toStream()
.to("order-count-stats",
Produced.with(WindowedSerdes.timeWindowedSerdeFrom(String.class),
Serdes.Long()));
orderAmountTable.toStream()
.to("order-amount-stats",
Produced.with(WindowedSerdes.timeWindowedSerdeFrom(String.class),
Serdes.Double()));
KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();
// 优雅关闭
Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
}
}
优势:
- 轻量级:直接嵌入 Java 应用。
- Exactly-Once:支持精确一次处理语义。
- 状态存储:内置状态存储(RocksDB),支持复杂计算。
- 容错:自动故障恢复。
8.4 消息幂等性与去重处理
问题: 网络重试可能导致 Producer 重复发送,Consumer 处理失败重试导致重复消费。
Producer 端幂等性:
- 启用
enable.idempotence=true(Kafka 0.11+)。 - 配合
max.in.flight.requests.per.connection=1(旧版本)或5(Kafka 2.1+,支持乱序)。 - 原理:Broker 通过
<ProducerID, SequenceNumber>识别并去重。
Consumer 端去重:
- 业务层幂等: 设计幂等的业务操作。
- 订单:
order_id作为数据库唯一键。 - 扣库存:基于
order_id和item_id判断是否已扣减。 - 发短信:记录
user_id+event是否已发送。
- 订单:
- 外部存储去重: 使用 Redis 记录已处理的
message_id或业务 ID。
public void consume(ConsumerRecord<String, String> record) {
String messageId = record.key(); // 或从 value 中提取业务 ID
Boolean added = redisTemplate.opsForSet().add("processed_messages", messageId);
if (!added) {
log.info("Message already processed, skipping: {}", messageId);
return; // 跳过
}
// 处理业务逻辑
processBusiness(record.value());
// 提交位移
consumer.commitSync();
}
事务性: Consumer 与下游数据库在同一个事务中提交(较少用)。
8.5 容错与高可用架构设计
目标: 确保 Kafka 集群和应用在节点故障时仍能正常服务。
Broker 高可用:
- 多副本:
replication.factor >= 3,Leader 故障时 ISR 中 Follower 自动选举为新 Leader。 - 跨机架/可用区部署:Broker 分布在不同物理机或云可用区,避免单点故障。
- Controller 高可用:Kafka Controller 本身也是多副本,ZooKeeper 选举。
Producer 容错:
acks=all:确保消息被 ISR 所有副本确认。retries=Integer.MAX_VALUE:无限重试(结合retry.backoff.ms)。enable.idempotence=true:防止重试导致重复。- 降级策略:发送失败时,可写入本地文件或内存队列,待恢复后重放。
Consumer 容错:
- 消费者组:多个 Consumer 实例组成 Group,自动负载均衡和故障转移。
- 位移管理:手动提交(
enable.auto.commit=false),确保处理成功再提交;异常时不提交位移,由 Kafka 重试。 - 死信队列(DLQ):处理失败的消息发送到专门的 Topic,避免阻塞主流程。
网络与 ZooKeeper:
- ZooKeeper 集群:至少 3 个节点,跨机架部署。
- 网络分区处理:合理设置
session.timeout.ms和heartbeat.interval.ms,避免误判。
监控与告警:
- 监控
UnderReplicatedPartitions、IsrShrinksPerSec、RequestHandlerAvgIdlePercent。 - 设置告警,及时发现并处理故障。
总结: 通过多副本、消费者组、幂等生产者、手动位移提交、DLQ 和全面监控,构建一个健壮的 Kafka 应用系统。
第九章:常见问题与最佳实践
9.1 常见错误与排查方法
| 问题现象 | 可能原因 | 排查方法 | 解决方案 |
|---|---|---|---|
| Producer 发送失败(Timeout、NotLeaderForPartition、NetworkException) | Broker 宕机或网络不通 | telnet <broker> 9092 检查网络连通性。 | 确认 Broker 正常运行,网络畅通。 |
| Topic 不存在或分区数变化 | kafka-topics.sh --describe 查看 Topic 状态和 Leader 分布。 | 创建 Topic 或增加分区。 | |
| 网络延迟/超时设置过短 | 检查 Broker 日志(server.log)是否有错误。 | 调大 request.timeout.ms、delivery.timeout.ms。 | |
acks=all 但 ISR 副本不足 | 监控 RequestLatencyMs 和 NetworkProcessorAvgIdlePercent。 | 确保 min.insync.replicas <= replication.factor,且 ISR 数量足够。 | |
| Consumer 消费延迟(Lag 增长) | Consumer 处理速度慢 | 监控 records-lag-max 和 fetch-rate。 | 优化 Consumer 业务逻辑,减少单次处理耗时。 |
| Consumer 实例数不足 | 使用 jstat 或 APM 工具分析 Consumer 应用 GC 情况。 | 增加 Consumer 实例(需增加分区)。 | |
| GC 停顿导致心跳超时触发再平衡 | 检查 max.poll.interval.ms 是否小于实际处理时间。 | 调大 max.poll.interval.ms 或拆分长任务。 | |
| 网络瓶颈 | 监控 Broker 的 BytesOutPerSec。 | 启用批量拉取(fetch.min.bytes)。 | |
| 集群负载不均(某些 Broker/CPU/磁盘高) | 分区 Leader 分布不均 | kafka-topics.sh --describe 检查各 Broker 上的 Leader 分区数。 | 执行 kafka-replica-assignment.sh 或使用 Cruise Control 进行 Leader 重平衡。 |
| 某些 Topic 流量过大且分区少 | 监控各 Broker 的 BytesInPerSec、RequestHandlerAvgIdlePercent。 | 为热点 Topic 增加分区数。 | |
| 数据倾斜(Key 分布不均) | 分析消息 Key 的分布。 | 优化 Key 的设计,使其更均匀分布。 | |
| ZooKeeper 连接频繁中断 | ZooKeeper 集群性能瓶颈或网络抖动 | 检查 ZooKeeper 服务器的 CPU、内存、网络和磁盘 I/O。 | 优化 ZooKeeper 配置和部署(独立机器,SSD)。 |
zookeeper.session.timeout.ms 设置过短 | 查看 ZooKeeper 日志是否有 Session expired 错误。 | 适当增大 zookeeper.connection.timeout.ms 和 zookeeper.session.timeout.ms。 | |
监控 ZooKeeperDisconnects 指标。 | |||
| 磁盘空间不足 | 日志保留策略过于宽松(retention.ms 太大) | du -sh /kafka/data 检查磁盘使用。 | 调整 log.retention.hours/bytes。 |
| 消费者组停滞导致数据无法清理 | kafka-log-dirs.sh --describe 查看各分区日志大小。 | 清理无用的消费者组或重启停滞的 Consumer。 | |
| 生产速度远大于消费速度 | 监控 LogEndOffset 和消费者位移,计算滞后量。 | 提升消费能力。 |
9.2 生产环境部署建议
| 维度 | 最佳实践 |
|---|---|
| 硬件选择 |
|
| 集群规划 |
|
| 操作系统调优 |
|
| Kafka 配置 |
|
| 监控告警 | 必须集成 Prometheus + Grafana + Alertmanager。 核心告警指标:
|
| 备份与恢复 |
|
9.3 安全配置(SSL/SASL)
Kafka 提供多层安全机制:
加密(SSL/TLS):
- 作用:加密客户端与 Broker、Broker 间、Broker 与 ZooKeeper 间的通信。
- 配置:
- 生成 CA 证书和密钥。
- 为每个 Broker 签发证书。
# Broker 配置(server.properties)
listeners=SSL://:9093
security.inter.broker.protocol=SSL
ssl.keystore.location=/path/to/server.keystore.jks
ssl.keystore.password=keystore_password
ssl.key.password=key_password
ssl.truststore.location=/path/to/server.truststore.jks
ssl.truststore.password=truststore_password
ssl.client.auth=required # 开启客户端认证
Client 配置需提供信任库。
认证(SASL):
- SASL/PLAIN:简单用户名/密码认证。必须配合 SSL 使用,否则密码明文传输。
# server.properties
sasl.enabled.mechanisms=PLAIN
sasl.mechanism.inter.broker.protocol=PLAIN
// JAAS 配置文件
KafkaServer {
org.apache.kafka.common.security.plain.PlainLoginModule required
username="admin"
password="admin-secret"
user_admin="admin-secret"
user_alice="alice-secret";
};
- SASL/SCRAM:比 PLAIN 更安全,支持凭证动态增删。
- SASL/GSSAPI(Kerberos):企业级认证,集成 AD/KDC。
授权(ACL):
- 控制用户对 Topic、Consumer Group 等资源的操作权限(Read、Write、Describe、Create…)。
- 启用:
authorizer.class.name=kafka.security.authorizer.AclAuthorizer
# 示例命令
kafka-acls.sh --add --allow-principal User:alice --operation Read --topic orders-topic
kafka-acls.sh --add --allow-principal User:bob --operation Write --topic logs-topic
部署建议: 生产环境应同时启用 SSL(加密)+ SASL/SCRAM(认证)+ ACL(授权)。
9.4 版本升级与兼容性注意事项
升级路径:
- Kafka 遵循语义化版本控制,主版本号变更可能引入不兼容改动。
- 推荐:逐个 Broker 滚动升级。先升级 Follower,最后升级 Controller。
- 使用
kafka-broker-api-versions.sh检查 API 兼容性。
兼容性:
- Producer/Consumer 与 Broker:通常向后兼容。新版本 Client 可连接旧 Broker,但可能无法使用新特性。旧 Client 连接新 Broker 一般没问题。
- 消息格式:
message.format.version需与inter.broker.protocol.version协调。升级前确保所有 Broker 版本一致后再提升协议版本。 - ZooKeeper:Kafka 3.0+ 开始逐步用 KRaft 替代 ZooKeeper。未来升级需考虑元数据存储的迁移。
关键步骤:
- 备份: 备份 ZooKeeper 数据和 Kafka 日志。
- 测试: 在预发环境充分测试新版本。
- 文档: 仔细阅读官方发布说明(Release Notes),了解 Breaking Changes。
- Rollback Plan: 准备好回滚到旧版本的方案。
9.5 Kafka 与其他中间件对比(RabbitMQ、RocketMQ)
| 特性 | Apache Kafka | RabbitMQ | Apache RocketMQ |
|---|---|---|---|
| 架构模型 | 发布/订阅(基于日志) | AMQP 0.9.1,支持多种模式(Pub/Sub、Point-to-Point、Routing、RPC) | 发布/订阅,点对点 |
| 吞吐量 | 极高(10w+ msg/s) | 中等(万级别) | 高(接近 Kafka) |
| 延迟 | 毫秒级(批处理时稍高) | 微秒级(低延迟) | 毫秒级 |
| 持久化 | 基于磁盘日志,顺序 I/O,极高效 | 内存+磁盘,可持久化 | 基于磁盘日志 |
| 伸缩性 | 极强,水平扩展容易,分区是并行单元 | 较弱,队列不易扩展,镜像队列非线性扩展 | 强,类似 Kafka |
| 消息顺序 | 分区内严格有序 | 单队列有序 | Topic 内有序 |
| 消息路由 | 简单(基于 Topic/Partition) | 极其灵活(Exchange + Binding + Routing Key) | 较灵活(Tag、SQL 过滤) |
| 功能丰富度 | 核心功能稳定,生态强大(Streams、Connect) | 功能最丰富(插件、管理界面、多种协议) | 功能全面(事务消息、定时消息、死信队列) |
| 运维复杂度 | 中等偏高(依赖 ZooKeeper/KRaft) | 低(单节点简单,集群需镜像) | 中等 |
| 适用场景 | 大数据管道、日志收集、流处理、高吞吐解耦 | 传统企业应用、复杂路由、RPC、低延迟任务 | 金融、电商(阿里系),需要事务/定时消息的场景 |
选型建议:
- 选 Kafka: 追求极致吞吐、大数据集成、流式处理、作为”数据中心”。
- 选 RabbitMQ: 需要复杂路由规则、低延迟、协议多样(MQTT、STOMP)、易上手。
- 选 RocketMQ: 需要强事务消息、定时/延时消息、在中国云环境下深度集成。