这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 Spring Kafka 是 Java 生态最常用的 Kafka 封装。本章从发送、消费、手动确认、错误处理讲到死信队列与批量监听,给出一套可以直接落地的工程模板。
9.1 依赖与配置
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>
Spring Boot 3.x 的依赖管理会自动匹配合适的 spring-kafka 版本。
application.yml:
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
linger.ms: 10
compression.type: zstd
consumer:
group-id: order-service
auto-offset-reset: earliest
enable-auto-commit: false
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
properties:
spring.json.trusted.packages: "com.example.dto.*"
isolation.level: read_committed
listener:
ack-mode: manual
concurrency: 3
要点:
concurrency是每个监听器容器的线程数,最大有效值 = 分区数;ack-mode: manual配合代码里Acknowledgment.acknowledge();- JSON 反序列化必须配置可信包,否则会拒绝反序列化未知类型。
9.2 发送消息:KafkaTemplate
@Service
@RequiredArgsConstructor
public class OrderEventPublisher {
private final KafkaTemplate<String, OrderCreated> kafkaTemplate;
public void publish(OrderCreated event) {
kafkaTemplate.send("order-events", event.orderId(), event)
.whenComplete((result, ex) -> {
if (ex != null) {
log.error("send failed, orderId={}", event.orderId(), ex);
// 落补偿表或触发重试
} else {
var meta = result.getRecordMetadata();
log.debug("sent partition={} offset={}",
meta.partition(), meta.offset());
}
});
}
}
同步发送(需要强确认时):
kafkaTemplate.send("order-events", key, event).get(3, TimeUnit.SECONDS);
9.3 消费消息:@KafkaListener
@Component
@RequiredArgsConstructor
@Slf4j
public class OrderEventListener {
private final OrderService orderService;
@KafkaListener(
topics = "order-events",
groupId = "order-service",
concurrency = "3"
)
public void onMessage(ConsumerRecord<String, OrderCreated> record,
Acknowledgment ack) {
try {
orderService.handle(record.value());
ack.acknowledge(); // 处理成功才提交
} catch (Exception e) {
log.error("handle failed, offset={}", record.offset(), e);
throw e; // 交给错误处理器决定重试或进 DLT
}
}
}
9.4 错误处理与死信队列
直接抛异常会触发容器重试,默认行为可能导致无限循环。生产推荐“有限重试 + 死信”:
@Configuration
public class KafkaErrorConfig {
@Bean
public DefaultErrorHandler errorHandler(KafkaTemplate<Object, Object> template) {
// 重试 3 次,间隔 1s,指数退避倍数 2
var recoverer = new DeadLetterPublishingRecoverer(
template,
(record, ex) -> {
// 死信目标:原 topic + ".DLT",保持原分区
return new TopicPartition(record.topic() + ".DLT", record.partition());
});
var backOff = new ExponentialBackOffWithMaxRetries(3);
backOff.setInitialInterval(1000L);
backOff.setMultiplier(2.0);
var handler = new DefaultErrorHandler(recoverer, backOff);
// 这些异常不重试,直接进死信
handler.addNotRetryableExceptions(
DeserializationException.class,
IllegalArgumentException.class);
return handler;
}
}
配套监听死信:
@KafkaListener(topics = "order-events.DLT", groupId = "order-dlt-watcher")
public void onDead(ConsumerRecord<String, OrderCreated> record) {
log.error("dead letter key={} offset={} value={}",
record.key(), record.offset(), record.value());
// 人工处理、告警、回补
}
死信设计建议:
- 死信消息带上原始 topic、partition、offset、异常堆栈头,方便追溯;
- 死信 topic 的保留期、副本策略要按业务重要性单独设置;
- 建死信看板与告警,否则问题只会在用户投诉时才暴露。
9.5 批量消费
下游是数据库或搜索引擎时,批量写入能带来数量级的性能提升:
@KafkaListener(topics = "order-events", batch = "true")
public void onBatch(List<ConsumerRecord<String, OrderCreated>> records,
Acknowledgment ack) {
List<OrderCreated> events = records.stream()
.map(ConsumerRecord::value).toList();
orderService.batchHandle(events);
ack.acknowledge();
}
spring.kafka.consumer:
max-poll-records: 200
fetch-min-size: 1KB
fetch-max-wait: 500ms
9.6 消费者多线程
concurrency=3 相当于起 3 个 KafkaMessageListenerContainer,各自拥有独立 consumer,各自分到一部分分区。这是首选的并发模型,因为它保留了标准的位移与再平衡语义。
只有当单条消息处理极重(比如调用多个慢接口)且分区数已到上限时,才考虑在监听器内部再套线程池。这时要自己解决:
- 位移提交粒度(一批任务全完成后再 ack);
- 优雅停机(关闭线程池、等待任务完成);
- 异常时的消息路由(重试/死信)。
9.7 事务生产者
需要“多条消息原子写入”或 consume-transform-produce 场景时:
kafkaTemplate.executeInTransaction(ops -> {
ops.send("order-events", key, orderCreated);
ops.send("audit-events", key, auditEvent);
return true;
});
配置:
spring:
kafka:
producer:
transaction-id-prefix: order-tx-
消费端记得 isolation.level: read_committed,否则会读到未提交数据(第 14 章展开)。
9.8 一个完整的消费端模板
@Component
@Slf4j
public class ReliableOrderListener {
private final OrderService orderService;
private final MeterRegistry metrics;
public ReliableOrderListener(OrderService orderService, MeterRegistry metrics) {
this.orderService = orderService;
this.metrics = metrics;
}
@KafkaListener(topics = "order-events", groupId = "order-service")
public void listen(ConsumerRecord<String, OrderCreated> record, Acknowledgment ack) {
long start = System.currentTimeMillis();
try {
orderService.handleIdempotent(record.value());
ack.acknowledge();
metrics.counter("kafka.consume.success").increment();
} catch (Exception e) {
metrics.counter("kafka.consume.failure").increment();
throw e;
} finally {
metrics.timer("kafka.consume.latency")
.record(System.currentTimeMillis() - start, TimeUnit.MILLISECONDS);
}
}
}
这套结构包含:手动 ack、幂等处理、指标埋点、异常上抛交给统一错误处理。可以直接作为业务消费者模板复制。
本章小结
concurrency创建多个真实 consumer,是 Spring Kafka 并发消费的标准方式;ack-mode: manual+ 处理后确认,实现至少一次语义;DefaultErrorHandler+DeadLetterPublishingRecoverer构成“重试有限次、失败进 DLT”的闭环;- 批量监听对数据库/搜索类下游收益巨大;
- 事务用
executeInTransaction,消费端配合read_committed。
思考题
concurrency=6但 topic 只有 3 个分区,会发生什么?- 监听器里
ack.acknowledge()之后进程立即崩溃,这条消息会怎样?不 ack 又会怎样? - 如何设计死信消息的元信息,才能让“重新投递”变得可操作?