【问题标题】:Check for the existence of a Kafka topic programatically in Java在 Java 中以编程方式检查是否存在 Kafka 主题
【发布时间】:2018-10-24 20:03:19
【问题描述】:

我如何知道主题是否已在 Kafka 集群中以编程方式创建,没有使用 CLI 工具,并且在尝试生成主题之前?

我遇到了一个主题不存在的问题,我们的应用程序正在尝试生成一个不存在的主题,但它仅在 90 秒后才收到通知(元数据超时)。我想知道是否有办法从 Java 代码中知道主题是否存在,以便我们可以在实际尝试发送消息之前进行检查。我想我可以看看 Kafka CLI utils 使用的代码,但我想知道是否有我可能错过的 API 或更简单的方法。

【问题讨论】:

  • CLI 工具只是您最终要编写的程序化 Java API 的包装器,那么为什么不看看这些代码呢?

标签: java apache-kafka


【解决方案1】:

您可以使用AdminClient#listTopics() 来检查给定主题是否存在,如下所示:

Properties props = new Properties();
props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
try (AdminClient client = AdminClient.create(props)) {
    ListTopicsOptions options = new ListTopicsOptions();
    options.listInternal(true); // includes internal topics such as __consumer_offsets
    ListTopicsResult topics = client.listTopics(options);
    Set<String> currentTopicList = topics.names().get();
    // do your filter logic here......
}

【讨论】:

  • 如果我们想检查单个主题是否存在并且 Kafka 上可能有数千个主题,这个 API 似乎效率不高。
  • 如果您担心潜在的性能低下,您可以改用AdminClient.describeTopics 并通过捕获 UnknownTopicOrPartitionException 来检查是否存在。
  • 我不这么认为。 AdminClient 在得到响应后迅速抛出UnknownTopicOrPartitionException 用于不存在的主题。您可以通过尝试捕获此异常来了解存在。就性能而言,您的意思是捕获异常会消耗更多的 JVM 指令吗?
  • 构造和展开异常是昂贵的。异常应该用于异常情况,而不是正常流程。见这里:shipilev.net/blog/2014/exceptional-performance
【解决方案2】:

您可以使用AdminUtils.topicExists(..) 方法来检查旧版kafka (1.0.0) 是否存在主题:

    int sessionTimeOutInMs = 15 * 1000;
    int connectionTimeOutInMs = 10 * 1000;
    String zkHost = "localhost:2181";
    ZkClient zkClient = new ZkClient(zkHost, sessionTimeOutInMs, connectionTimeOutInMs, ZKStringSerializer$.MODULE$);
    ZkUtils zkUtils = new ZkUtils(zkClient, new ZkConnection(zkHost), false);
    System.out.println(AdminUtils.topicExists(zkUtils, "TopicName"));

AdminUtils 在最近的 Kafka 版本中已被弃用。所以你可以将AdminClient 用于kafka 1.0 +:

    Properties prop = new Properties();
    prop.setProperty("bootstrap.servers", "localhost:9092");
    AdminClient admin = AdminClient.create(prop);
    boolean topicExists = admin.listTopics().names().get().stream().anyMatch(topicName -> topicName.equalsIgnoreCase("tealium.topic"));

【讨论】:

  • 所有这些 API 都已弃用,因此我不建议使用它们。而是使用AdminClient
猜你喜欢
  • 2015-09-05
  • 2015-08-14
  • 2017-10-02
  • 1970-01-01
  • 2021-12-10
  • 1970-01-01
  • 2019-11-24
  • 2011-05-05
  • 2021-01-13
相关资源
最近更新 更多