【问题标题】:Kafka Consumer with Circuit Breaker, Retry Patterns using Resilience4j带有断路器的 Kafka 消费者,使用 Resilience4j 重试模式
【发布时间】:2021-05-09 23:39:33
【问题描述】:

我需要一些帮助来了解如何使用 Spring boot、Kafka、Resilence4J 提出解决方案,以实现来自我的 Kafka Consumer 的微服务调用。假设如果微服务关闭,那么我需要使用断路器模式通知我的 Kafka 消费者停止获取消息/事件,直到微服务启动并运行。

【问题讨论】:

  • 您好,澄清一下。您有一个服务 A 正在使用来自某个主题的消息并为每个新消息调用另一个服务 B?如果服务 B 关闭,您想停止使用该主题的消息吗?
  • @RobertWinkler 绝对正确。并在打开电路之前尝试重试

标签: spring-boot apache-kafka spring-cloud-stream resilience4j


【解决方案1】:

使用 Spring Kafka,您可以根据 CircuitBreaker 状态转换使用 pauseresume 方法。我为此找到的最好方法是将其定义为带有 @Configuration 注释的“主管”。还使用了 Resilience4j。

@Configuration
public class CircuitBreakerConsumerConfiguration {

public CircuitBreakerConsumerConfiguration(CircuitBreakerRegistry circuitBreakerRegistry, KafkaManager kafkaManager) {
    circuitBreakerRegistry.circuitBreaker("yourCBName").getEventPublisher().onStateTransition(event -> {
  
        switch (event.getStateTransition()) {
            case CLOSED_TO_OPEN:
            case CLOSED_TO_FORCED_OPEN:
            case HALF_OPEN_TO_OPEN:
                kafkaManager.pause();
                break;
            case OPEN_TO_HALF_OPEN:
            case HALF_OPEN_TO_CLOSED:
            case FORCED_OPEN_TO_CLOSED:
            case FORCED_OPEN_TO_HALF_OPEN:
                kafkaManager.resume();
                break;
            default:
                throw new IllegalStateException("Unknown transition state: " + event.getStateTransition());
        }
    });
   }
}

这是我与带有 @Component 注释的 KafkaManager 结合使用的。

@Component
public class KafkaManager {
  private final KafkaListenerEndpointRegistry registry;

  public KafkaManager(KafkaListenerEndpointRegistry registry) {
    this.registry = registry;
  }
  public void pause() {   
    registry.getListenerContainers().forEach(MessageListenerContainer::pause);
  }

  public void resume() {
    registry.getListenerContainers().forEach(MessageListenerContainer::resume);
  }
}

此外,我的消费者服务如下所示:

  @KafkaListener(topics = "#{'${topic.name}'}", concurrency = "1", id = "CBListener")
public void receive(final ConsumerRecord<String, ReplayData> replayData, Acknowledgment acknowledgment) throws
        Exception {

    try {
        httpClientServiceCB.receiveHandleCircuitBreaker(replayData);
        acknowledgement.acknowledge();
    } catch (Exception e) {
        acknowledgment.nack(1000);
    }
}

还有@CircuitBreaker注解:

@CircuitBreaker(name = "yourCBName")
public void receiveHandleCircuitBreaker(ConsumerRecord<String, ReplayData> replayData) throws
        Exception {
    try {
        String response = restTemplate.getForObject("http://localhost:8081/item", String.class);
    } catch (Exception e                                                                       ) {
       
        // throwing the exception is needed to trigger the Circuit Breaker state change
        throw new Exception();
    }
}

另外补充以下application.properties

  resilience4j.circuitbreaker.instances.yourCBName.failure-rate-threshold=80
  resilience4j.circuitbreaker.instances.yourCBName.sliding-window-type=COUNT_BASED
  resilience4j.circuitbreaker.instances.yourCBName.sliding-window-size=5
  resilience4j.circuitbreaker.instances.yourCBName.wait-duration-in-open-state=10000
  resilience4j.circuitbreaker.instances.yourCBName.automatic-transition-from-open-to-half-open-enabled=true
  spring.kafka.consumer.enable.auto.commit = false
  spring.kafka.listener.ack-mode = MANUAL_IMMEDIATE

也可以看看https://resilience4j.readme.io/docs/circuitbreaker

【讨论】:

    【解决方案2】:

    如果您使用的是 Spring Kafka,则可以使用 ConcurrentMessageListenerContainer 类的 pauseresume 方法。 您可以将 EventListener 附加到 CircuitBreaker,它侦听状态转换并暂停或恢复事件处理。将 CircuitBreakerRegistry 注入你的 bean:

    circuitBreakerRegistry.circuitBreaker("yourCBName").getEventPublisher().onStateTransition(
                            event -> {
                                switch (event.getStateTransition()) {
                                    case CLOSED_TO_OPEN:
                                        container.pause();
                                    case OPEN_TO_HALF_OPEN:
                                        container.resume();
                                    case HALF_OPEN_TO_CLOSED:
                                        container.resume();
                                    case HALF_OPEN_TO_OPEN:
                                        container.pause();
                                    case CLOSED_TO_FORCED_OPEN:
                                        container.pause();
                                    case FORCED_OPEN_TO_CLOSED:
                                        container.resume();
                                    case FORCED_OPEN_TO_HALF_OPEN:
                                        container.resume();
                                    default:
                                }
                            }
                    );
    

    【讨论】:

      猜你喜欢
      • 2022-01-09
      • 2014-07-30
      • 2020-01-27
      • 2019-02-25
      • 1970-01-01
      • 2016-09-01
      • 2017-06-22
      • 2019-03-24
      • 2020-04-02
      相关资源
      最近更新 更多