为 Spring Kafka 设置 authorizationExceptionRetryInterval
Setting authorizationExceptionRetryInterval for Spring Kafka
任何人都知道如何设置新的 属性: authorizationExceptionRetryInterval 而无需手动创建 ConcurrentKafkaListenerContainerFactory。
我本来想说...
@Component
class ContainerFactoryCustomizer {
ContainerFactoryCustomizer(AbstractKafkaListenerContainerFactory<?, ?, ?> factory) {
factory.setContainerCustomizer(
container -> container.getContainerProperties()
.setAuthorizationExceptionRetryInterval(Duration.ofSeconds(10L)));
}
}
但这不起作用,due to a bug (the container customizer is not set up)。
这是一个解决方法:
@SpringBootApplication
public class So60054097Application {
public static void main(String[] args) {
SpringApplication.run(So60054097Application.class, args);
}
@KafkaListener(id = "so60054097", topics = "so60054097", autoStartup = "false")
public void listen(String in) {
System.out.println(in);
}
@Bean
public NewTopic topic() {
return TopicBuilder.name("so60054097").partitions(1).replicas(1).build();
}
@Bean
public ApplicationRunner runner(KafkaListenerEndpointRegistry registry) {
return args -> {
MessageListenerContainer container = registry.getListenerContainer("so60054097");
container.getContainerProperties()
.setAuthorizationExceptionRetryInterval(Duration.ofSeconds(10L));
container.start();
};
}
}
(将 autoStartup
设置为 false;修复 属性 并启动容器)。
任何人都知道如何设置新的 属性: authorizationExceptionRetryInterval 而无需手动创建 ConcurrentKafkaListenerContainerFactory。
我本来想说...
@Component
class ContainerFactoryCustomizer {
ContainerFactoryCustomizer(AbstractKafkaListenerContainerFactory<?, ?, ?> factory) {
factory.setContainerCustomizer(
container -> container.getContainerProperties()
.setAuthorizationExceptionRetryInterval(Duration.ofSeconds(10L)));
}
}
但这不起作用,due to a bug (the container customizer is not set up)。
这是一个解决方法:
@SpringBootApplication
public class So60054097Application {
public static void main(String[] args) {
SpringApplication.run(So60054097Application.class, args);
}
@KafkaListener(id = "so60054097", topics = "so60054097", autoStartup = "false")
public void listen(String in) {
System.out.println(in);
}
@Bean
public NewTopic topic() {
return TopicBuilder.name("so60054097").partitions(1).replicas(1).build();
}
@Bean
public ApplicationRunner runner(KafkaListenerEndpointRegistry registry) {
return args -> {
MessageListenerContainer container = registry.getListenerContainer("so60054097");
container.getContainerProperties()
.setAuthorizationExceptionRetryInterval(Duration.ofSeconds(10L));
container.start();
};
}
}
(将 autoStartup
设置为 false;修复 属性 并启动容器)。