【发布时间】:2016-05-11 11:16:39
【问题描述】:
在 Flume 中,当使用 Kafka 通道时,有没有办法影响事件发送到哪个分区?
对于 Kafka sink,key FlumeEvent 标头显然用于选择分区,但我找不到任何有关 Kafka channel 分区的文档。 p>
【问题讨论】:
标签: apache-kafka flume flume-ng
在 Flume 中,当使用 Kafka 通道时,有没有办法影响事件发送到哪个分区?
对于 Kafka sink,key FlumeEvent 标头显然用于选择分区,但我找不到任何有关 Kafka channel 分区的文档。 p>
【问题讨论】:
标签: apache-kafka flume flume-ng
Flume 的 Kafka 通道不支持像 KafkaSink 那样将 Event 标头映射到开箱即用的分区键。
但是,对其进行修改以使其能够做到这一点并不太复杂。由于我不确定我是否可以分享代码,我将给出指示:
serializedEvents 类型中的 byte[] 替换为也可以为每个事件(内部类,甚至是 Kafka KeyedMessage<String, byte[]>)持有一个 String 键的东西KafkaTransaction.doPut(Event event)中,从headers中获取key并与序列化消息一起存储在serializedEvents中KafkaTransaction.doCommit()中,使用与序列化事件一起存储的密钥而不是batchUUID。注意事务中的事件将不再保证由通道消费者端的单个 KafkaChannel 实例处理,因此您必须检查它是否与您的用例(关于事务大小等)。
【讨论】:
通道不必担心分区。因为通道是写它的,通道是消费消息的,所以不需要对消息进行分区。这就是flume-kafka-channel创建消息以进行写入的方式。
new KeyedMessage<String, byte[]>(topic.get(), null,
batchUUID, event)
但是,如果您的主题有多个分区,那么缺少密钥会导致消息被喷射到可用分区中。
如果您想更好地控制消息在分区中的分布方式,那么您可能需要研究 Kafka 的自定义分区器概念,因此您可以创建一个实现 org.apache.kafka.clients.producer.Partitioner 接口的类,并将 partitioner.class 属性设置为值相等为您的类命名,并确保您的自定义分区器在您的类路径中可用。这样,您可以在发布之前控制每条消息,并且可以决定消息应该转到哪个分区。您可以在您的水槽通道配置中设置属性 kafka.partitioner.class 以便它被拾取
【讨论】:
key 参数始终为null,所以消息总是被喷射到可用分区中。对吗?
kafka.partitioner.class 技巧,因为它是解决方案的一部分。
kafka.producer.Partitioner。这是不幸的,因为方法public int partition(Object key, int a_numPartitions) 不接收消息作为参数,只接收密钥。所以我想唯一的解决方案是修改 Kafka 通道源以将键 FlumeEvent 标头映射到 Kafka 中的消息键,就像对 Kafka Sink 所做的那样。