Kafka 原理与实战
Kafka 概述
LinkedIn 开源背景
Apache Kafka 最初由 LinkedIn 开发,于 2011 年开源并捐赠给 Apache 软件基金会。LinkedIn 当时面临一个典型的大数据基础设施问题:需要将来自不同系统的海量数据(用户活动、系统指标、日志等)实时地、可靠地传输到多个下游系统中。传统的点对点集成方案难以应对数据量增长和系统多样性的挑战。
Kafka 的设计目标与 LinkedIn 的需求紧密相关:
- 高吞吐量:每秒处理数十万条消息
- 可持久化:消息写入磁盘,支持数据回放
- 分布式:天然支持水平扩展
- 多订阅者:一条消息可被多个消费者独立消费
Kafka 的设计理念深受 日志(Log) 概念的影响——将数据流视为仅追加的日志,消费者只需维护自己的读取偏移量即可。
事件流平台概念
Kafka 本质上是一个分布式事件流平台(Event Streaming Platform),其核心能力包括:
| 能力 | 描述 |
|---|---|
| 发布/订阅 | 像消息队列一样发布和订阅事件流 |
| 持久化存储 | 可靠地持久化事件流,可无限期保留 |
| 流式处理 | 在事件发生时实时处理和分析 |
与传统的消息队列(如 RabbitMQ、ActiveMQ)不同,Kafka 的设计更侧重日志抽象和大吞吐量,消息被持久化到磁盘,消费者可以控制读取位置,这使得它既能做消息中间件,也能做数据存储层。
应用场景
Kafka 在企业级架构中的典型应用场景包括:
- 日志聚合:收集分布式服务的日志,统一输送到日志系统(如 ELK Stack)
- 指标监控:收集系统和业务指标数据,实时推送至监控系统(如 Prometheus、Grafana)
- 微服务异步解耦:服务间通过事件驱动通信,降低耦合
- 事件溯源(Event Sourcing):以事件序列的形式记录业务状态变更
- CDC(Change Data Capture):通过 Debezium 等工具捕获数据库变更事件
- 流式 ETL:对实时数据流进行清洗、转换、聚合后写入目标存储
- 消息队列:作为高吞吐量的消息中间件,连接生产者和消费者
架构
Kafka 的整体架构由多个组件协作构成,下图展示了核心组件及其关系。
+-----------+
| Producer |
+-----+-----+
|
+-----v-----+
| Broker | <-- Kafka 集群
| +------+ |
| |Topic | |
| |/Part. | |
| +------+ |
+-----+-----+
|
+---------------v---------------+
| Consumer Group |
| Consumer 1 | Consumer 2 | ... |
+--------------------------------+Producer(生产者)
生产者负责向 Kafka Topic 发送消息。生产者可以指定消息发往哪个分区,也可以由分区器根据 Key 自动决定。生产者具备重试、幂等、事务等机制。
Broker(代理服务器)
Broker 是 Kafka 集群中的服务器节点。每个 Broker 都有一个唯一 ID,可以同时管理多个 Partition。一个 Kafka 集群通常由多个 Broker 组成,用于实现负载均衡和高可用。
Consumer(消费者)
消费者从 Broker 拉取消息。消费者以消费组(Consumer Group) 的形式组织,同一个消费组内的消费者共同消费一个 Topic 的消息,每个分区只能由组内的一个消费者消费。
Topic(主题)
Topic 是消息的逻辑分类。生产者将消息发送到特定的 Topic,消费者从 Topic 订阅消息。一个 Topic 可以被多个消费组独立消费。
Partition(分区)
Partition 是 Topic 的物理分片。每个 Topic 可以包含多个 Partition,消息被均匀分布在各个 Partition 中。Partition 内部的消息是有序的。
Offset(偏移量)
Offset 是分区内消息的唯一标识,是一个单调递增的整数。消费者消费消息时,通过 Offset 记录消费进度,支持从任意位置开始消费。
Zookeeper / KRaft
Kafka 早期依赖 ZooKeeper 进行集群元数据管理(Broker 注册、Topic 元数据、Leader 选举等)。从 Kafka 2.8 开始引入了 KRaft 模式,通过基于 Raft 协议的共识机制替代 ZooKeeper,实现了去 ZooKeeper 依赖。
| 特性 | ZooKeeper 模式 | KRaft 模式 |
|---|---|---|
| 依赖 | 需要独立部署 ZooKeeper | 无需额外依赖 |
| 元数据存储 | ZooKeeper 树形节点 | KRaft 日志 |
| Controller 选举 | ZooKeeper 选举 | Raft 共识 |
| 成熟度 | 生产级稳定 | 较新,持续完善中 |
分区机制
分区策略
Kafka 的分区设计是其高吞吐量的核心。每个 Topic 被划分为多个 Partition,每个 Partition 是一个有序的、不可变的日志序列。
分区的作用:
- 水平扩展:Partition 分布在多个 Broker 上,突破单机瓶颈
- 并行度:多个消费者可并行消费不同 Partition
- 有序性保证:单个 Partition 内消息有序
分区写入
生产者消息到 Partition 的映射策略由分区器(Partitioner) 决定:
默认分区器逻辑:
// 如果消息指定了分区号,直接使用
if (record.partition() != null) {
return record.partition();
}
// 如果消息有 Key,对 Key 哈希取模
if (record.key() != null) {
return Utils.toPositive(Utils.murmur2(record.key())) % numPartitions;
}
// 无 Key:Round-Robin 或 Sticky 策略
return nextBatchIndex.getAndIncrement() % numPartitions;自定义分区器示例:
public class OrderIdPartitioner implements Partitioner {
@Override
public int partition(String topic, Object key, byte[] keyBytes,
Object value, byte[] valueBytes, Cluster cluster) {
if (key == null) return 0;
String orderId = key.toString();
// 按订单号前缀分桶
String prefix = orderId.substring(0, 4);
return Math.abs(prefix.hashCode()) % cluster.partitionCountForTopic(topic);
}
@Override
public void close() {}
@Override
public void configure(Map<String, ?> configs) {}
}分区读取
消费者从 Partition 读取消息时,遵循以下原则:
- 一个 Partition 只能被同一个消费组内的一个消费者消费
- 消费者通过 Offset 控制读取位置
- 支持顺序消费(单个分区内)和消息回放(重置 Offset)
指定 Offset 消费:
// 从指定 Offset 开始消费
consumer.assign(Arrays.asList(topicPartition));
consumer.seek(topicPartition, 100L);
// 从分区开头消费
consumer.seekToBeginning(Arrays.asList(topicPartition));
// 从分区末尾消费
consumer.seekToEnd(Arrays.asList(topicPartition));Leader / Follower
每个 Partition 在物理存储上存在多个副本,其中一个为 Leader,其余为 Follower。
- Leader:负责处理所有读写请求
- Follower:从 Leader 同步数据,作为热备
ISR 机制
ISR(In-Sync Replica) 是与 Leader 保持同步的副本集合。Kafka 通过 ISR 机制实现一致性保证:
- 只有 ISR 集合中的 Follower 才有资格在 Leader 宕机时被选为新 Leader
- Follower 若同步落后超过
replica.lag.time.max.ms(默认 30s),将被踢出 ISR - ISR 机制允许在可用性和一致性之间做权衡
# ISR 相关配置
replica.lag.time.max.ms=30000
# 最小 ISR 数量,配合 acks=all 使用
min.insync.replicas=2副本与高可用
ISR(In-Sync Replica)
ISR(In-Sync Replica)是 Kafka 实现高可用的核心机制。ISR 包含 Leader 以及所有与 Leader 保持同步的 Follower。
ISR 判定条件:
Follower 与 Leader 的最后同步时间差 < replica.lag.time.max.ms当 Follower 满足以下条件时,被视为"同步":
- 在配置的超时时间内成功拉取过 Leader 的数据
- 没有落后 Leader 太多数据
ACK 配置(0 / 1 / -1)
生产者通过 acks 参数控制消息的可靠性级别:
| acks 值 | 行为 | 可靠性 | 延迟 | 适用场景 |
|---|---|---|---|---|
acks=0 | 生产者发送后不等待确认 | 最低:可能丢消息 | 最低 | 日志、监控等可丢失场景 |
acks=1 | Leader 写入成功后即确认 | 中:Leader 宕机可能丢 | 中 | 一般业务场景(默认值) |
acks=-1(或 all) | Leader 和所有 ISR 都写入后才确认 | 最高:不丢消息 | 最高 | 金融、订单等关键数据 |
配置示例:
Properties props = new Properties();
props.put("acks", "all"); // 最可靠的确认模式
props.put("retries", 3); // 重试次数
props.put("min.insync.replicas", 2); // 最小同步副本数
props.put("enable.idempotence", true); // 启用幂等重要说明: 当 acks=all 且 min.insync.replicas 配置不满足时,生产者会收到 NotEnoughReplicasException。
Leader 选举
当 Partition Leader 宕机时,Kafka Controller 负责从 ISR 中选择一个 Follower 作为新 Leader。
选举流程:
- Controller 检测到 Leader 心跳超时
- 从 ISR 中选取第一个副本作为新 Leader
- 所有副本更新元数据
- 生产者、消费者重新连接新 Leader
如果 ISR 为空(全部不可用):
| 配置 | 行为 |
|---|---|
unclean.leader.election.enable=false(默认) | 不允许非 ISR 副本成为 Leader,分区不可用 |
unclean.leader.election.enable=true | 允许非 ISR 副本成为 Leader,可能丢数据 |
# server.properties
unclean.leader.election.enable=false生产者
发送流程
Kafka 生产者的消息发送流程分为以下几个步骤:
Producer -> Serializer -> Partitioner -> RecordAccumulator -> Sender Thread -> Network -> Broker- 拦截器(Interceptor):对消息进行预处理或修改
- 序列化器(Serializer):将 Key 和 Value 转为字节数组
- 分区器(Partitioner):决定消息发往哪个 Partition
- RecordAccumulator:将消息批量缓存到内存队列
- Sender 线程:从 Accumulator 拉取批次数据,发送到 Broker
序列化器
Kafka 支持内置序列化器和自定义序列化器。
内置序列化器:
key.serializer=org.apache.kafka.common.serialization.StringSerializer
value.serializer=org.apache.kafka.common.serialization.StringSerializer自定义序列化器示例:
public class UserSerializer implements Serializer<User> {
@Override
public byte[] serialize(String topic, User data) {
if (data == null) return null;
byte[] bytes = new byte[8 + data.getName().getBytes().length];
ByteBuffer buffer = ByteBuffer.wrap(bytes);
buffer.putLong(data.getId());
buffer.put(data.getName().getBytes());
return buffer.array();
}
}推荐使用 Avro / Protobuf 格式:
value.serializer=io.confluent.kafka.serializers.KafkaAvroSerializer
schema.registry.url=http://schema-registry:8081分区器
生产者通过分区器决定消息写入哪个 Partition。Kafka 提供了默认的 DefaultPartitioner,同时也支持自定义分区器。
# 自定义分区器
partitioner.class=com.example.OrderIdPartitioner重试
生产者支持自动重试,由 retries 和 retry.backoff.ms 参数控制:
Properties props = new Properties();
props.put("retries", Integer.MAX_VALUE); // 无限重试
props.put("retry.backoff.ms", 500); // 重试间隔
props.put("max.in.flight.requests.per.connection", 5);注意: 当
max.in.flight.requests.per.connection > 1且启用重试时,可能导致消息乱序。启用幂等后可以解决此问题。
幂等
幂等生产者确保消息在发送到 Broker 时不会被重复持久化,即使生产者重试也不会产生重复消息。
props.put("enable.idempotence", true);幂等原理:
- 每个生产者会话分配一个唯一的 Producer ID(PID)
- 每条消息分配一个单调递增的 Sequence Number
- Broker 根据 PID + Sequence Number 去重
启用幂等后的约束:
max.in.flight.requests.per.connection≤ 5(默认值 5)retries > 0acks必须为all
事务
Kafka 事务支持跨分区、跨 Topic 的原子写入,确保多条消息要么全部成功,要么全部失败。
// 初始化事务
producer.initTransactions();
// 事务内发送
try {
producer.beginTransaction();
producer.send(new ProducerRecord<>("topic1", "key1", "value1"));
producer.send(new ProducerRecord<>("topic2", "key2", "value2"));
producer.commitTransaction();
} catch (Exception e) {
producer.abortTransaction();
}事务配置:
transactional.id=my-transactional-id-001
enable.idempotence=true关键点:
transactional.id必须唯一,生产者宕机重启后仍能保证事务状态一致。
消费者
消费组
消费者通过 消费组(Consumer Group) 机制实现并行消费和容错。
核心规则:
- 同一个消费组内的消费者共同消费一个 Topic
- 每个 Partition 只能由组内的一个消费者消费
- 组内的消费者数量不应超过分区数(多余的消费者会空闲)
Topic: my-topic (6 partitions)
消费组A (3 consumers): 消费组B (2 consumers):
Consumer-1: P0, P1 Consumer-1: P0, P1, P2
Consumer-2: P2, P3 Consumer-2: P3, P4, P5
Consumer-3: P4, P5Rebalance
Rebalance 是指 Partition 在消费组成员间重新分配的过程。当消费者加入或离开消费组时触发。
触发条件:
- 消费者加入或离开消费组
- Topic 分区数发生变化
- 消费者心跳超时
Rebalance 类型:
| 类型 | 机制 | 说明 |
|---|---|---|
| Eager Rebalance | 全部暂停 → 重新分配 | 所有消费者停止消费,STW 风格(旧版本默认) |
| Cooperative Rebalance | 增量式重新分配 | 分批次移动分区,尽量减少中断(新版本推荐) |
配置优雅 Rebalance 的重要参数:
# 心跳间隔,影响 Rebalance 检测速度
heartbeat.interval.ms=3000
# 两次心跳之间的最大时间
session.timeout.ms=45000
# Rebalance 时最大等待时间
max.poll.interval.ms=300000位移提交(自动 / 手动)
位移提交是消费者记录当前消费进度的机制。
自动提交:
Properties props = new Properties();
props.put("enable.auto.commit", "true");
props.put("auto.commit.interval.ms", "5000"); // 每 5 秒自动提交风险: 自动提交可能导致重复消费(两次提交间隔内消费者宕机)或消息丢失(处理时间超过提交间隔)。
手动提交:
// 同步提交(阻塞)
consumer.commitSync();
// 异步提交(非阻塞,有回调)
consumer.commitAsync((offsets, exception) -> {
if (exception != null) {
log.error("提交失败: {}", offsets, exception);
}
});
// 组合提交:异步 + 同步兜底
try {
consumer.commitAsync();
} finally {
consumer.commitSync(); // 保证最终提交成功
}消费模式
Kafka 支持三种消费模式:
1. 自动分配(Subscribe):
consumer.subscribe(Arrays.asList("my-topic"));2. 手动分配(Assign):
TopicPartition tp = new TopicPartition("my-topic", 0);
consumer.assign(Arrays.asList(tp));3. 指定 Offset 消费:
// 重置到最早
consumer.seekToBeginning(Arrays.asList(tp));
// 重置到最晚
consumer.seekToEnd(Arrays.asList(tp));
// 重置到指定时间
Map<TopicPartition, Long> timestampsToSearch = new HashMap<>();
timestampsToSearch.put(tp, System.currentTimeMillis() - 3600 * 1000); // 1小时前
Map<TopicPartition, OffsetAndTimestamp> offsets = consumer.offsetsForTimes(timestampsToSearch);
if (offsets.get(tp) != null) {
consumer.seek(tp, offsets.get(tp).offset());
}完整消费者示例:
public class SimpleConsumer {
public static void main(String[] args) {
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "my-group");
props.put("key.deserializer",
"org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer",
"org.apache.kafka.common.serialization.StringDeserializer");
props.put("enable.auto.commit", "false");
props.put("auto.offset.reset", "earliest");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("my-topic"));
try {
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
for (ConsumerRecord<String, String> record : records) {
System.out.printf("offset=%d, key=%s, value=%s%n",
record.offset(), record.key(), record.value());
}
consumer.commitAsync();
}
} finally {
consumer.close();
}
}
}存储机制
Segment(分段日志)
Kafka 将每个 Partition 的日志分为多个 Segment,每个 Segment 是一个顺序文件。
partition-0/
├── 00000000000000000000.log # Segment 日志文件
├── 00000000000000000000.index # 偏移量索引文件
├── 00000000000000000000.timeindex # 时间戳索引文件
├── 00000000000001000000.log # 下一个 Segment
├── 00000000000001000000.index
└── 00000000000001000000.timeindexSegment 生命周期:
- 消息先写入当前活跃的 Segment
- 当 Segment 达到
log.segment.bytes(默认 1GB)或log.segment.ms限制时,滚动到新 Segment - 旧的 Segment 根据保留策略被清理或删除
相关配置:
# Segment 滚动大小
log.segment.bytes=1073741824
# Segment 滚动时间
log.segment.ms=604800000
# 日志保留时长
log.retention.hours=168
# 日志保留大小
log.retention.bytes=-1索引文件
Kafka 为每个 Segment 维护索引文件,用于快速定位消息。
偏移量索引(.index):
相对偏移量 | 物理位置
0 | 0
100 | 4096
200 | 8192索引采用稀疏存储,每隔 log.index.interval.bytes(默认 4096 字节)写入一个索引条目。
索引查找流程:
给定 Offset = 150
1. 二分查找找到所属 Segment(00000000000000000000.log 段)
2. 在 .index 文件中二分查找 ≤150 的索引条目
3. 找到 (100, 4096),从物理位置 4096 开始顺序扫描时间戳索引(.timeindex):
时间戳 | 偏移量
1609459200000 | 100
1609459200100 | 200用于按时间戳查找消息。
日志清理(delete / compact)
Kafka 支持两种日志清理策略:
1. 删除策略(delete)—— 默认策略
根据时间或大小清理旧的 Segment:
# 启用删除策略(默认)
log.cleanup.policy=delete
# 保留 7 天
log.retention.hours=168
# 保留最大 10GB
log.retention.bytes=107374182402. 压缩策略(compact)
保留每个 Key 的最新值,清理旧版本:
# 启用压缩策略
log.cleanup.policy=compact
# 最少保留的脏数据比例
log.cleaner.min.cleanable.ratio=0.5压缩原理:
压缩前: 压缩后:
Key=1, Value=A (offset=0) Key=1, Value=D (offset=3)
Key=1, Value=B (offset=1) Key=2, Value=Y (offset=5)
Key=2, Value=X (offset=2) Key=3, Value=Z (offset=4)
Key=1, Value=D (offset=3)
Key=3, Value=Z (offset=4)
Key=2, Value=Y (offset=5)适用场景: 配置变更日志、数据库 CDC 事件、需要保留每个 Key 最新状态的场景。
混合策略:
# 先按大小删除,再对剩余数据做压缩
log.cleanup.policy=[delete,compact]Kafka 监控
JMX 指标
Kafka 通过 JMX(Java Management Extensions)暴露丰富的运行时指标,可通过 JMX 客户端(如 JConsole、JVisualVM)或 Prometheus JMX Exporter 采集。
关键 Broker 指标:
| 指标 | MBean 名称 | 说明 |
|---|---|---|
| 消息流入速率 | kafka.server:type=BrokerTopicMetrics,name=MessagesInPerSec | 每秒消息数 |
| 字节流入 | kafka.server:type=BrokerTopicMetrics,name=BytesInPerSec | 每秒流入字节数 |
| 字节流出 | kafka.server:type=BrokerTopicMetrics,name=BytesOutPerSec | 每秒流出字节数 |
| 分区数 | kafka.server:type=ReplicaManager,name=PartitionCount | Broker 管理的分区数 |
| Leader 数 | kafka.server:type=ReplicaManager,name=LeaderCount | Broker 管理的 Leader 分区数 |
| ISR 收缩 | kafka.server:type=ReplicaManager,name=IsrShrinksPerSec | ISR 收缩速率 |
| ISR 扩张 | kafka.server:type=ReplicaManager,name=IsrExpandsPerSec | ISR 扩张速率 |
| 请求队列大小 | kafka.network:type=RequestChannel,name=RequestQueueSize | 待处理请求数 |
消费者指标:
| 指标 | 说明 |
|---|---|
records-lag-max | 消费者最大滞后量 |
fetch-rate | 拉取请求速率 |
bytes-consumed-rate | 字节消费速率 |
启用 JMX 监控:
# 启动 Kafka 时开启 JMX
export JMX_PORT=9999
export KAFKA_JMX_OPTS="-Dcom.sun.management.jmxremote \
-Dcom.sun.management.jmxremote.authenticate=false \
-Dcom.sun.management.jmxremote.ssl=false \
-Djava.rmi.server.hostname=localhost"Prometheus + Grafana 方案:
# Prometheus JMX Exporter 配置片段
---
startDelaySeconds: 0
hostPort: localhost:9999
username:
password:
rules:
- pattern: kafka.server<type=BrokerTopicMetrics, name=BytesInPerSec><>Count
name: kafka_bytes_in_total
type: GAUGEKafka Eagle
Kafka Eagle(现更名为 EFAK,Eagle For Apache Kafka)是一个开源 Web 监控系统,提供:
- 集群健康状态监控
- Topic 和消费者可视化
- 分区偏移量和 Lag 监控
- 告警通知(邮件、钉钉等)
- 历史数据趋势图表
部署配置示例(conf/system-config.properties):
# 集群别名
kafka.eagle.zk.cluster.alias=cluster1
# ZooKeeper 连接
cluster1.zk.list=zk1:2181,zk2:2181,zk3:2181
# 数据库存储(SQLite 或 MySQL)
kafka.eagle.driver=org.sqlite.JDBC
kafka.eagle.url=jdbc:sqlite:/opt/kafka-eagle/db/ke.dbBMF(Burrow + Monitor Framework)
BMF 是 LinkedIn 开源的消费者 Lag 监控方案,以 Burrow 为核心。Burrow 是一个独立的消费者偏移量检查服务:
- 多集群支持:同时监控多个 Kafka 集群
- 无需配置消费者:直接从 Kafka 内部消费组协调机制获取数据
- Lag 评估模型:基于滑动窗口计算消费者的健康状况,输出
NOTFOUND/OK/WARNING/ERROR/STOP状态
Burrow 配置示例:
[zookeeper]
hostname="zk1:2181,zk2:2181"
timeout=6
[kafka "production"]
broker="broker1:9092,broker2:9092"
offset.storage="kafka"
group-whitelist=".*"Cruise Control
Cruise Control 是 LinkedIn 开源的 Kafka 集群自动化运维工具,核心能力包括:
- 集群负载均衡:自动调整 Partition 和 Leader 分布
- 容量规划:基于指标预测集群容量
- 异常检测:检测 Broker 故障和不平衡
- 运维窗口:支持指定时间段执行运维操作
# 运行负载均衡建议
./kafka-cruise-control-start.sh config/cruisecontrol.properties 9090
# API 获取集群状态
curl http://localhost:9090/kafkacruisecontrol/stateSpring Boot 集成
依赖配置
Spring Boot 通过 spring-kafka 项目提供 Kafka 支持:
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>application.yml 配置
spring:
kafka:
bootstrap-servers: localhost:9092
# 生产者配置
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.apache.kafka.common.serialization.StringSerializer
acks: all
retries: 3
properties:
enable.idempotence: true
max.in.flight.requests.per.connection: 5
# 消费者配置
consumer:
group-id: my-group
auto-offset-reset: earliest
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
enable-auto-commit: false
# 监听器配置
listener:
ack-mode: manual_immediate
concurrency: 3生产者代码示例
@Component
public class KafkaProducerService {
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
public void sendMessage(String topic, String message) {
kafkaTemplate.send(topic, message);
}
public void sendMessage(String topic, String key, String value) {
kafkaTemplate.send(topic, key, value);
}
public SendResult<String, String> sendSync(String topic, String message) {
try {
return kafkaTemplate.send(topic, message).get(5, TimeUnit.SECONDS);
} catch (Exception e) {
throw new RuntimeException("发送消息失败", e);
}
}
// 发送带回调的消息
public void sendWithCallback(String topic, String message) {
ListenableFuture<SendResult<String, String>> future = kafkaTemplate.send(topic, message);
future.addCallback(
result -> log.info("发送成功: offset={}", result.getRecordMetadata().offset()),
ex -> log.error("发送失败", ex)
);
}
// 事务消息
@KafkaTransaction
public void sendInTransaction(String topic, String message) {
kafkaTemplate.send(topic, message);
}
}消费者代码示例
@Component
public class KafkaConsumerService {
// 简单消费
@KafkaListener(topics = "my-topic", groupId = "my-group")
public void listen(String message) {
log.info("接收消息: {}", message);
}
// 手动提交偏移量
@KafkaListener(topics = "my-topic", groupId = "my-group")
public void listenWithAck(ConsumerRecord<String, String> record,
Acknowledgment ack) {
try {
process(record.value());
ack.acknowledge(); // 手动提交
} catch (Exception e) {
log.error("处理消息失败", e);
// nack 后移交给重试或 DLT
ack.nack(3000); // 3 秒后重试
}
}
// 批量消费
@KafkaListener(topics = "my-topic", groupId = "my-group")
public void listenBatch(List<ConsumerRecord<String, String>> records,
Acknowledgment ack) {
for (ConsumerRecord<String, String> record : records) {
process(record.value());
}
ack.acknowledge();
}
// 异常处理与重试
@DltHandler
public void handleDlt(ConsumerRecord<String, String> record) {
log.error("死信队列消息: key={}, value={}", record.key(), record.value());
}
@RetryableTopic(
attempts = "3",
backoff = @Backoff(value = 1000, multiplier = 2.0),
dltTopicSuffix = "-dlt"
)
@KafkaListener(topics = "my-topic", groupId = "my-group")
public void listenWithRetry(String message) {
throw new RuntimeException("模拟处理异常");
}
}自定义配置类
@Configuration
public class KafkaConfig {
@Bean
public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>>
batchFactory(ConsumerFactory<String, String> consumerFactory) {
ConcurrentKafkaListenerContainerFactory<String, String> factory =
new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory);
factory.setConcurrency(3);
factory.setBatchListener(true);
factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL);
return factory;
}
@Bean
public KafkaTemplate<String, String> kafkaTemplate(
ProducerFactory<String, String> producerFactory) {
return new KafkaTemplate<>(producerFactory);
}
}最佳实践
分区数设置
分区数是 Kafka 性能调优的关键参数,设置时需综合考虑以下因素:
| 因素 | 建议 |
|---|---|
| 吞吐量需求 | 分区数 = 目标吞吐量 / 单个分区吞吐量 |
| 消费者并行度 | 分区数 ≥ 消费者数(多余的消费者会空闲) |
| 文件句柄数 | 每个分区对应多个文件,分区过多会消耗文件句柄 |
| Leader 选举 | 分区数越多,Leader 选举时间越长 |
| 端到端延迟 | 分区数过多可能增加延迟 |
推荐参考:
# 一般建议:分区数 = max(消费者数 × 2, 目标 TPS / 10MBps 单分区吞吐)
# 单个 Topic 分区数建议不超过 1000
num.partitions=3 # Topic 默认分区数副本因子
副本因子决定数据冗余程度:
# 生产环境建议
replication.factor=3
# 最低 ISR 数量
min.insync.replicas=2| 副本数 | 优势 | 劣势 |
|---|---|---|
| 1 | 节省存储,最高性能 | 单点故障,数据易丢 |
| 2 | 折中方案 | 一台宕机后 ISR 只剩 1,无法容忍再丢失 |
| 3(推荐) | 可容忍 1 台 Broker 故障 | 需 3 倍存储空间 |
| 3+ | 更高容错性 | 存储和网络开销大 |
消息大小
# Broker 端最大消息大小
message.max.bytes=10485760 # 10MB
# 生产者端最大请求大小
max.request.size=10485760 # 10MB
# 消费者端最大拉取大小
fetch.message.max.bytes=10485760 # 10MB原则: Broker、Producer、Consumer 三端的消息大小限制必须保持一致。
消息大小的选择权衡:
- 小于 1KB 的消息:批量效率低,建议合并发送
- 1KB~1MB 的消息:最通用的范围
- 大于 10MB 的消息:建议拆分或使用对象存储(如 S3)并发送引用链接
重试策略
生产者重试配置:
# 基础重试
retries=3
retry.backoff.ms=500
# 启用幂等防止重复
enable.idempotence=true
# 控制重试乱序
max.in.flight.requests.per.connection=5消费者重试与死信队列:
建议消费者采用指数退避 + 死信队列策略:
| 重试次数 | 等待时间 | 说明 |
|---|---|---|
| 第 1 次 | 1 秒 | 立即重试 |
| 第 2 次 | 2 秒 | 加倍等待 |
| 第 3 次 | 4 秒 | 继续加倍 |
| 第 4 次 | 8 秒 | ... |
| 超过最大次数 | 转入 DLT | 人工介入处理 |
Spring Boot 重试示例:
@RetryableTopic(
attempts = "4",
backoff = @Backoff(delay = 1000, multiplier = 2.0),
dltTopicSuffix = "-dlt"
)
@KafkaListener(topics = "my-topic")
public void processMessage(String message) {
// 业务处理
}其他重要实践
1. 消息 Key 设计
- 需要顺序保证的消息使用相同的 Key
- 避免 Key 分布极度不均匀导致数据倾斜
- 必要时增加随机前缀进行分区
2. 幂等消费
消费者侧实现幂等性,可以通过以下方式:
- 利用数据库唯一索引去重
- 使用 Redis 记录已处理消息 ID
- 在业务主键上做幂等
3. Topic 命名规范
{应用}.{领域}.{事件类型}
示例:
order.payment.paid
user.account.registered4. 配置管理
- Topic 配置通过
kafka-configs.sh动态调整,避免重启 - 关键配置(如
min.insync.replicas)设置后不可自动降低
5. 容量规划
磁盘空间 = 消息流入速率 × 保留时长 × 副本因子 × 1.2(缓冲)
示例(日 1TB 数据,保留 7 天,3 副本):
1TB × 7 × 3 × 1.2 = 25.2 TB 集群总存储6. 操作系统优化
# 页缓存刷盘间隔(默认 5s,减少刷盘频率)
sysctl -w vm.dirty_ratio=80
sysctl -w vm.dirty_background_ratio=5
# 文件描述符限制
ulimit -n 100000
# 网络缓冲区
sysctl -w net.core.rmem_max=16777216
sysctl -w net.core.wmem_max=167772167. 版本升级策略
Kafka 版本升级建议遵循:0.10.x → 0.11.x → 1.x → 2.x → 3.x
避免跨大版本升级,推荐逐步升级 Broker 并保持兼容。参考来源: