跳到正文
Joeplover
后端开发·2025-10-14·约 3 分钟阅读

Spring Boot + Kafka 实战:生产者消费者完整示例

Spring Boot 集成 Spring Kafka 的完整示例,包含 HTTP 触发消息发送、KafkaTemplate 批量生产和 @KafkaListener 消费。

实时数据流

Spring Boot + Kafka 实战:生产与消费

Spring Kafka 集成

Spring Kafka 在原生 Kafka Client API 上做了封装,让生产者和消费者的使用体验大幅简化。

KafkaTemplate 生产者

@Service
public class OrderProducer {
    
    @Autowired
    private KafkaTemplate<String, OrderMessage> kafkaTemplate;
    
    public void sendOrderCreated(Order order) {
        OrderMessage msg = OrderMessage.builder()
            .orderId(order.getId())
            .userId(order.getUserId())
            .amount(order.getTotalAmount())
            .timestamp(System.currentTimeMillis())
            .build();
        
        // 发送消息
        ListenableFuture<SendResult<String, OrderMessage>> future = 
            kafkaTemplate.send("order-events", msg.getOrderId(), msg);
        
        // 异步处理结果
        future.addCallback(
            result -> log.info("订单消息发送成功: {}", msg.getOrderId()),
            ex -> log.error("订单消息发送失败", ex)
        );
    }
}

@KafkaListener 消费者

@Component
public class OrderConsumer {
    
    @KafkaListener(
        topics = "order-events",
        groupId = "order-processing-group",
        containerFactory = "batchFactory"
    )
    public void consume(List<ConsumerRecord<String, OrderMessage>> records, 
                        Acknowledgment ack) {
        try {
            for (ConsumerRecord<String, OrderMessage> record : records) {
                processOrder(record.value());
            }
            // 处理完成后手动提交 offset
            ack.acknowledge();
        } catch (Exception e) {
            log.error("消费订单消息失败", e);
            // 不提交 offset,下次重新消费
        }
    }
}

序列化与反序列化

Kafka 传输的是字节数组,Java 对象需要序列化:

public class OrderMessageSerializer implements Serializer<OrderMessage> {
    @Override
    public byte[] serialize(String topic, OrderMessage data) {
        return JsonUtils.toJson(data).getBytes(StandardCharsets.UTF_8);
    }
}

public class OrderMessageDeserializer implements Deserializer<OrderMessage> {
    @Override
    public OrderMessage deserialize(String topic, byte[] data) {
        return JsonUtils.fromJson(new String(data, StandardCharsets.UTF_8), 
                                   OrderMessage.class);
    }
}

对应的生产者配置:

spring:
  kafka:
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: com.example.config.OrderMessageSerializer
    consumer:
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: com.example.config.OrderMessageDeserializer
      properties:
        spring.json.trusted.packages: "*"

幂等消费:唯一业务 ID

Kafka 是 at-least-once 语义,意味着消费者可能收到重复消息。消费端需要幂等:

@Service
public class OrderProcessor {
    
    @Autowired
    private RedisTemplate<String, String> redisTemplate;
    
    public void processOrder(OrderMessage msg) {
        // 先检查是否处理过
        String processedKey = "processed:order:" + msg.getOrderId();
        Boolean existed = redisTemplate.opsForValue().setIfAbsent(
            processedKey, "1", Duration.ofHours(24));
        
        if (Boolean.FALSE.equals(existed)) {
            log.info("订单 {} 已处理过,跳过", msg.getOrderId());
            return;  // 幂等
        }
        
        // 核心业务逻辑
        orderService.processOrder(msg);
    }
}

异常处理与死信队列

@Bean
public ConcurrentKafkaListenerContainerFactory<String, OrderMessage> factory() {
    ConcurrentKafkaListenerContainerFactory<String, OrderMessage> factory = 
        new ConcurrentKafkaListenerContainerFactory<>();
    
    // 重试配置:重试 3 次,间隔 2 秒
    factory.setRetryTemplate(new RetryTemplate() {{
        setRetryOperations(new SimpleRetryPolicy(3));
        setBackOffPolicy(new FixedBackOffPolicy() {{
            setBackOffPeriod(2000);
        }});
    }});
    
    // 重试耗尽后进入死信队列
    factory.setErrorHandler((record, exception) -> {
        kafkaTemplate.send("order-events-dlq", record.value());
        log.error("消息处理失败,已转入死信队列", exception);
    });
    
    return factory;
}

消费者组与 offset

消费者组内每个消费者消费不同的 Partition。offset 提交方式:

  • 自动提交(enable.auto.commit=true):定时提交,可能导致重复消费或丢消息
  • 手动提交(enable.auto.commit=false):业务处理完成后手动 ack。推荐生产使用
spring:
  kafka:
    consumer:
      enable-auto-commit: false
      auto-offset-reset: earliest  # 新消费者从最早的消息开始

版本兼容

Kafka 服务端版本和客户端版本不匹配可能导致连接失败。Spring Kafka 自动适配 Kafka 服务端版本,但如果用原生 Kafka Client API,需要确保客户端版本 ≤ 服务端版本。

消息顺序

Kafka 的消息顺序只在单 Partition 内保证。如果业务要求严格顺序(如订单状态流转),需要确保同一状态的订单发到同一个 Partition(用订单 ID 作为 Key)。