【问题标题】:Consume messages from Kafka only if some condition true仅在某些条件为真时才使用来自 Kafka 的消息
【发布时间】:2019-12-16 05:29:48
【问题描述】:
我们有特定的主题,只有在条件 consumeEnabled=true 时才需要消费消息。
所以,它应该像这样工作:
- 如果应用程序正在启动并且consumeEnabled=true,则分配
分区给消费者并使用来自主题的消息。
- 如果应用程序正在启动并且consumeEnabled=false,则不要将分区分配给消费者,也不要使用来自主题的消息。
- 如果应用程序已在使用 consumeEnabled=false 的情况下运行,但在运行时属性变为 consumeEnabled=true,则在运行时将分区分配给消费者并从主题消费消息。
应用正在消费消息的情况,但随后consumeEnabled变为false,无需考虑。
请您定义使用 Spring Kafka 和/或 Kafka Java 客户端实现决策的最佳方式
【问题讨论】:
标签:
java
apache-kafka
spring-kafka
【解决方案1】:
如果您使用的是@KafkaListener 那么
@KafkaListener(id = "foo", ... , autoStartup="${consume.enabled}")
consume.enabled 是一个属性。
要在运行时启动/停止容器,请使用 KafkaListenerEndpointRegistry bean。
registry.getListenerContainer("foo").start();
【解决方案2】:
您可以将您的消费者置于一个简单的线程中,以切换消费者对象的轮询状态。
public class EnabledConsumer implements Runnable {
private Consumer consumer;
private boolean enabled;
public EnabledConsumer(Consumer consumer, boolean enabled) {
this.consumer = consumer;
this.enabled = enabled;
}
public void setEnabled(boolean enable) {
this.enabled = enable;
}
@Override
public void run() {
while(enabled) {
ConsumerRecords records = consumer.poll(...);
...
}
}