【问题标题】:Kafka: Delete idle consumer group idKafka:删除空闲的消费者组ID
【发布时间】:2021-01-16 02:43:40
【问题描述】:
在某些情况下,我使用 Kafka-stream 对主题的小内存(哈希图)投影进行建模。 K,V 缓存确实需要一些操作,因此它不是 GlobalKTable 的好案例。在这样的“缓存”场景中,我希望我的所有兄弟实例都拥有相同的缓存,因此我需要绕过消费者组机制。
要启用此功能,我通常只需使用随机生成的应用程序 ID 启动我的应用程序,因此每个应用程序每次重新启动时都会重新加载主题。唯一需要注意的是,我最终得到了一些在 kafka 代理上孤立的消费者组,直到 offsets.retention.minutes,这对于我们的操作监控工具来说并不理想。
知道如何解决这个问题吗?
- 我们能否将 applicationId 配置为临时的,以便在应用死机后使其消失?
- 或者我们可以强制消费者仅在本地管理其偏移量吗?
- 或者是否有一些 java adminApi 可以在正常关闭应用程序时用来清理我的 consumer-group-id?
谢谢
【问题讨论】:
标签:
apache-kafka
kafka-consumer-api
apache-kafka-streams
【解决方案1】:
AdminClient 中有一个名为 deleteConsumerGroups 的 Java API,可用于删除单个 ConsumerGroup。
您可以将其与 Kafka 2.5.0 一起使用,如下所示。
import java.util.Arrays;
import java.util.Properties;
import java.util.concurrent.ExecutionException;
import org.apache.kafka.clients.admin.*;
import org.apache.kafka.common.KafkaFuture;
public class DeleteConsumerGroups {
public static void main(String[] args) {
System.out.println("*** Starting AdminClient to delete a Consumer Group ***");
final Properties properties = new Properties();
properties.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
properties.put(AdminClientConfig.REQUEST_TIMEOUT_MS_CONFIG, "1000");
properties.put(AdminClientConfig.DEFAULT_API_TIMEOUT_MS_CONFIG, "5000");
AdminClient adminClient = AdminClient.create(properties);
String consumerGroupToBeDeleted = "console-consumer-65092";
DeleteConsumerGroupsResult deleteConsumerGroupsResult = adminClient.deleteConsumerGroups(Arrays.asList(consumerGroupToBeDeleted));
KafkaFuture<Void> resultFuture = deleteConsumerGroupsResult.all();
try {
resultFuture.get();
} catch (InterruptedException e) {
e.printStackTrace();
} catch (ExecutionException e) {
e.printStackTrace();
}
adminClient.close();
}
}
运行上述代码前的ConsumerGroups List
$ kafka-consumer-groups --bootstrap-server localhost:9092 --list
console-consumer-65092
console-consumer-53268
上面代码运行后的ConsumerGroups List
$ kafka-consumer-groups --bootstrap-server localhost:9092 --list
console-consumer-53268