【问题标题】:Does Kafka support priority for topic or message?Kafka 是否支持主题或消息的优先级?
【发布时间】:2015-08-19 17:52:29
【问题描述】:

我正在探索 Kafka 是否支持任何队列或消息处理的优先级。

它似乎不支持任何这样的事情。我用谷歌搜索,发现这个邮件档案也支持这个: http://mail-archives.apache.org/mod_mbox/incubator-kafka-users/201206.mbox/%3CCAOeJiJhVHsr=d6aSTihPsqWVg6vK5xYLam6yMDcd6UAUoXf-DQ@mail.gmail.com%3E

这里有没有人配置了 Kafka 来确定任何主题或消息的优先级?

【问题讨论】:

标签: apache-kafka


【解决方案1】:

我将在此处添加@Skyanswer 的Java 版本供大家参考。注意:我没有使用 KafkaStreams,而是使用普通的 KafkaConsumer 实现。

从生产者的角度,你可以根据优先级将消息发布到各自的主题。

从消费者的角度来看,您可以尝试实现以下内容。请注意,这不是生产就绪的实现。此解决方案是单线程的,可能会很慢。

与桶优先级模式不同,此代码将继续处理来自高优先级主题的消息,直到处理完所有消息。当高优先级主题没有消息时,将退回到下一个优先级,依此类推。

import lombok.*;
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.context.event.ApplicationReadyEvent;
import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Component;
import java.util.function.Consumer;

import javax.annotation.PostConstruct;
import java.time.Duration;
import java.util.*;

@Component
@Slf4j
public class PriorityBasedConsumer {

    @Value("${spring.kafka.consumer.group-id}")
    private String consumerGroupId;

    @Value("${spring.kafka.bootstrap-servers}")
    private String bootstrapServers;

    private final List<TopicConsumer> consumersInPriorityOrder = new ArrayList<>();
    
    @RequiredArgsConstructor(staticName = "of")
    @Getter
    private static class TopicConsumer {
        private final String topic;
        private final KafkaConsumer<String, String> kafkaConsumer;
        private final Consumer<ConsumerRecords<String, String>> consumerLogic;
    }

    private void highPriorityConsumer(ConsumerRecords<String, String> records) {
        // high priority processing...
    }

    private void mediumPriorityConsumer(ConsumerRecords<String, String> records) {
        // medium priority processing...
    }

    private void lowPriorityConsumer(ConsumerRecords<String, String> records) {
        // low priority processing...
    }

    @PostConstruct
    public void init() {
        Map<String, Consumer<ConsumerRecords<String, String>>> topicVsConsumerLogic = new HashMap<>();
        topicVsConsumerLogic.put("high_priority_queue", this::highPriorityConsumer);
        topicVsConsumerLogic.put("medium_priority_queue", this::mediumPriorityConsumer);
        topicVsConsumerLogic.put("low_priority_queue", this::lowPriorityConsumer);
        // if you're taking the topic names from external configuration, make sure to order it based on priority.
        for (String topic : Arrays.asList("high_priority_queue", "medium_priority_queue", "low_priority_queue")) {
            Properties consumerProperties = new Properties();
            consumerProperties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
            consumerProperties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
            consumerProperties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
            consumerProperties.put(ConsumerConfig.GROUP_ID_CONFIG, consumerGroupId);
            // add other properties.
            KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerProperties);
            consumer.subscribe(Collections.singletonList(topic));
            consumersInPriorityOrder.add(TopicConsumer.of(topic, consumer, topicVsConsumerLogic.get(topic)));
        }
    }

    @EventListener(ApplicationReadyEvent.class) // To execute once the application is ready.
    public void startConsumers() {
        // For illustration purposes, I just wrote this synchronous code. Use thread pools where ever 
        // necessary for high performance.
        while (true) { // poll infinitely
            try {
                // Consumers iterated based on priority.
                for (TopicConsumer topicConsumer : consumersInPriorityOrder) {
                    ConsumerRecords<String, String> records
                            = topicConsumer.getKafkaConsumer().poll(Duration.ofMillis(100));
                    if (!records.isEmpty()) {
                        topicConsumer.getConsumerLogic().accept(records);
                        break;  // To start consuming again based on priority.
                    }
                }
            } catch (Exception e) {
                // on any unknown runtime exceptions, ignoring here. You can add your proper logic.
                log.error("Unknown exception occurred.", e);
            }
        }
    }
}

【讨论】:

    【解决方案2】:

    Implementing Message Prioritization in Apache Kafka 上有来自 Confluent 的博客,其中描述了如何实现消息优先级。

    首先,重要的是要了解 Kafka 的设计不允许开箱即用的解决方案来确定消息的优先级。主要原因是:

    • 存储:Kafka 被设计为一个仅追加提交日志,其中包含反映现实事件及时发生的不可变消息。
    • 消费者:Kafka 主题中的消息可以同时被多个消费者消费。每个消费者可能有不同的优先级,这使得无法提前对主题内的消息进行排序。

    建议的解决方案是使用GitHub 上提供的Bucket Priority Pattern,可以通过其自述文件中的图表进行最佳描述。您可以通过自定义生产者的partitioner和消费者的分配策略来使用具有多个分区的单个主题,而不是为不同的优先级使用多个主题。

    根据messages key,生产者将消息写入正确的优先级桶中:

    另一方面,消费者组将自定义其分配策略,并优先从具有最高分区的分区中读取消息:

    在您的客户端代码(生产者和消费者)中,您需要启动并调整以下客户端配置。

    # Producer
    configs.setProperty(ProducerConfig.PARTITIONER_CLASS_CONFIG,
       BucketPriorityPartitioner.class.getName());
    configs.setProperty(BucketPriorityConfig.BUCKETS_CONFIG, "Platinum, Gold");
    configs.setProperty(BucketPriorityConfig.ALLOCATION_CONFIG, "70%, 30%");
    
    # Consumer
    configs.setProperty(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG,
       BucketPriorityAssignor.class.getName());
    configs.setProperty(BucketPriorityConfig.BUCKETS_CONFIG, "Platinum, Gold");
    configs.setProperty(BucketPriorityConfig.ALLOCATION_CONFIG, "70%, 30%");
    
    

    【讨论】:

      【解决方案3】:

      解决方案是根据优先级创建 3 个不同的主题。

      • 高优先级主题
      • 中等优先主题
      • 低优先级主题

      作为一般经验法则,高优先级主题的消费者数量 > 中优先级主题的消费者数量 > 低优先级主题的消费者数量

      这样可以保证到达高优先级主题的消息将比低优先级主题更快地处理。

      【讨论】:

        【解决方案4】:

        Kafka 是一种快速、可扩展、分布式的设计、分区和复制的提交日志服务。因此没有主题或消息的优先级。

        我也遇到了和你一样的问题。解决方法很简单。在 kafka 队列中创建主题,比方说:

        1. high_priority_queue

        2. medium_priority_queue

        3. 低优先级队列

        在 high_priority_queue 中发布高优先级消息,在 medium_priority_queue 中发布中优先级消息。

        现在您可以创建 kafka 消费者并为所有主题打开流。

          // this is scala code 
          val props = new Properties()
          props.put("group.id", groupId)
          props.put("zookeeper.connect", zookeeperConnect)
          val config = new ConsumerConfig(props)
          val connector = Consumer.create(config)
          val topicWithStreamCount = Map(
               "high_priority_queue" -> 1,
               "medium_priority_queue" ->  1, 
               "low_priority_queue" -> 1
          )
          val streamsMap = connector.createMessageStreams(topicWithStreamCount)
        

        你得到每个主题的流。现在你可以先阅读 high_priority 主题,如果主题没有任何消息,然后回退到 medium_priority_queue 主题。如果 medium_priority_queue 为空,则读取 low_priority 队列。

        这个技巧对我来说很好用。可能对你有帮助!

        【讨论】:

        • 嘿天空,感谢您的回答。我在消费多个主题时面临一个问题。我正在使用 ConsumerIterator 来使用流。但是,一旦我创建了 ConsumerIterator,我的代码就会被阻止在这里。你能解释一下你是如何消费这两个流的吗?我正在用 java 编写消费者和发布者
        • @aviundefined 可以使用线程池进行并行消费。看看:cwiki.apache.org/confluence/display/KAFKA/…
        • 它看起来像旧的消费者 API - 有没有使用新消费者 API 的推荐方式?我注意到方法 pause 和 resume,但不知道如何找出何时是调用 pause 的正确时间 - 更具体地说,如何找出优先级更高的主题中有新消息?
        • @Sky - 如果您必须将 low_priority_queue 中现有消息的优先级设置为 high_priority_queue,它可能不起作用。
        • 它将继续处理高优先级主题,并且永远不会转到低优先级主题。这就是我们的意图@beinghuman
        【解决方案5】:

        您可以结帐priority-kafka-client 以获得主题的优先消费。

        基本思路如下(复制/粘贴部分README):

        在此上下文中,优先级是一个正整数 (N),优先级为 0 &lt; 1 &lt; ... &lt; N-1

        PriorityKafkaProducer (implements org.apache.kafka.clients.producer.Producer):

        实现采用额外的优先级参数Future&lt;RecordMetadata&gt; send(int priority, ProducerRecord&lt;K, V&gt; record)。这表明在该优先级上产生记录。 Future&lt;RecordMetadata&gt; send(int priority, ProducerRecord&lt;K, V&gt; record) 默认在最低优先级 0 上记录生产。对于每个逻辑主题 XYZ - 优先级 0 XYZ-i

        CapacityBurstPriorityKafkaConsumer (implements org.apache.kafka.clients.consumer.Consumer):

        实现为每个优先级 0 XYZ-i 987654329@。这与 PriorityKafkaProducer 协同工作。

        max.poll.records 属性根据maxPollRecordsDistributor 划分为优先主题消费者 - 默认为ExpMaxPollRecordsDistributor。其余的 KafkaConsumer 配置按原样传递给每个优先主题消费者。定义 max.partition.fetch.bytesfetch.max.bytesmax.poll.interval.ms 时必须小心,因为这些值将按原样用于所有优先主题消费者。

        致力于将max.poll.records 属性分配给每个优先主题消费者作为他们的保留容量。记录是从所有优先级主题消费者中按顺序获取的,这些消费者配置了分布式max.poll.records 值。分配必须为更高的优先级保留更高的容量或处理速率。

        注意 1 - 如果我们在优先级主题中存在倾斜分区,例如优先级 2 分区中的 10K 条记录,优先级 1 分区中的 100 条记录,优先级 0 分区中的 10 条记录分配给不同的消费者线程,那么实现将不会在这些消费者之间同步以调节容量,因此将无法遵守优先级.所以生产者必须确保没有倾斜的分区(例如使用循环 - 这“可能”意味着没有消息排序假设,消费者可以通过分离获取和处理关注点来选择并行处理记录)。

        注意 2 - 如果我们在优先级主题中有空分区,例如分配的优先级 2 和 1 分区中没有待处理记录,优先级 0 分区中的 10K 记录分配给同一个消费者线程,那么我们希望优先级 0 主题分区消费者将其容量突增到max.poll.records 并且不将自身限制为保留容量基于maxPollRecordsDistributor 否则整体容量将被充分利用。

        此实现将尝试解决上述注意事项。每个消费者对象都有单独的优先级主题消费者,每个优先级消费者都具有基于 maxPollRecordsDistributor 的预留容量。如果满足以下所有条件,每个优先级主题消费者将尝试突入群组中其他优先级主题消费者的容量:

        • 它有资格爆发 - 这是如果在最后一次 max.poll.history.window.size 尝试 poll() 至少 min.poll.window.maxout.threshold 次时它收到的记录数等于分配的 max.poll.records maxPollRecordsDistributor。这表明该分区有更多的传入记录需要处理。

        • 更高优先级不符合突发条件 - 根据上述逻辑,没有更高优先级的主题消费者符合突发条件。基本上让位于更高的优先级。

        如果上述情况属实,那么优先级主题消费者将爆发到所有其他优先级主题消费者容量。每个优先级主题消费者的突发量等于poll() 的最后一次max.poll.history.window.size 尝试中的最少未使用容量。

        【讨论】:

          【解决方案6】:

          您需要有一个单独的主题并根据它们的优先级进行流式传输

          【讨论】:

            猜你喜欢
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            • 2015-02-11
            • 1970-01-01
            • 1970-01-01
            • 2016-06-02
            • 1970-01-01
            • 1970-01-01
            相关资源
            最近更新 更多