use*_*814 10 java spring spring-boot spring-kafka
我有一个春季启动kafka应用程序。我的经纪人每隔几天就会被回收一次。旧的经纪人已取消配置,新的经纪人已配置。
我有一个调度程序,每隔几个小时要检查一次代理。我想确保一旦有了新的经纪人,我们就应该重新加载所有与Spring Kafka相关的bean。与KafkaAutoConfiguration非常相似,除了我希望触发代理值更改并以编程方式加载自动配置。
每当将旧的代理替换为新的代理时,如何以编程方式调用自动配置?
您的要求听起来像是Spring Cloud中的Config Server:https : //cloud.spring.io/spring-cloud-static/Greenwich.SR2/multi/multi__spring_cloud_config_2.html#_spring_cloud_config_2及其@RefreshScope功能:https : //cloud.spring.io /spring-cloud-static/Greenwich.SR2/multi/multi__spring_cloud_context_application_context_services.html#refresh-scope。
因此,您需要指定自己的bean并用该注释标记它们:
@Bean
@RefreshScope
public ConsumerFactory<?, ?> kafkaConsumerFactory() {
return new DefaultKafkaConsumerFactory<>(this.properties.buildConsumerProperties());
}
@Bean
@RefreshScope
public ProducerFactory<?, ?> kafkaProducerFactory() {
DefaultKafkaProducerFactory<?, ?> factory = new DefaultKafkaProducerFactory<>(
this.properties.buildProducerProperties());
String transactionIdPrefix = this.properties.getProducer().getTransactionIdPrefix();
if (transactionIdPrefix != null) {
factory.setTransactionIdPrefix(transactionIdPrefix);
}
return factory;
}
Run Code Online (Sandbox Code Playgroud)
这两个bean依赖于配置属性来连接到Apache Kafka代理,这确实足够使它们可刷新。每当ContextRefreshedEvent发生这种情况时,将使用新的配置属性来重新初始化这些bean。
我认为ConsumerFactory使用者(MessageListenerContainer和KafkaListenerEndpointRegistry)也必须在该事件上重新启动。关键是MessageListenerContainer开始一个长期的过程,因此KafkaConsumer为此poll目的缓存了一个实例。
所有ProducerFactory使用者不需要重新启动。即使KafkaProducer 被缓存在阶段DefaultKafkaProducerFactory,也要重新初始化@RefreshScope。
更新
我不使用配置服务器。我从领事目录服务中获得了新主机。
是的,我并不是说您使用的是配置服务器。那只是以相似的方式寻找我。因此,从高处看,我真的会研究Consul目录解决方案的Config Client实现。
不过,您仍然可以发出,RefreshEvent这将触发您所有的@RefreshScope'd豆重新加载。为此,ApplicationEventPublisherAware无论何时从Consul更新,您都需要实现并发出该事件。切记:必须重新启动Kafka侦听器容器。为此,您可以听,RefreshScopeRefreshedEvent因为您只有在所有@RefreshScope刷新后才真正对重启感兴趣。
有关刷新范围的更多信息:https : //gist.github.com/dsyer/a43fe5f74427b371519af68c5c4904c7
| 归档时间: |
|
| 查看次数: |
424 次 |
| 最近记录: |