POM依赖

<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>
  • application.yaml文件添加
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";
    }
}

生产环境配置

  • 生产者yaml文件
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);
                // 可触发自定义告警或转存到死信表
            }
        });
    }
}
  • 消费者yaml文件
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; // 抛出异常触发重试
        }
    }
}