【问题标题】:Checking the existence of topic in kafka before creating in Java在 Java 中创建之前检查 kafka 中是否存在主题
【发布时间】:2015-08-14 19:48:45
【问题描述】:

我正在尝试使用以下方法在 kafka 0.8.2 中创建一个主题:

AdminUtils.createTopic(zkClient, myTopic, 2, 1, properties);

如果我在本地多次运行代码进行测试,则会失败,因为主题已经创建。有没有办法在创建主题之前检查主题是否存在? TopicCommand api 似乎没有为 listTopicsdescribeTopic 返回任何内容 .

【问题讨论】:

    标签: java apache-kafka


    【解决方案1】:
    public static void createKafkaTopic(String sourceTopicName, String sinkTopicName, String responseTopicName, String kafkaUrl) {
    
        try {
            Properties properties = new Properties();
            properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaUrl);
            AdminClient kafkaAdminClient = KafkaAdminClient.create(properties);
            ListTopicsResult topics = kafkaAdminClient.listTopics();
            Set <String> names = topics.names().get();
    
            boolean containsSourceTopic = names.contains(sourceTopicName);
            boolean containsSinkTopic = names.contains(sinkTopicName);
            boolean containsResponseTopic = names.contains(responseTopicName);
    
            if (!containsResponseTopic && !containsSinkTopic && !containsSourceTopic) {
                CreateTopicsResult result = kafkaAdminClient.createTopics(
                        Stream.of(sourceTopicName, sinkTopicName, responseTopicName).map(
                                name -> new NewTopic(name, 1, (short) 1)
                        ).collect(Collectors.toList())
                );
                result.all().get();
                LOG.info("new sourceTopicName: {}, sinkTopicName: {}, responseTopicName: {} are created",
                        sourceTopicName, sinkTopicName, responseTopicName);
            }
        } catch (ExecutionException | InterruptedException e) {
            LOG.info("Error message {}", e.getMessage());
        }
    }
    

    【讨论】:

    • 虽然这段代码 sn-p 可以解决问题,including an explanation 确实有助于提高您的帖子质量。请记住,您是在为将来的读者回答问题,而这些人可能不知道您提出代码建议的原因。
    【解决方案2】:

    您可以使用 kakfa-client 版本 0.11.0.0 中的 AdminClient

    示例代码:

        Properties config = new Properties();
        config.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "localhist:9091");
    
        AdminClient admin = AdminClient.create(config);
        ListTopicsResult listTopics = admin.listTopics();
        Set<String> names = listTopics.names().get();
        boolean contains = names.contains("TEST_6");
        if (!contains) {
            List<NewTopic> topicList = new ArrayList<NewTopic>();
            Map<String, String> configs = new HashMap<String, String>();
            int partitions = 5;
            Short replication = 1;
            NewTopic newTopic = new NewTopic("TEST_6", partitions, replication).configs(configs);
            topicList.add(newTopic);
            admin.createTopics(topicList);
        }
    

    【讨论】:

      【解决方案3】:

      为此,您可以使用方法AdminUtils.topicExists(ZkUtils zkClient, String topic),如果主题已经存在,它将返回true,否则返回false

      您的代码将是这样的:

      if (!AdminUtils.topicExists(zkClient, myTopic)){
          AdminUtils.createTopic(zkClient, myTopic, 2, 1, properties);
      }
      

      【讨论】:

        猜你喜欢
        • 2015-09-05
        • 1970-01-01
        • 2020-07-29
        • 2017-10-02
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2013-09-03
        • 1970-01-01
        相关资源
        最近更新 更多