【问题标题】: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 ...
}