【问题标题】:spring kafka consumer with circuit breaker functionality using Resilience4j library使用 Resilience4j 库的具有断路器功能的 spring kafka 消费者
【发布时间】:2022-01-09 09:22:24
【问题描述】:

我正在尝试实现 spring kafka 消费者,在处理事件时需要在某个异常后暂停(例如:将事件信息存储到 DB 时,DB 已关闭)。

我们如何使用 Resilience4j 断路器方法和 spring boot - 2.3.8 (spring kafka) 来处理这种情况

寻找一些关于消费者暂停和恢复的例子。

@Component
public class CircuitBreakerManager {

    private CircuitBreaker circuitBreaker;

    @Autowired
    private KafkaListenerEndpointRegistry kafkaListenerEndpointRegistry;

    public CircuitBreakerManager() {
        CircuitBreakerConfig circuitBreakerConfig = CircuitBreakerConfig.custom()
                .slidingWindowType(CircuitBreakerConfig.SlidingWindowType.COUNT_BASED)
                .enableAutomaticTransitionFromOpenToHalfOpen()
                .minimumNumberOfCalls(5)
                .permittedNumberOfCallsInHalfOpenState(3)
                .slidingWindowSize(10)
                .failureRateThreshold(50)
                .slowCallRateThreshold(60.0f)
                .slowCallDurationThreshold(Duration.ofSeconds(3))
                .build();
        CircuitBreakerRegistry registry = CircuitBreakerRegistry.of(circuitBreakerConfig);
        this.circuitBreaker = registry.circuitBreaker("serialization_exception");
        this.circuitBreaker.getEventPublisher().onStateTransition(this::onStateChange);
    }

    private void onStateChange(CircuitBreakerOnStateTransitionEvent circuitBreakerEvent) {
        CircuitBreaker.State toState = circuitBreakerEvent.getStateTransition()
                .getToState();
        System.out.println("Change in Circuit Breaker state " + toState);
        switch (toState) {
            case OPEN:
                    kafkaListenerEndpointRegistry.getListenerContainer("my_listener_id").stop();
                break;
            case CLOSED:
                break;
            case HALF_OPEN:
                kafkaListenerEndpointRegistry.getListenerContainer("my_listener_id").start();
                break;
        }
    }


}

在 kafka 监听器只是想捕获解析错误。如果我们得到超过 5 个解析错误,则需要停止监听器。但我不确定断路器将如何被触发。

@CircuitBreaker(name = RESILIENCE4J_INSTANCE_NAME)
    private Event getParsedEvent(ConsumerRecord consumerRecord) {
        Event event = getEvent(consumerRecord);
        
            if (StringUtils.isEmpty(event)) {
                
                throw new RuntimeException("Serialization Exception occurred");
            }
        }
        return event;
    }

【问题讨论】:

  • 如果您只是想暂停消费,那么您也可以通过在 DB 关闭时不确认偏移量来实现。这样就不会有额外的轮询,但如果达到总会话超时,它可能会重新平衡
  • 感谢 Suraj 的意见。似乎断路器方法更方便、更干净,可以让听众在一段时间内度过。

标签: spring-boot spring-kafka resilience4j


【解决方案1】:

Pausing and Resuming Listener Containers

请注意,在当前轮询返回的所有记录都已处理完毕(或侦听器抛出异常,只要默认错误处理程序到位)后,暂停才会生效。

【讨论】:

  • 嗨,加里。谢谢!我想出了代码。将在实际帖子中显示。
  • 嗨,Gary.. 只是想和你核实一下,考虑到断路器模式,我如何在应用程序中启动和停止监听器?我已经添加了代码,但不确定我们如何触发它。
  • 我对resilience4j不熟悉;您正在正确使用端点注册表;但是,使用 stop/start 而不是 pause/resume 将完全停止容器并将分区分配给另一个实例(如果存在)。暂停使分区分配保持原样,但阻止接收更多记录(直到恢复)。
  • 感谢加里的回复。将断路器放在一边,暂停和恢复将如何在批处理模式下工作?
  • 同记录模式,暂停在当前批处理完成或抛出异常后生效。
猜你喜欢
  • 2021-05-09
  • 2020-09-08
  • 2020-01-27
  • 2018-10-13
  • 1970-01-01
  • 2020-11-02
  • 2023-03-20
  • 2020-08-11
  • 1970-01-01
相关资源
最近更新 更多