【发布时间】:2020-06-10 22:55:35
【问题描述】:
Kafka 消费者 API 非常好,可以隐藏任何暂时的连接错误,并且如果 Kafka 代理死亡并再次出现,只需从其当前偏移量中读取数据。
但在某些应用程序中,如果整个 Kafka 集群已关闭(即所有代理),则发出警报并停止处理数据(来自其他来源)很重要。 我浏览了杂项。 API,这似乎不是一项功能。
我最接近的方法是提交管理员调用,并根据超时得出 Kafka 集群已关闭的结论:
Properties properties = ... // Load properties from somewhere.
int timeout = 5_000; // 5 second timeout
AdminClient adminClient = AdminClient.create(properties);
try {
adminClient.listTopics(new ListTopicsOptions().timeoutMs(timeout)).listings().get();
// Here we know the cluster is up as call returned within timeout.
} catch (ExecutionException ex) {
// Here we know that the cluster is down as the call timed out.
}
这是最好的方法吗?
另一种方法是查询 ZooKeeper,但上述方法也适用于应用程序和 Kafka 之间存在网络问题的情况。
【问题讨论】:
-
您的具体用例是什么?通常,“Kafka 集群已关闭”是一项监控任务,而不是针对消费者软件的任务。
-
假设我有一个通过 Kafka 控制的应用程序,如果 Kafka 出现故障,它需要停止运行。 (单独监控很容易)。
-
消费者无论如何都会抛出异常,但会继续重试,假设你使用了一个while循环并且没有断路器
-
我会避免查询 Zookeeper 并观看 kip500
标签: java apache-kafka