【发布时间】: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