【问题标题】:Flume: Routing events to the proper topic partition with Kafka channelFlume:使用 Kafka 通道将事件路由到正确的主题分区
【发布时间】:2016-05-11 11:16:39
【问题描述】:

在 Flume 中,当使用 Kafka 通道时,有没有办法影响事件发送到哪个分区?

对于 Kafka sinkkey FlumeEvent 标头显然用于选择分区,但我找不到任何有关 Kafka channel 分区的文档。 p>

【问题讨论】:

    标签: apache-kafka flume flume-ng


    【解决方案1】:

    Flume 的 Kafka 通道不支持像 KafkaSink 那样将 Event 标头映射到开箱即用的分区键。

    但是,对其进行修改以使其能够做到这一点并不太复杂。由于我不确定我是否可以分享代码,我将给出指示:

    1. 为将映射到分区键的标头名称添加配置键
    2. 在内部类 KafkaTransaction 中,将成员 serializedEvents 类型中的 byte[] 替换为也可以为每个事件(内部类,甚至是 Kafka KeyedMessage<String, byte[]>)持有一个 String 键的东西
    3. 在方法KafkaTransaction.doPut(Event event)中,从headers中获取key并与序列化消息一起存储在serializedEvents
    4. 在方法KafkaTransaction.doCommit()中,使用与序列化事件一起存储的密钥而不是batchUUID

    注意事务中的事件将不再保证由通道消费者端的单个 KafkaChannel 实例处理,因此您必须检查它是否与您的用例(关于事务大小等)。

    【讨论】:

      【解决方案2】:

      通道不必担心分区。因为通道是写它的,通道是消费消息的,所以不需要对消息进行分区。这就是flume-kafka-channel创建消息以进行写入的方式。

      new KeyedMessage<String, byte[]>(topic.get(), null,
                    batchUUID, event)
      

      但是,如果您的主题有多个分区,那么缺少密钥会导致消息被喷射到可用分区中。

      如果您想更好地控制消息在分区中的分布方式,那么您可能需要研究 Kafka 的自定义分区器概念,因此您可以创建一个实现 org.apache.kafka.clients.producer.Partitioner 接口的类,并将 partitioner.class 属性设置为值相等为您的类命名,并确保您的自定义分区器在您的类路径中可用。这样,您可以在发布之前控制每条消息,并且可以决定消息应该转到哪个分区。您可以在您的水槽通道配置中设置属性 kafka.partitioner.class 以便它被拾取

      【讨论】:

      • 嗯,我确实在 Flume 之外还有其他消费者群体来讨论相关主题,所以我关心分区。我希望做的是将标头(原始主机)映射到密钥。我对代码的理解是 KeyedMessage<...> 构造函数的key 参数始终为null,所以消息总是被喷射到可用分区中。对吗?
      • 是的,消息被喷了。您可能想尝试的一种解决方案是 Kafka 具有自定义分区器的概念,因此您可以创建一个类实现接口,并将 partitioner.class 属性设置为与您的类名称相等的值,并确保您的自定义分区器在您的类路径中可用.这样,您可以在发布之前控制每条消息,并且可以确定消息应该转到哪个分区。您可以在您的水槽通道配置中设置属性 kafka.partitioner.class 以便它被拾取
      • 谢谢,会考虑您的建议。您应该编辑您的答案以包含 kafka.partitioner.class 技巧,因为它是解决方案的一部分。
      • 我猜应该实现的接口是kafka.producer.Partitioner。这是不幸的,因为方法public int partition(Object key, int a_numPartitions) 不接收消息作为参数,只接收密钥。所以我想唯一的解决方案是修改 Kafka 通道源以将键 FlumeEvent 标头映射到 Kafka 中的消息键,就像对 Kafka Sink 所做的那样。
      • 我最终分叉了 Flume 的 KafkaChannel 以添加对将 Event 标头映射到 Kafka 中的分区键的支持。我不确定我是否可以共享代码,但我的修改是微不足道的。但是,原始代码将内部事务 UID 映射到分区键,以便事务中的所有事件都由单个 KafkaChannel 实例的另一端使用。我对这个功能不感兴趣(在我的情况下,通道的另一端没有接收器,我只是使用 KafkaChannel 作为 Source->Channel->Kafka Sink 的快捷方式)但我不确定它有什么影响一般情况下可能有。
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-03-23
      • 2016-10-01
      • 2016-10-15
      相关资源
      最近更新 更多