09. Spring 集成消息队列

Spring 集成 RabbitMQ、Kafka、RocketMQ 的完整方案,涵盖事务消息、顺序消费、死信队列与消费幂等性设计

消息队列是微服务解耦与异步化的核心基础设施。Spring 提供了统一的抽象层,但每种消息中间件在消息语义、事务保证与消费模型上各有侧重。本文基于实战场景对比集成方案。

1. 三剑客选型对比

维度RabbitMQKafkaRocketMQ
协议AMQP自定义协议自定义协议
吞吐量万级百万级十万级
延迟微秒级毫秒级毫秒级
消费模型Push(推)Pull(拉)Push + Pull
消息顺序单队列有序Partition 内有序Queue 内有序
事务消息本地事务 + Confirm事务生产者(幂等)✅ 原生支持
延迟消息死信队列 + TTL不太适合✅ 原生支持 18 级
消息追踪插件较复杂较完善
管理界面完善需 Kowl/AKHQConsole 完善
最佳场景路由复杂、可靠投递日志/流处理/大数据金融级事务、定时任务

2. Spring AMQP + RabbitMQ

2.1 基础集成

spring:
  rabbitmq:
    host: localhost
    port: 5672
    username: guest
    password: guest
    listener:
      simple:
        concurrency: 5
        max-concurrency: 20
        acknowledge-mode: manual
        prefetch: 10
        retry:
          enabled: true
          max-attempts: 3
          initial-interval: 1000ms
@Configuration
public class RabbitConfig {
    
    @Bean
    public Queue orderQueue() {
        return QueueBuilder.durable("order.queue")
            .withArgument("x-dead-letter-exchange", "order.dlx")
            .withArgument("x-dead-letter-routing-key", "order.failed")
            .build();
    }
    
    @Bean
    public DirectExchange orderExchange() {
        return new DirectExchange("order.exchange");
    }
    
    @Bean
    public Binding orderBinding(Queue orderQueue, DirectExchange orderExchange) {
        return BindingBuilder.bind(orderQueue).to(orderExchange).with("order.created");
    }
}

@Component
public class OrderListener {
    @RabbitListener(queues = "order.queue")
    public void handleOrder(@Payload OrderMessage message, 
                           Channel channel,
                           @Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException {
        try {
            processOrder(message);
            channel.basicAck(tag, false);
        } catch (Exception e) {
            channel.basicNack(tag, false, !message.isRedelivered()); // 重试一次后入死信
        }
    }
}

2.2 延迟消息(死信 + TTL)

@Bean
public Queue delayQueue() {
    return QueueBuilder.durable("order.delay.queue")
        .withArgument("x-message-ttl", 60000)           // 60 秒后过期
        .withArgument("x-dead-letter-exchange", "order.exchange")
        .withArgument("x-dead-letter-routing-key", "order.timeout")
        .build();
}

// 业务:订单 30 分钟未支付自动取消
rabbitTemplate.convertAndSend("order.delay.exchange", "order.delay", order, msg -> {
    msg.getMessageProperties().setExpiration("1800000"); // 30 分钟 TTL
    return msg;
});

RocketMQ 原生延迟:18 个固定级别(1s/5s/10s/30s/1m/…/2h)。

3. Spring Kafka

3.1 生产者配置

spring:
  kafka:
    bootstrap-servers: kafka1:9092,kafka2:9092
    producer:
      acks: all                    # 0/1/all
      retries: 3
      batch-size: 16384
      linger-ms: 5
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
      properties:
        enable.idempotence: true    # 幂等生产者(防止重复发送)
        max.in.flight.requests.per.connection: 5
@Service
public class KafkaOrderProducer {
    @Autowired private KafkaTemplate<String, OrderEvent> kafkaTemplate;
    
    public void sendOrderCreated(OrderEvent event) {
        ProducerRecord<String, OrderEvent> record = new ProducerRecord<>(
            "order-events", 
            event.getOrderId(),    // 相同 key → 同一 partition → 顺序保证
            event
        );
        
        ListenableFuture<SendResult<String, OrderEvent>> future = 
            kafkaTemplate.send(record);
        
        future.addCallback(
            result -> log.info("Sent: {}", result.getRecordMetadata().offset()),
            failure -> log.error("Failed: {}", failure.getMessage())
        );
    }
}

3.2 消费者配置

spring:
  kafka:
    consumer:
      group-id: order-service
      auto-offset-reset: earliest
      enable-auto-commit: false
      max-poll-records: 50
      isolation-level: read_committed  // 事务消费者
    listener:
      ack-mode: manual_immediate
      concurrency: 3                    // 每个 topic 的并发消费者数
@Component
public class KafkaOrderConsumer {
    
    @KafkaListener(topics = "order-events", groupId = "order-service")
    public void consume(@Payload OrderEvent event,
                       @Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
                       @Header(KafkaHeaders.OFFSET) long offset,
                       Acknowledgment ack) {
        try {
            // 幂等校验:orderId + eventType 已处理过?
            if (idempotencyService.isProcessed(event.getId())) {
                ack.acknowledge();
                return;
            }
            
            processEvent(event);
            idempotencyService.markProcessed(event.getId());
            ack.acknowledge();
        } catch (Exception e) {
            // 不 ack,自动重试(配置 `retry`)或入死信 topic
            throw e;
        }
    }
}

3.3 Kafka 事务

@Service
public class OrderService {
    @Autowired private KafkaTemplate<String, Object> kafkaTemplate;
    
    @Transactional("kafkaTransactionManager")
    public void createOrderWithEvent(Order order) {
        // 1. 保存订单到 DB
        orderRepository.save(order);
        
        // 2. 发送事件(同一事务中)
        kafkaTemplate.send("order-events", new OrderCreatedEvent(order));
        // 若 DB 回滚,消息也不会发送
    }
}

4. RocketMQ Spring Starter

4.1 事务消息

RocketMQ 是原生支持事务消息最完善的方案。

@Service
public class OrderRocketService {
    @Autowired private RocketMQTemplate rocketMQTemplate;
    
    public void createOrder(Order order) {
        Message<Order> message = MessageBuilder.withPayload(order)
            .setHeader("KEYS", order.getId())
            .build();
        
        TransactionSendResult result = rocketMQTemplate.sendMessageInTransaction(
            "order-producer-group",
            "order-topic",
            message,
            order  // 本地事务参数
        );
        
        if (!result.getLocalTransactionState().equals(LocalTransactionState.COMMIT_MESSAGE)) {
            throw new OrderException("订单创建失败");
        }
    }
}

// 本地事务监听器
@RocketMQTransactionListener
public class OrderTransactionListener implements RocketMQLocalTransactionListener {
    @Override
    public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        try {
            Order order = (Order) arg;
            orderService.save(order);  // 执行本地事务
            return RocketMQLocalTransactionState.COMMIT_MESSAGE;
        } catch (Exception e) {
            return RocketMQLocalTransactionState.ROLLBACK_MESSAGE;
        }
    }
    
    @Override
    public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
        // 事务回查:网络等原因导致 RocketMQ 未收到确认
        String orderId = msg.getKeys();
        if (orderService.exists(orderId)) {
            return RocketMQLocalTransactionState.COMMIT_MESSAGE;
        }
        return RocketMQLocalTransactionState.ROLLBACK_MESSAGE;
    }
}

5. 消费幂等性设计

5.1 幂等实现方案

消息唯一标识:msgId(中间件生成)或业务键(bizId)

方案 1:数据库唯一索引
  CREATE TABLE idempotency (
    idempotency_key VARCHAR(64) PRIMARY KEY,
    created_at TIMESTAMP DEFAULT NOW()
  );
  INSERT IGNORE INTO idempotency (idempotency_key) VALUES ('order:123:created');

方案 2:Redis SETNX(带 TTL)
  SET idempotency:order:123:created "1" NX EX 86400

方案 3:业务状态机防重
  订单状态: PENDING → PAID → SHIPPED
  PAID 状态下重复收到 PAYMENT_SUCCESS 消息 → 直接忽略

6. 死信队列与监控

// Spring Kafka 死信 topic
@KafkaListener(topics = "order-events", groupId = "order-service")
public void consumeWithDLQ(OrderEvent event) {
    // 重试 3 次后自动转发到 order-events.DLT
}

// DLT 监听(人工/自动化处理)
@KafkaListener(topics = "order-events.DLT", groupId = "order-dlt-service")
public void handleDeadLetter(OrderEvent event, 
                              @Header(KafkaHeaders.DLT_EXCEPTION_MESSAGE) String error) {
    alertService.send("DLQ Alert", event, error);
    // 人工介入或自动补偿逻辑
}

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「java-enterprise」更多文章

  1. 限流算法深度解析:令牌桶、漏桶与滑动窗口计数
  2. Java 代码质量:SonarQube、Checkstyle 与 SpotBugs 工程化实践
  3. Spring IoC 容器与依赖注入原理深度剖析