【问题标题】:Spring cloud stream kafka pause/resume bindersSpring Cloud Stream kafka 暂停/恢复活页夹
【发布时间】:2019-04-27 19:39:26
【问题描述】:
【问题讨论】:
标签:
java
apache-kafka
spring-cloud-stream
spring-kafka
circuit-breaker
【解决方案1】:
对不起,我看错了你的问题。
您可以自动连接 BindingsEndpoint,但不幸的是,它的 State 枚举是私有的,因此您不能以编程方式调用 changeState()。
我有opened an issue for this。
编辑
你可以用反射来做,但它有点难看......
@SpringBootApplication
@EnableBinding(Sink.class)
public class So53476384Application {
public static void main(String[] args) {
SpringApplication.run(So53476384Application.class, args);
}
@Autowired
BindingsEndpoint binding;
@Bean
public ApplicationRunner runner() {
return args -> {
Class<?> clazz = ClassUtils.forName("org.springframework.cloud.stream.endpoint.BindingsEndpoint$State",
So53476384Application.class.getClassLoader());
ReflectionUtils.doWithMethods(BindingsEndpoint.class, method -> {
try {
method.invoke(this.binding, "input", clazz.getEnumConstants()[2]); // PAUSE
}
catch (InvocationTargetException e) {
e.printStackTrace();
}
}, method -> method.getName().equals("changeState"));
};
}
@StreamListener(Sink.INPUT)
public void listen(String in) {
}
}