【问题标题】:Is the number of Kafka partitions accessible when using Camel-Kafka?使用 Camel-Kafka 时可以访问的 Kafka 分区数吗?
【发布时间】:2017-03-09 16:50:13
【问题描述】:

当使用 camel-kafka 时,有没有办法 a) 调用默认的 Kafka 分区器,和/或 b) 确定用于分区算法的主题的分区数?在我目前拥有的骆驼处理器中

exchange.getIn().setHeader(KafkaConstants.PARTITION_KEY, partition);

但我无法在不知道有多少分区可用的情况下为partition 选择一个值。我可以使用 Camel 来执行此操作,还是需要“转为原生”并调用 KafkaProducer.partitionsFor

【问题讨论】:

    标签: apache-camel apache-kafka


    【解决方案1】:

    我找到了答案。我可以通过省略设置KafkaConstants.PARTITION_KEY 标头来调用默认分区程序。

    我可以通过将 Partitioner 指定为 Camel 选项来使用它的自定义实现:

    from(...)
        ...
        .to("kafka:{{kafka.host}}:{{kafka.port}}"
            + "?topic={{kafka.topic}}"
            + "&partitioner=my.package.MyPartitioner");
    

    在 MyPartitioner 中,我可以通过 cluster 参数获取分区信息:

    @Override
    public int partition(String topic,
        Object key, byte[] keyBytes, Object value, byte[] valueBytes,
        Cluster cluster) {
    
      List<PartitionInfo> partitions = cluster.availablePartitionsForTopic(topic);
    
      // use partitions.size()...
    
      return ...    
    }
    

    【讨论】:

      猜你喜欢
      • 2020-10-08
      • 2021-07-27
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-04-10
      • 1970-01-01
      • 1970-01-01
      • 2021-03-10
      相关资源
      最近更新 更多