【问题标题】:Can a Kafka producer create topics and partitions?Kafka 生产者可以创建主题和分区吗?
【发布时间】:2017-04-27 22:03:58
【问题描述】:

目前我正在评估不同的消息传递系统。 有一个与 Apache Kafka 相关的问题,我自己无法回答。

Kafka 生产者是否可以动态创建主题和分区(也在现有主题上)? 如果是,有什么缺点吗?

提前致谢

【问题讨论】:

    标签: apache-kafka messaging kafka-producer-api


    【解决方案1】:

    更新:

    kafka broker 有一个属性: auto.create.topics.enable

    如果您将其设置为 true,如果生产者使用新主题名称向主题发布消息,它将自动为您创建一个主题。

    Confluent 团队建议不要这样做,因为主题的爆炸式增长(取决于您的环境)可能会变得笨拙,并且主题创建在创建时将始终具有相同的默认值。复制因子至少为 3 非常重要,以确保您的主题在发生磁盘故障时的持久性。

    【讨论】:

    • 谢谢,我想为每个设备(生产者)设置一个主题/分区。我不知道会有多少设备,所以我想动态添加它们。上述解决方案听起来有点“迟钝”。我想经典的 Pub/Sub 系统可能会更好。
    【解决方案2】:

    当您启动 kafka 代理时,您可以在 conf/server.properties 文件中定义一堆属性。其中一个属性是auto.create.topics.enable,如果您将其设置为true(默认情况下),当您向不存在的主题发送消息时,kafka 将自动创建一个主题。分区号将由同一​​文件中的默认设置定义。

    缺点:据我所知,以这种方式创建的主题将始终具有相同的默认设置(分区、副本...)。

    【讨论】:

      【解决方案3】:

      如果需要,您可以从 java 创建主题。是否推荐,取决于用例。例如。如果您的主题名称是生产者的传入有效负载的函数,则它可能很有用。以下是在 kafka 0.10.x 中工作的代码 sn-p

      void createTopic(String zookeeperConnect, String topicName) throws InterruptedException {
          int sessionTimeoutMs = <some-int-value>;
          int connectionTimeoutMs = <some-int-value>;
      
          ZkClient zkClient = new ZkClient(zookeeperConnect, sessionTimeoutMs, connectionTimeoutMs, ZKStringSerializer$.MODULE$);
      
          boolean isSecureKafkaCluster = false;
          ZkUtils zkUtils = new ZkUtils(zkClient, new  ZkConnection(zookeeperConnect), isSecureKafkaCluster);
      
          Properties topicConfig = new Properties();
          try {
            AdminUtils.createTopic(zkUtils, topicName, 1, 1, topicConfig,
            RackAwareMode.Disabled$.MODULE$);
          } catch (TopicExistsException ex) {
          //log it 
          }
          zkClient.close();
      }
      

      注意:只允许增加编号。的分区。

      【讨论】:

      • 我们使用类似的方法来动态创建主题。分区呢?
      • @user2105282 AdminUtils.createTopic() 方法将分区数和复制数作为参数。因此,您可以相应地选择它们。
      【解决方案4】:

      对于任何消息传递系统,我认为不推荐由生产者动态创建主题/分区或任何队列。

      对于您的用例,您可能可以使用 device_id 作为分区键来区分消息。这样您就可以使用一个主题。

      【讨论】:

      • 我想过。问题是,我不知道所有设备/设备 ID。或者换句话说,我想添加动态发布数据的设备。
      • 我认为您不必担心预测密钥(即设备)。默认情况下,Kafka 会随机分配分区。如果您想按设备(即密钥)进行分隔,您可以创建一个按密钥名称过滤的流。
      • @Girdhar 如果使用 device_id,消费者必须先读取主题中的所有消息,然后通过 device_id 过滤才能获取相关数据,是吗?
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2018-07-27
      • 2017-04-28
      • 2020-03-30
      • 2016-12-31
      • 2018-02-05
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多