最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

SpringBoot整合Kafka實(shí)現(xiàn)高可用消息隊(duì)列集群詳解

 更新時(shí)間:2026年01月08日 14:29:38   作者:悟空碼字  
Apache?Kafka是一個(gè)分布式流處理平臺(tái),這篇文章主要介紹了SpringBoot如何整合Kafka實(shí)現(xiàn)高可用消息隊(duì)列集群,文中的示例代碼講解詳細(xì),感興趣的小伙伴可以了解下

一、Kafka簡(jiǎn)介

1.1 什么是Kafka

Apache Kafka是一個(gè)分布式流處理平臺(tái),具有以下核心特性:

  • 高吞吐量:支持每秒百萬(wàn)級(jí)消息處理
  • 可擴(kuò)展性:支持水平擴(kuò)展,可動(dòng)態(tài)添加節(jié)點(diǎn)
  • 持久化存儲(chǔ):消息可持久化到磁盤,支持?jǐn)?shù)據(jù)保留策略
  • 高可用性:通過(guò)副本機(jī)制保證數(shù)據(jù)不丟失
  • 分布式架構(gòu):支持多生產(chǎn)者和消費(fèi)者

1.2 Kafka核心概念

  • Broker:Kafka集群中的單個(gè)節(jié)點(diǎn)
  • Topic:消息的分類主題
  • Partition:Topic的分區(qū),實(shí)現(xiàn)并行處理
  • Replica:分區(qū)副本,保證高可用
  • Producer:消息生產(chǎn)者
  • Consumer:消息消費(fèi)者
  • Consumer Group:消費(fèi)者組

二、搭建Kafka高可用集群

集群架構(gòu)規(guī)劃

建議至少3個(gè)節(jié)點(diǎn)的Kafka集群 + 3個(gè)節(jié)點(diǎn)的Zookeeper集群:

  • Zookeeper集群:zk1:2181, zk2:2181, zk3:2181
  • Kafka集群:kafka1:9092, kafka2:9092, kafka3:9092

三、SpringBoot整合Kafka詳細(xì)步驟

3.1 創(chuàng)建SpringBoot項(xiàng)目

使用Spring Initializr創(chuàng)建項(xiàng)目,添加依賴:

<!-- pom.xml -->
<dependencies>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-web</artifactId>
    </dependency>
    
    <dependency>
        <groupId>org.springframework.kafka</groupId>
        <artifactId>spring-kafka</artifactId>
    </dependency>
    
    <dependency>
        <groupId>org.projectlombok</groupId>
        <artifactId>lombok</artifactId>
        <optional>true</optional>
    </dependency>
    
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-validation</artifactId>
    </dependency>
</dependencies>

3.2 配置文件

# application.yml
spring:
  kafka:
    # Kafka集群配置(高可用)
    bootstrap-servers: kafka1:9092,kafka2:9092,kafka3:9092
    
    # 生產(chǎn)者配置
    producer:
      retries: 3  # 發(fā)送失敗重試次數(shù)
      acks: all   # 所有副本確認(rèn)才認(rèn)為發(fā)送成功
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
      properties:
        compression.type: snappy  # 壓縮類型
        linger.ms: 5  # 等待時(shí)間,批量發(fā)送提高吞吐量
    
    # 消費(fèi)者配置
    consumer:
      group-id: ${spring.application.name}-group
      auto-offset-reset: earliest
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
      properties:
        spring.json.trusted.packages: "com.example.kafka.dto"
        max.poll.records: 500  # 一次拉取最大記錄數(shù)
        session.timeout.ms: 10000  # 會(huì)話超時(shí)時(shí)間
        heartbeat.interval.ms: 3000  # 心跳間隔
    
    # 監(jiān)聽(tīng)器配置
    listener:
      concurrency: 3  # 并發(fā)消費(fèi)者數(shù)量
      ack-mode: batch  # 批量確認(rèn)
      missing-topics-fatal: false  # 主題不存在時(shí)不報(bào)錯(cuò)
      
    # 高可用配置
    properties:
      # 分區(qū)副本配置
      replication.factor: 3
      min.insync.replicas: 2
      # 生產(chǎn)者的高可用配置
      enable.idempotence: true  # 冪等性
      max.in.flight.requests.per.connection: 5

# 自定義配置
kafka:
  topics:
    order-topic: order-topic
    payment-topic: payment-topic
    retry-topic: retry-topic
  retry:
    max-attempts: 3
    backoff-interval: 1000

3.3 配置類

// KafkaConfig.java
@Configuration
@EnableKafka
@Slf4j
public class KafkaConfig {
    
    @Value("${kafka.topics.order-topic}")
    private String orderTopic;
    
    @Value("${kafka.topics.payment-topic}")
    private String paymentTopic;
    
    @Value("${kafka.topics.retry-topic}")
    private String retryTopic;
    
    @Bean
    public KafkaAdmin kafkaAdmin() {
        Map<String, Object> configs = new HashMap<>();
        configs.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, 
                   "kafka1:9092,kafka2:9092,kafka3:9092");
        return new KafkaAdmin(configs);
    }
    
    @Bean
    public NewTopic orderTopic() {
        // 創(chuàng)建Topic:3個(gè)分區(qū),3個(gè)副本
        return new NewTopic(orderTopic, 3, (short) 3);
    }
    
    @Bean
    public NewTopic paymentTopic() {
        return new NewTopic(paymentTopic, 2, (short) 3);
    }
    
    @Bean
    public NewTopic retryTopic() {
        return new NewTopic(retryTopic, 1, (short) 3);
    }
    
    // 死信隊(duì)列配置
    @Bean
    public DeadLetterPublishingRecoverer dlqRecoverer(KafkaTemplate<String, Object> template) {
        return new DeadLetterPublishingRecoverer(template, 
            (record, ex) -> {
                log.error("消息處理失敗,發(fā)送到死信隊(duì)列: {}", record.value(), ex);
                return new TopicPartition("dlq-topic", record.partition());
            });
    }
    
    @Bean
    public DefaultErrorHandler errorHandler(DeadLetterPublishingRecoverer dlqRecoverer) {
        // 重試3次后進(jìn)入死信隊(duì)列
        DefaultErrorHandler handler = new DefaultErrorHandler(dlqRecoverer, 
            new FixedBackOff(1000L, 3));
        handler.addNotRetryableExceptions(IllegalArgumentException.class);
        return handler;
    }
    
    // 生產(chǎn)者工廠增強(qiáng)配置
    @Bean
    public ProducerFactory<String, Object> producerFactory() {
        Map<String, Object> configProps = new HashMap<>();
        configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, 
                       "kafka1:9092,kafka2:9092,kafka3:9092");
        configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, 
                       StringSerializer.class);
        configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, 
                       JsonSerializer.class);
        configProps.put(ProducerConfig.ACKS_CONFIG, "all");  // 所有副本確認(rèn)
        configProps.put(ProducerConfig.RETRIES_CONFIG, 3);    // 重試次數(shù)
        configProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);  // 冪等性
        configProps.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5);
        return new DefaultKafkaProducerFactory<>(configProps);
    }
}

3.4 消息實(shí)體類

// OrderMessage.java
@Data
@NoArgsConstructor
@AllArgsConstructor
@Builder
public class OrderMessage implements Serializable {
    private String orderId;
    private String userId;
    private BigDecimal amount;
    private String productName;
    private Integer quantity;
    private LocalDateTime createTime;
    private MessageStatus status;
    
    public enum MessageStatus {
        PENDING, PROCESSING, SUCCESS, FAILED
    }
}

// PaymentMessage.java
@Data
@NoArgsConstructor
@AllArgsConstructor
@Builder
public class PaymentMessage {
    private String paymentId;
    private String orderId;
    private BigDecimal amount;
    private PaymentMethod paymentMethod;
    private PaymentStatus status;
    private LocalDateTime paymentTime;
    
    public enum PaymentMethod {
        ALIPAY, WECHAT, CREDIT_CARD
    }
    
    public enum PaymentStatus {
        INIT, PROCESSING, SUCCESS, FAILED
    }
}

3.5 生產(chǎn)者服務(wù)

// KafkaProducerService.java
@Service
@Slf4j
public class KafkaProducerService {
    
    @Autowired
    private KafkaTemplate<String, Object> kafkaTemplate;
    
    @Value("${kafka.topics.order-topic}")
    private String orderTopic;
    
    @Value("${kafka.topics.payment-topic}")
    private String paymentTopic;
    
    /**
     * 發(fā)送訂單消息(同步)
     */
    public SendResult<String, Object> sendOrderSync(OrderMessage orderMessage) {
        try {
            // 設(shè)置消息頭
            MessageHeaders headers = new MessageHeaders(Map.of(
                "message-id", UUID.randomUUID().toString(),
                "message-time", String.valueOf(System.currentTimeMillis())
            ));
            
            Message<OrderMessage> message = MessageBuilder
                .withPayload(orderMessage)
                .copyHeaders(headers)
                .build();
            
            // 同步發(fā)送,等待確認(rèn)
            ListenableFuture<SendResult<String, Object>> future = 
                kafkaTemplate.send(orderTopic, orderMessage.getOrderId(), message);
            
            // 等待發(fā)送結(jié)果
            SendResult<String, Object> result = future.get(5, TimeUnit.SECONDS);
            log.info("訂單消息發(fā)送成功: topic={}, partition={}, offset={}", 
                    result.getRecordMetadata().topic(),
                    result.getRecordMetadata().partition(),
                    result.getRecordMetadata().offset());
            return result;
            
        } catch (Exception e) {
            log.error("訂單消息發(fā)送失敗: {}", orderMessage, e);
            throw new RuntimeException("消息發(fā)送失敗", e);
        }
    }
    
    /**
     * 發(fā)送訂單消息(異步)
     */
    public void sendOrderAsync(OrderMessage orderMessage) {
        kafkaTemplate.send(orderTopic, orderMessage.getOrderId(), orderMessage)
            .addCallback(new ListenableFutureCallback<SendResult<String, Object>>() {
                @Override
                public void onSuccess(SendResult<String, Object> result) {
                    log.info("異步發(fā)送成功: topic={}, offset={}", 
                            result.getRecordMetadata().topic(),
                            result.getRecordMetadata().offset());
                }
                
                @Override
                public void onFailure(Throwable ex) {
                    log.error("異步發(fā)送失敗: {}", orderMessage, ex);
                    // 可以添加重試邏輯或?qū)懭氡镜匚募?
                }
            });
    }
    
    /**
     * 批量發(fā)送消息
     */
    public void batchSendOrders(List<OrderMessage> orderMessages) {
        orderMessages.forEach(message -> {
            kafkaTemplate.send(orderTopic, message.getOrderId(), message);
        });
        kafkaTemplate.flush(); // 確保所有消息都發(fā)送
    }
    
    /**
     * 發(fā)送到指定分區(qū)
     */
    public void sendToPartition(OrderMessage orderMessage, int partition) {
        kafkaTemplate.send(orderTopic, partition, 
                          orderMessage.getOrderId(), orderMessage);
    }
    
    /**
     * 事務(wù)消息發(fā)送
     */
    @Transactional(transactionManager = "kafkaTransactionManager")
    public void sendTransactionalMessage(OrderMessage orderMessage) {
        // 數(shù)據(jù)庫(kù)操作
        // orderRepository.save(order);
        
        // Kafka消息發(fā)送(與數(shù)據(jù)庫(kù)操作在同一個(gè)事務(wù)中)
        kafkaTemplate.send(orderTopic, orderMessage.getOrderId(), orderMessage);
        
        // 其他業(yè)務(wù)操作
    }
}

3.6 消費(fèi)者服務(wù)

// KafkaConsumerService.java
@Service
@Slf4j
public class KafkaConsumerService {
    
    private static final String ORDER_CONTAINER_FACTORY = "orderContainerFactory";
    private static final String PAYMENT_CONTAINER_FACTORY = "paymentContainerFactory";
    
    /**
     * 訂單消息消費(fèi)者 - 批量消費(fèi)
     */
    @KafkaListener(
        topics = "${kafka.topics.order-topic}",
        containerFactory = ORDER_CONTAINER_FACTORY,
        groupId = "order-consumer-group"
    )
    public void consumeOrderMessages(List<OrderMessage> messages) {
        log.info("收到批量訂單消息,數(shù)量: {}", messages.size());
        
        for (OrderMessage message : messages) {
            try {
                processOrderMessage(message);
            } catch (Exception e) {
                log.error("訂單處理失敗: {}", message.getOrderId(), e);
                // 記錄失敗消息,可以發(fā)送到重試隊(duì)列
            }
        }
    }
    
    /**
     * 單個(gè)訂單消息消費(fèi)
     */
    @KafkaListener(
        topics = "${kafka.topics.order-topic}",
        groupId = "order-single-consumer-group"
    )
    public void consumeSingleOrderMessage(
            @Payload OrderMessage message,
            @Header(KafkaHeaders.RECEIVED_TOPIC) String topic,
            @Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
            @Header(KafkaHeaders.OFFSET) long offset) {
        
        log.info("收到單個(gè)訂單消息: topic={}, partition={}, offset={}, orderId={}", 
                topic, partition, offset, message.getOrderId());
        
        try {
            // 業(yè)務(wù)處理邏輯
            processOrderMessage(message);
            
            // 處理成功后,可以發(fā)送確認(rèn)消息到下游
            sendPaymentMessage(message);
            
        } catch (Exception e) {
            log.error("訂單處理失敗: {}", message.getOrderId(), e);
            throw e; // 拋出異常會(huì)觸發(fā)重試機(jī)制
        }
    }
    
    /**
     * 支付消息消費(fèi)者
     */
    @KafkaListener(
        topics = "${kafka.topics.payment-topic}",
        containerFactory = PAYMENT_CONTAINER_FACTORY,
        groupId = "payment-consumer-group"
    )
    public void consumePaymentMessage(PaymentMessage message) {
        log.info("收到支付消息: {}", message.getPaymentId());
        
        // 支付處理邏輯
        try {
            processPayment(message);
        } catch (Exception e) {
            log.error("支付處理失敗: {}", message.getPaymentId(), e);
        }
    }
    
    private void processOrderMessage(OrderMessage message) {
        // 模擬業(yè)務(wù)處理
        log.info("處理訂單: {},金額: {}", message.getOrderId(), message.getAmount());
        
        // 業(yè)務(wù)邏輯,如:
        // 1. 驗(yàn)證訂單
        // 2. 扣減庫(kù)存
        // 3. 記錄日志
        // 4. 更新訂單狀態(tài)
        
        // 模擬處理時(shí)間
        try {
            Thread.sleep(100);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
    
    private void sendPaymentMessage(OrderMessage orderMessage) {
        PaymentMessage paymentMessage = PaymentMessage.builder()
            .paymentId(UUID.randomUUID().toString())
            .orderId(orderMessage.getOrderId())
            .amount(orderMessage.getAmount())
            .paymentMethod(PaymentMessage.PaymentMethod.ALIPAY)
            .status(PaymentMessage.PaymentStatus.INIT)
            .paymentTime(LocalDateTime.now())
            .build();
        
        // 這里可以使用KafkaTemplate發(fā)送支付消息
    }
    
    private void processPayment(PaymentMessage message) {
        // 支付處理邏輯
        log.info("處理支付: {},訂單: {}", message.getPaymentId(), message.getOrderId());
    }
}

3.7 消費(fèi)者容器工廠配置

// ConsumerConfig.java
@Configuration
public class ConsumerConfig {
    
    @Value("${spring.kafka.bootstrap-servers}")
    private String bootstrapServers;
    
    // 訂單消費(fèi)者容器工廠(批量消費(fèi))
    @Bean(ORDER_CONTAINER_FACTORY)
    public ConcurrentKafkaListenerContainerFactory<String, OrderMessage> 
            orderContainerFactory() {
        
        ConcurrentKafkaListenerContainerFactory<String, OrderMessage> factory =
            new ConcurrentKafkaListenerContainerFactory<>();
        
        factory.setConsumerFactory(orderConsumerFactory());
        factory.setConcurrency(3); // 并發(fā)消費(fèi)者數(shù)量
        factory.getContainerProperties().setPollTimeout(3000);
        factory.setBatchListener(true); // 啟用批量消費(fèi)
        factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.BATCH);
        
        // 設(shè)置批量消費(fèi)參數(shù)
        factory.getContainerProperties().setIdleBetweenPolls(1000);
        
        return factory;
    }
    
    @Bean
    public ConsumerFactory<String, OrderMessage> orderConsumerFactory() {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "order-consumer-group");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
        props.put(JsonDeserializer.TRUSTED_PACKAGES, "com.example.kafka.dto");
        props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 500); // 批量拉取數(shù)量
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
        
        // 高可用配置
        props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 10000);
        props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 3000);
        
        return new DefaultKafkaConsumerFactory<>(props);
    }
}

3.8 監(jiān)控和管理端點(diǎn)

// KafkaMonitorController.java
@RestController
@RequestMapping("/api/kafka")
@Slf4j
public class KafkaMonitorController {
    
    @Autowired
    private KafkaAdmin kafkaAdmin;
    
    @Autowired
    private KafkaTemplate<String, Object> kafkaTemplate;
    
    /**
     * 獲取Topic列表
     */
    @GetMapping("/topics")
    public ResponseEntity<List<String>> getTopics() throws Exception {
        try (AdminClient adminClient = AdminClient.create(kafkaAdmin.getConfigurationProperties())) {
            ListTopicsResult topicsResult = adminClient.listTopics();
            Set<String> topicNames = topicsResult.names().get();
            return ResponseEntity.ok(new ArrayList<>(topicNames));
        }
    }
    
    /**
     * 獲取Topic詳情
     */
    @GetMapping("/topics/{topic}/details")
    public ResponseEntity<Map<Integer, List<Integer>>> getTopicDetails(
            @PathVariable String topic) throws Exception {
        
        try (AdminClient adminClient = AdminClient.create(kafkaAdmin.getConfigurationProperties())) {
            DescribeTopicsResult describeResult = adminClient.describeTopics(Collections.singleton(topic));
            TopicDescription topicDescription = describeResult.values().get(topic).get();
            
            Map<Integer, List<Integer>> partitionInfo = new HashMap<>();
            for (TopicPartitionInfo partition : topicDescription.partitions()) {
                List<Integer> replicas = partition.replicas().stream()
                    .map(Node::id)
                    .collect(Collectors.toList());
                partitionInfo.put(partition.partition(), replicas);
            }
            
            return ResponseEntity.ok(partitionInfo);
        }
    }
    
    /**
     * 發(fā)送測(cè)試消息
     */
    @PostMapping("/send-test")
    public ResponseEntity<String> sendTestMessage(@RequestParam String topic) {
        OrderMessage testMessage = OrderMessage.builder()
            .orderId("TEST-" + System.currentTimeMillis())
            .userId("test-user")
            .amount(new BigDecimal("100.00"))
            .productName("測(cè)試商品")
            .quantity(1)
            .createTime(LocalDateTime.now())
            .status(OrderMessage.MessageStatus.PENDING)
            .build();
        
        kafkaTemplate.send(topic, testMessage.getOrderId(), testMessage);
        return ResponseEntity.ok("測(cè)試消息發(fā)送成功");
    }
    
    /**
     * 獲取消費(fèi)者組信息
     */
    @GetMapping("/consumer-groups")
    public ResponseEntity<Map<String, Object>> getConsumerGroups() throws Exception {
        try (AdminClient adminClient = AdminClient.create(kafkaAdmin.getConfigurationProperties())) {
            ListConsumerGroupsResult groupsResult = adminClient.listConsumerGroups();
            Collection<ConsumerGroupListing> groups = groupsResult.all().get();
            
            Map<String, Object> result = new HashMap<>();
            result.put("consumerGroups", groups);
            result.put("count", groups.size());
            
            return ResponseEntity.ok(result);
        }
    }
}

3.9 異常處理和重試機(jī)制

// KafkaExceptionHandler.java
@Component
@Slf4j
public class KafkaExceptionHandler {
    
    /**
     * 全局Kafka監(jiān)聽(tīng)器異常處理
     */
    @EventListener
    public void handleException(ListenerContainerConsumerFailedEvent event) {
        log.error("Kafka消費(fèi)者異常: {}", event.getContainer().getListenerId(), event.getException());
        
        // 記錄異常信息
        // 發(fā)送告警
        // 寫入錯(cuò)誤日志
    }
    
    /**
     * 自定義重試策略
     */
    @Bean
    public RetryTemplate kafkaRetryTemplate() {
        RetryTemplate retryTemplate = new RetryTemplate();
        
        // 重試策略:最多重試3次,每次間隔1秒
        SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy();
        retryPolicy.setMaxAttempts(3);
        
        // 退避策略
        FixedBackOffPolicy backOffPolicy = new FixedBackOffPolicy();
        backOffPolicy.setBackOffPeriod(1000L);
        
        retryTemplate.setRetryPolicy(retryPolicy);
        retryTemplate.setBackOffPolicy(backOffPolicy);
        
        return retryTemplate;
    }
}

3.10 健康檢查

// KafkaHealthIndicator.java
@Component
public class KafkaHealthIndicator implements HealthIndicator {
    
    @Autowired
    private KafkaTemplate<String, Object> kafkaTemplate;
    
    @Override
    public Health health() {
        try {
            // 嘗試發(fā)送一個(gè)測(cè)試消息來(lái)檢查Kafka連接
            kafkaTemplate.send("health-check-topic", "health-check", "ping")
                .get(5, TimeUnit.SECONDS);
            
            return Health.up()
                .withDetail("status", "Kafka集群連接正常")
                .withDetail("timestamp", LocalDateTime.now())
                .build();
            
        } catch (Exception e) {
            return Health.down()
                .withDetail("status", "Kafka集群連接異常")
                .withDetail("error", e.getMessage())
                .withDetail("timestamp", LocalDateTime.now())
                .build();
        }
    }
}

四、高可用性保障措施

4.1 集群配置建議

# kafka-server.properties 關(guān)鍵配置
broker.id=1
listeners=PLAINTEXT://:9092
advertised.listeners=PLAINTEXT://kafka1:9092
num.network.threads=3
num.io.threads=8
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600

# 日志配置
log.dirs=/data/kafka-logs
num.partitions=3
num.recovery.threads.per.data.dir=1

# 副本和ISR配置
offsets.topic.replication.factor=3
transaction.state.log.replication.factor=3
transaction.state.log.min.isr=2
default.replication.factor=3
min.insync.replicas=2

# 日志保留
log.retention.hours=168
log.segment.bytes=1073741824
log.retention.check.interval.ms=300000

# Zookeeper配置
zookeeper.connect=zk1:2181,zk2:2181,zk3:2181
zookeeper.connection.timeout.ms=6000

4.2 生產(chǎn)環(huán)境部署建議

硬件配置

  • 至少3個(gè)Kafka節(jié)點(diǎn) + 3個(gè)Zookeeper節(jié)點(diǎn)
  • SSD磁盤提高IO性能
  • 充足的內(nèi)存和CPU資源

網(wǎng)絡(luò)配置

  • 使用專用網(wǎng)絡(luò)
  • 配置合理的防火墻規(guī)則

監(jiān)控告警

  • 使用Kafka Manager或Confluent Control Center
  • 監(jiān)控指標(biāo):吞吐量、延遲、副本同步狀態(tài)

五、測(cè)試示例

// KafkaIntegrationTest.java
@SpringBootTest
@Slf4j
class KafkaIntegrationTest {
    
    @Autowired
    private KafkaProducerService producerService;
    
    @Test
    void testSendAndReceiveMessage() throws InterruptedException {
        // 創(chuàng)建測(cè)試消息
        OrderMessage orderMessage = OrderMessage.builder()
            .orderId("TEST-" + UUID.randomUUID())
            .userId("user-001")
            .amount(new BigDecimal("199.99"))
            .productName("測(cè)試商品")
            .quantity(2)
            .createTime(LocalDateTime.now())
            .status(OrderMessage.MessageStatus.PENDING)
            .build();
        
        // 發(fā)送消息
        SendResult<String, Object> result = producerService.sendOrderSync(orderMessage);
        
        assertNotNull(result);
        assertNotNull(result.getRecordMetadata());
        
        log.info("消息發(fā)送成功,分區(qū): {}, offset: {}", 
                result.getRecordMetadata().partition(),
                result.getRecordMetadata().offset());
        
        // 等待消費(fèi)者處理
        Thread.sleep(2000);
    }
    
    @Test
    void testBatchSend() {
        List<OrderMessage> messages = new ArrayList<>();
        for (int i = 0; i < 100; i++) {
            OrderMessage message = OrderMessage.builder()
                .orderId("BATCH-" + i)
                .userId("user-" + i)
                .amount(new BigDecimal(i * 10))
                .productName("商品" + i)
                .quantity(1)
                .createTime(LocalDateTime.now())
                .status(OrderMessage.MessageStatus.PENDING)
                .build();
            messages.add(message);
        }
        
        producerService.batchSendOrders(messages);
    }
}

六、總結(jié)

6.1 實(shí)現(xiàn)的高可用特性

  • 數(shù)據(jù)冗余:通過(guò)副本機(jī)制(Replication Factor=3)保證數(shù)據(jù)安全
  • 故障轉(zhuǎn)移:Leader選舉機(jī)制確保節(jié)點(diǎn)故障時(shí)自動(dòng)切換
  • 負(fù)載均衡:分區(qū)機(jī)制實(shí)現(xiàn)水平擴(kuò)展和負(fù)載均衡
  • 容錯(cuò)處理:死信隊(duì)列和重試機(jī)制保障消息不丟失
  • 監(jiān)控告警:完善的健康檢查和監(jiān)控體系

6.2 最佳實(shí)踐建議

  • 合理規(guī)劃分區(qū):根據(jù)業(yè)務(wù)吞吐量和消費(fèi)者數(shù)量設(shè)置分區(qū)數(shù)
  • 監(jiān)控副本同步:確保ISR(In-Sync Replicas)數(shù)量足夠
  • 配置重試機(jī)制:針對(duì)網(wǎng)絡(luò)波動(dòng)和臨時(shí)故障進(jìn)行重試
  • 實(shí)施消息冪等:避免重復(fù)消費(fèi)問(wèn)題
  • 定期清理數(shù)據(jù):設(shè)置合理的消息保留策略

6.3 性能優(yōu)化建議

  • 批量操作:使用批量發(fā)送和批量消費(fèi)提高吞吐量
  • 壓縮傳輸:?jiǎn)⒂孟嚎s減少網(wǎng)絡(luò)帶寬消耗
  • 合理批大小:根據(jù)業(yè)務(wù)場(chǎng)景調(diào)整批量大小
  • 異步確認(rèn):非關(guān)鍵業(yè)務(wù)使用異步發(fā)送提高響應(yīng)速度

通過(guò)以上方案,SpringBoot整合Kafka實(shí)現(xiàn)了高可用的消息隊(duì)列集群,具備生產(chǎn)級(jí)的可靠性、可擴(kuò)展性和容錯(cuò)能力,能夠滿足企業(yè)級(jí)應(yīng)用的需求。

以上就是SpringBoot整合Kafka實(shí)現(xiàn)高可用消息隊(duì)列集群詳解的詳細(xì)內(nèi)容,更多關(guān)于SpringBoot Kafka實(shí)現(xiàn)消息隊(duì)列集群的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • java依賴混亂存在的問(wèn)題與解決方案

    java依賴混亂存在的問(wèn)題與解決方案

    這篇文章主要為大家介紹了java依賴混亂存在的問(wèn)題與解決方案,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2023-09-09
  • 值得收藏!教你如何在IDEA中快速查看Java字節(jié)碼

    值得收藏!教你如何在IDEA中快速查看Java字節(jié)碼

    開(kāi)發(fā)中如果我們想看JVM虛擬機(jī)怎么編譯我們的Java文件,生成字節(jié)碼的,用IDEA工具就可以查看,本篇文章就給大家詳細(xì)介紹,對(duì)正在學(xué)習(xí)java的小伙伴們很有幫助,需要的朋友可以參考下
    2021-05-05
  • java并發(fā)等待條件的實(shí)現(xiàn)原理詳解

    java并發(fā)等待條件的實(shí)現(xiàn)原理詳解

    這篇文章主要介紹了java并發(fā)等待條件的實(shí)現(xiàn)原理詳解,還是比較不錯(cuò)的,這里分享給大家,供需要的朋友參考。
    2017-11-11
  • JAVA不使用線程池來(lái)處理的異步的方法詳解

    JAVA不使用線程池來(lái)處理的異步的方法詳解

    這篇文章主要介紹了JAVA不使用線程池來(lái)處理的異步的方法,在這個(gè)示例中,asyncTask方法創(chuàng)建了一個(gè)新的線程來(lái)執(zhí)行異步任務(wù),這個(gè)新線程會(huì)立即開(kāi)始執(zhí)行,而主線程則會(huì)繼續(xù)執(zhí)行后續(xù)的代碼,感興趣的朋友跟隨小編一起看看吧
    2024-05-05
  • MyBatis?Generator使用小結(jié)

    MyBatis?Generator使用小結(jié)

    本文主要介紹了MyBatis?Generator使用小結(jié),它能夠根據(jù)數(shù)據(jù)庫(kù)表,自動(dòng)生成java實(shí)體類、dao層接口及mapper.xml文件,具有一定的參考價(jià)值,感興趣的可以了解一下
    2023-11-11
  • Java基礎(chǔ)知識(shí)總結(jié)之繼承

    Java基礎(chǔ)知識(shí)總結(jié)之繼承

    這一篇我們來(lái)學(xué)習(xí)面向?qū)ο蟮牡诙€(gè)特征——繼承,文中有非常詳細(xì)的基礎(chǔ)知識(shí)總結(jié),對(duì)正在學(xué)習(xí)java的小伙伴們很有幫助,需要的朋友可以參考下
    2021-06-06
  • mybatis-plus @DS實(shí)現(xiàn)動(dòng)態(tài)切換數(shù)據(jù)源原理

    mybatis-plus @DS實(shí)現(xiàn)動(dòng)態(tài)切換數(shù)據(jù)源原理

    本文主要介紹了mybatis-plus @DS實(shí)現(xiàn)動(dòng)態(tài)切換數(shù)據(jù)源原理,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2022-07-07
  • Java初始化塊及執(zhí)行過(guò)程解析

    Java初始化塊及執(zhí)行過(guò)程解析

    這篇文章主要介紹了Java初始化塊及執(zhí)行過(guò)程解析,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2019-09-09
  • 在java中http請(qǐng)求帶cookie的例子

    在java中http請(qǐng)求帶cookie的例子

    今天小編就為大家分享一篇在java中http請(qǐng)求帶cookie的例子,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過(guò)來(lái)看看吧
    2019-08-08
  • java協(xié)變返回類型使用示例

    java協(xié)變返回類型使用示例

    在面向?qū)ο蟪绦蛟O(shè)計(jì)中,協(xié)變返回類型指的是子類中的成員函數(shù)的返回值類型不必嚴(yán)格等同于父類中被重寫的成員函數(shù)的返回值類型,而可以是更"狹窄"的類型
    2014-02-02

最新評(píng)論

平安县| 江都市| 石柱| 永登县| 万安县| 景洪市| 衡阳市| 平和县| 岱山县| 新田县| 璧山县| 丰城市| 南充市| 灵丘县| 永靖县| 岑巩县| 福清市| 井研县| 玉树县| 资溪县| 东莞市| 武安市| 凤山县| 茶陵县| 民勤县| 浦江县| 闻喜县| 临汾市| 探索| 绥德县| 江孜县| 万载县| 龙门县| 崇礼县| 阿拉善右旗| 富裕县| 太和县| 开化县| 山丹县| 宕昌县| 新兴县|