<dependencyManagement>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-dependencies</artifactId>
<version>2.7.18</version>
<type>pom</type>
<scope>import</scope>
</dependency>
</dependencies>
</dependencyManagement>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
<version>2.9.11</version>
</dependency>
</dependencies>
spring:
kafka:
bootstrap-servers: localhost:9092
@Component
public class KafkaProducer {
@Resource
private KafkaTemplate<String, String> kafkaTemplate;
public void sendMessage(String topic, String key, String value) {
// 创建消息
ProducerRecord<String, String> record = new ProducerRecord<>(topic, key, value);
// 发送消息
kafkaTemplate.send(record);
}
}
@Component
public class KafkaConsumer {
@KafkaListener(topics = "test", groupId = "test-group")
public void receiveMessage(String message) {
System.out.println("接收到消息:" + message);
}
}
@RestController
public class TestController {
@Resource
private KafkaProducer kafkaProducer;
@GetMapping("/test")
public String test() {
kafkaProducer.sendMessage("test", "test", "test");
return "hello world";
}
}
spring:
kafka:
bootstrap-servers: ${KAFKA_BOOTSTRAP_SERVERS:localhost:9092}
producer:
# 确认机制:all 表示所有副本都写入才确认,最强可靠性
acks: all
# 发送失败重试次数,建议大于 0,同时配合 retry.backoff.ms 控制重试间隔
retries: 10
# 批量大小,适当调大可提升吞吐(默认 16384)
batch-size: 16384
# 发送延迟,为了凑足 batch-size 或时间到达即发送
linger-ms: 5
# 缓冲区内存大小
buffer-memory: 33554432
# 生产者客户端 ID,用于监控
client-id: ${spring.application.name}-producer
# 键值序列化器
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
# 可选:启用事务(如果需要 exactly-once 语义)
transaction-id-prefix: tx-${spring.application.name}-
properties:
# 幂等性(防止重试导致的重复),必须开启(>=0.11)
enable.idempotence: true
# 最大请求大小(默认 1MB,如果消息体大需要调大)
max.request.size: 10485760
@Component
public class KafkaProducer {
@Resource
private KafkaTemplate<String, Object> kafkaTemplate;
public void sendMessage(String topic, String key, Object value) {
ProducerRecord<String, Object> record = new ProducerRecord<>(topic, key, value);
ListenableFuture<SendResult<String, Object>> future = kafkaTemplate.send(record);
future.addCallback(new ListenableFutureCallback<SendResult<String, Object>>() {
@Override
public void onSuccess(SendResult<String, Object> result) {
// 记录成功日志或 Metrics
log.info("消息发送成功 topic={} partition={} offset={}",
result.getRecordMetadata().topic(),
result.getRecordMetadata().partition(),
result.getRecordMetadata().offset());
}
@Override
public void onFailure(Throwable ex) {
// 记录失败并告警,必要时进行补偿(如持久化到本地后重试)
log.error("消息发送失败 topic={} key={}", topic, key, ex);
// 可触发自定义告警或转存到死信表
}
});
}
}
spring:
kafka:
consumer:
# 消费者组 ID,建议按业务区分
group-id: ${spring.application.name}-group
# 自动提交关闭,改为手动提交(保证至少一次)
enable-auto-commit: false
# 从何处开始消费:latest / earliest / none
auto-offset-reset: latest
# 键值反序列化器
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
# 允许反序列化不信任的包(使用 JSON 时必须)
properties:
spring.json.trusted.packages: "*"
# 每次拉取的最大记录数
max-poll-records: 50
# 客户端 ID
client-id: ${spring.application.name}-consumer
listener:
# 手动提交模式:手动调用 Acknowledgment.acknowledge()
ack-mode: manual
# 并发消费者数量(对应分区数)
concurrency: 3
# 当消费者出现异常时的行为:默认抛出异常停止,可配置重试
properties:
# 一次 poll 的最大等待时间(毫秒)
fetch.max.wait.ms: 500
# 每次 fetch 的最小数据量(字节)
fetch.min.bytes: 1024
@Component
@Slf4j
public class KafkaConsumer {
@KafkaListener(topics = "test", groupId = "test-group")
public void receiveMessage(String message, Acknowledgment ack) {
try {
// 业务处理
log.info("接收到消息:{}", message);
// 处理成功,手动提交偏移量
ack.acknowledge();
} catch (Exception e) {
log.error("处理消息失败,消息内容:{}", message, e);
// 不提交偏移量,消息会重试
// 为防止无限重试,可配置重试次数 + 死信队列
throw e; // 抛出异常触发重试
}
}
}