第 199 章:Kafka 进阶特性
学习目标
- 掌握 Kafka 事务
- 学会消息幂等设计
- 深入分区与副本机制
- 了解 Kafka Streams
一、Kafka 幂等性
1.1 什么是幂等 Producer
默认 Producer 可能因重试导致重复消息,幂等 Producer 保证单分区单会话内不重复。
yaml
spring:
kafka:
producer:
properties:
enable.idempotence: true # 开启幂等性
acks: all
retries: 5
max.in.flight.requests.per.connection: 51.2 原理
- ProducerId:每个 Producer 实例唯一
- SequenceNumber:每条消息单调递增
- Broker 端去重:已存在的 (PID, seq) 直接丢弃
二、Kafka 事务
2.1 事务 Producer
java
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
public void sendInTransaction(List<String> messages) {
kafkaTemplate.executeInTransaction(ops -> {
for (String msg : messages) {
ops.send("topic1", msg);
}
// 原子提交:要么全成功,要么全失败
return true;
});
}2.2 多 Topic 事务
java
kafkaTemplate.executeInTransaction(ops -> {
ops.send("order-topic", order);
ops.send("audit-topic", auditLog);
ops.send("notification-topic", notification);
// 三条消息要么全成功,要么全回滚(对消费者而言)
return true;
});2.3 读已提交(read_committed)
yaml
spring:
kafka:
consumer:
properties:
isolation.level: read_committed # 只读已提交事务中的消息对消费者不可见,直到事务提交。
三、消费幂等性
3.1 业务层面幂等
java
@KafkaListener(topics = "order-events")
public void handleOrder(OrderEvent event, Acknowledgment ack) {
// 1. 唯一键去重(Redis / DB)
String dedupKey = "kafka:order:" + event.getOrderId();
if (!redisTemplate.opsForValue().setIfAbsent(dedupKey, "1", 24, TimeUnit.HOURS)) {
log.info("重复消息,跳过: {}", event.getOrderId());
ack.acknowledge();
return;
}
// 2. 业务处理
inventoryService.deduct(event);
ack.acknowledge();
}3.2 数据库唯一约束
java
@Transactional
public void handle(OrderEvent event) {
// 唯一约束兜底
try {
orderRepository.save(event.toEntity());
} catch (DuplicateKeyException e) {
log.warn("重复消息");
}
}四、消息积压处理
4.1 排查
bash
kafka-consumer-groups --bootstrap-server localhost:9092 \
--describe --group inventory-group
# 看 LAG(滞后量)
# CURRENT-OFFSET LOG-END-OFFSET LAG
# 100 10000 9900 ← 积压4.2 解决方案
扩容消费者
yaml
spring:
kafka:
listener:
concurrency: 10 # 10 个消费者线程批量消费
yaml
spring:
kafka:
consumer:
max-poll-records: 500 # 一次拉 500 条
listener:
type: batchjava
@KafkaListener(topics = "events", containerFactory = "batchFactory")
public void handleBatch(List<Event> events, Acknowledgment ack) {
// 批量处理
service.batchProcess(events);
ack.acknowledge();
}加分区
bash
kafka-topics --bootstrap-server localhost:9092 \
--alter --topic events --partitions 10五、DLQ(Dead Letter Queue)
处理失败的消息放到 DLQ,主流程不受影响。
java
@Configuration
public class KafkaConfig {
@Bean
public DefaultErrorHandler errorHandler(KafkaTemplate<String, String> template) {
// 失败 3 次后发到 DLQ
DeadLetterPublishingRecoverer recoverer =
new DeadLetterPublishingRecoverer(template,
(record, ex) -> new TopicPartition("events.DLT", record.partition()));
return new DefaultErrorHandler(recoverer, new FixedBackOff(1000L, 3));
}
@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(
ConsumerFactory<String, String> consumerFactory,
DefaultErrorHandler errorHandler) {
ConcurrentKafkaListenerContainerFactory<String, String> factory =
new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory);
factory.setCommonErrorHandler(errorHandler);
return factory;
}
}java
@KafkaListener(topics = "events.DLT")
public void handleDLT(String message) {
log.error("DLQ 消息: {}", message);
// 人工处理或重新入队
}六、分区策略
6.1 默认 Hash
java
// KafkaTemplate.send(topic, key, value)
// key 为 null → 轮询
// key 不为 null → hash(key) % partitions6.2 自定义分区器
java
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();
if (key == null) {
return ThreadLocalRandom.current().nextInt(numPartitions);
}
// 自定义规则:VIP 用户到独立分区
String keyStr = key.toString();
if (keyStr.startsWith("vip_")) {
return 0; // VIP 分区
}
return Math.abs(keyStr.hashCode()) % numPartitions;
}
}七、Kafka Streams
7.1 引入
xml
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-streams</artifactId>
</dependency>7.2 单词计数
java
@Configuration
public class StreamsConfig {
@Bean
public KStream<String, Long> wordCount(StreamsBuilder builder) {
KStream<String, String> source = builder.stream("text-input");
KTable<String, Long> counts = source
.flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\s+")))
.groupBy((key, value) -> value)
.count();
counts.toStream().to("word-count-output", Produced.with(Serdes.String(), Serdes.Long()));
return source;
}
}八、Spring Cloud Stream(屏蔽底层)
xml
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-stream-binder-kafka</artifactId>
</dependency>yaml
spring:
cloud:
stream:
bindings:
output:
destination: orders
input:
destination: orders
group: inventory-group
kafka:
binder:
brokers: localhost:9092java
// 生产
@Service
public class OrderService {
@Autowired
private StreamBridge streamBridge;
public void send(OrderEvent event) {
streamBridge.send("output", event);
}
}
// 消费
@Component
public class InventoryConsumer {
@Bean
public Consumer<OrderEvent> input() {
return event -> inventoryService.deduct(event);
}
}九、Kafka 监控
9.1 JMX 指标
| 指标 | 含义 |
|---|---|
kafka.server:type=BrokerTopicMetrics,name=MessagesInPerSec | 入速率 |
kafka.server:type=ReplicaManager,name=UnderReplicatedPartitions | 副本不足分区 |
kafka.server:type=BrokerTopicMetrics,name=BytesInPerSec | 入字节 |
kafka.network:type=RequestMetrics,name=TotalTimeMs | 请求延迟 |
9.2 关键告警
promql
# Lag 告警
kafka_consumergroup_lag{group="my-group"} > 10000
# 副本不足
sum(kafka_topic_partition_in_sync_replica_count) by (topic) == 0
# 磁盘使用率
100 - (kafka_log_log_size / kafka_log_log_capacity * 100) < 10十、性能调优
10.1 Producer 调优
yaml
spring:
kafka:
producer:
batch-size: 65536 # 64KB 批量
linger-ms: 10 # 等 10ms 凑批
compression-type: snappy # 压缩
buffer-memory: 67108864 # 64MB 缓冲区10.2 Broker 调优
properties
# 内存
KAFKA_HEAP_OPTS=-Xms4g -Xmx4g
# 磁盘
log.dirs=/data/kafka
log.retention.hours=72
# 网络
num.network.threads=8
num.io.threads=16
# 分区数预估
# 期望吞吐量 / 单分区吞吐 ≈ 分区数10.3 Consumer 调优
yaml
spring:
kafka:
consumer:
fetch.min.bytes: 1024 # 至少 1KB 才返回
fetch.max.wait.ms: 500 # 最多等 500ms
max.partition.fetch.bytes: 1048576 # 1MB十一、本章小结
| 特性 | 解决问题 |
|---|---|
| 幂等性 | Producer 重试重复 |
| 事务 | 跨 topic 原子性 |
| DLQ | 失败消息隔离 |
| 批量消费 | 提高吞吐 |
| Kafka Streams | 实时流处理 |
动手练习
- 实现幂等 Producer(防重复)
- 用事务保证多 topic 原子性
- 配 DLQ 处理失败消息
- 监控 Consumer Lag 并告警
推荐阅读
下一章:第 200 章:RabbitMQ 入门与实战 →