Kafka 集成
约 874 字大约 3 分钟
布欧-Lewyon
2026-05-17
首页 › Spring Boot › 消息与集成 › Kafka 集成
Kafka 是基于分布式日志(Commit Log)的消息系统,核心模型是 Topic → Partition → Offset。消息顺序追加到 Partition,消费者通过 Offset 控制读取位置,已消费的消息按保留策略持久化不删除。适合高吞吐事件流、日志采集、数据管道、事件溯源等场景。
快速入门
依赖
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-kafka</artifactId>
</dependency>配置
spring:
kafka:
bootstrap-servers: localhost:9092
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
acks: all
properties:
enable.idempotence: true
consumer:
group-id: order-group
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
properties:
spring.json.trusted.packages: "com.example.demo.dto"
enable-auto-commit: false
auto-offset-reset: earliest发送消息
@Component
public class OrderPublisher {
private final KafkaTemplate<String, Object> kafkaTemplate;
public CompletableFuture<SendResult<String, Object>> send(Order order) {
// 用订单号作 key→同一订单路由到同一 Partition,保证顺序
return kafkaTemplate.send("order-events", order.getOrderNo(), order)
.whenComplete((result, ex) -> {
if (ex != null) {
log.error("发送失败", ex);
} else {
log.info("发送成功: partition={}, offset={}",
result.getRecordMetadata().partition(),
result.getRecordMetadata().offset());
}
});
}
}消费消息
@Component
public class OrderConsumer {
@KafkaListener(topics = "order-events", groupId = "order-group")
public void handle(
@Payload Order order,
@Header(KafkaHeaders.OFFSET) long offset,
Acknowledgment ack) {
try {
orderService.process(order);
ack.acknowledge();
} catch (Exception e) {
log.error("消费失败 offset={}", offset, e);
// 不 ack → 偏移量不提交,下次重启重新消费
}
}
}核心概念
| 概念 | 说明 | 关键参数 |
|---|---|---|
| Topic | 消息的逻辑分类 | @KafkaListener(topics = "name") |
| Partition | 物理分片,有序追加 | Topic 创建时指定 --partitions N |
| Offset | 分区内消息的唯一序号 | 消费者通过 Acknowledgment 提交 |
| Consumer Group | 组内消费者协同消费分区 | spring.kafka.consumer.group-id |
| ISR | 同步副本集 | acks=all + min.insync.replicas=2 |
重试与死信 Topic
@Configuration
public class KafkaConsumerConfig {
@Bean
public ConcurrentKafkaListenerContainerFactory<String, Object> dltFactory(
ConsumerFactory<String, Object> cf,
KafkaTemplate<String, Object> template) {
var factory = new ConcurrentKafkaListenerContainerFactory<String, Object>();
factory.setConsumerFactory(cf);
// 重试 3 次后发送到 "原Topic.DLT" 死信
var recoverer = new DeadLetterPublishingRecoverer(template,
(rec, ex) -> new TopicPartition(rec.topic() + ".DLT", rec.partition()));
factory.setCommonErrorHandler(
new DefaultErrorHandler(recoverer, new FixedBackOff(2000L, 3L)));
return factory;
}
}@Component
public class OrderDltConsumer {
@KafkaListener(topics = "order-events.DLT", groupId = "order-dlt")
public void handleDlt(@Payload Order order) {
log.warn("死信: orderNo={}", order.getOrderNo());
// 持久化到错误表,人工介入
}
}幂等消费
@KafkaListener(topics = "order-events")
public void handle(Order order, Acknowledgment ack) {
try {
// 利用业务唯一键去重(idempotent_key 在数据库有唯一约束)
eventService.processIfNotExists(order.getEventId(), () -> {
orderService.process(order);
});
ack.acknowledge();
} catch (DuplicateKeyException e) {
log.info("重复事件: {}", order.getEventId());
ack.acknowledge(); // 仍要提交偏移量,避免无限重放
}
}架构:ConcurrentMessageListenerContainer
@EnableKafka
└── KafkaListenerAnnotationBeanPostProcessor
└── 为每个 @KafkaListener 创建 ConcurrentMessageListenerContainer
└── 按 concurrency=N 创建 N 个 KafkaMessageListenerContainer
└── 每个拥有一个 KafkaConsumer 线程
├── 分配到的分区
├── poll() → 拉取一批消息
├── 遍历调用 @KafkaListener 方法
└── 提交偏移量| 分区数 | concurrency | 结果 |
|---|---|---|
| 3 | 3 | 每个消费者 1 个分区——最优 |
| 6 | 3 | 每个消费者 2 个分区——合理 |
| 3 | 6 | 3 个线程闲置——concurrency 不能超过分区数 |
生产意识:concurrency 不能大于分区数,否则闲置。auto-offset-reset=earliest 在新消费者组或分区再均衡时会重新消费历史数据。Topic 分区数只能增不能减,预先规划好预期吞吐量。
小结
- Kafka 是分布式日志:消息持久化到磁盘,消费者通过 Offset 控制读取位置。
- 生产者
acks=all+enable.idempotence=true保证投递可靠;消费者手动 ack 实现精确控制。 - 重试耗尽后进死信 Topic(
.DLT),避免消息无限重试。 - 消费端去重依赖业务唯一键——幂等生产者只保证生产端不重复。
下一节:RabbitMQ 集成 · RabbitMQ 与 Kafka 对比
上一节:RabbitMQ 集成
