【问题标题】:Load Balance Queues for variable rate consumers可变速率消费者的负载平衡队列
【发布时间】:2018-11-06 05:17:19
【问题描述】:

我有一个生产者和消费者框架。每个生产者推送到一个队列,消费者从队列中消费。在任何时间点都可以有一个或多个队列,每个消费者从单个队列中消费。但是生产者可以生产到任何队列。如果消费者很慢,它会不断堆积消息。我正在尝试提供一个框架,我可以在其中对消费者进行负载平衡,以便所有消费者队列具有几乎相等的消息,而不管消费者的速度如何。

例子:

这里的队列 Q1-Q3 应该有几乎相等的消息,而与 C1-C3 消费者的比率无关。我现在使用的默认策略是针对生产者的循环,但如果任何消费者速度较慢,它会继续将消息添加到队列中。所有消息都属于同一类型,因此它会进入任何队列。

任何开始的建议都是有帮助的。

【问题讨论】:

    标签: java performance distributed-computing


    【解决方案1】:

    简单 - 添加到具有最低项目数的队列。

    【讨论】:

    • 我使用了相同的逻辑并使用计时器对其进行了一些即兴创作。请检查。
    • @AmarendraReddy:在 SO 上验证您的代码是个坏主意。不建议。读者可能会误解的东西太多了。
    【解决方案2】:

    以下是我实施的解决方案。使用的算法如下。

    1. 每 30 秒查找所有队列的平均值。
    2. 如果滞后 消费者 w.r.t 意味着大于特定阈值忽略该 队列/消费者。

    生产者代码:

    import java.util.ArrayList;
    import java.util.List;
    import java.util.Random;
    import java.util.concurrent.BlockingQueue;
    
    public class Producer implements Runnable{
    
        private List<BlockingQueue<Integer>> blockingQueues = new ArrayList<>();
        private List<Integer> fullPartitions;
        private List<Integer> activePartitions;
        long timer = System.currentTimeMillis();
        int THRESHOLD = 10000;
        int currentQueue = 0;
    
        public Producer(List<BlockingQueue<Integer>> blockingQueues, List<Integer> fullPartitions, List<Integer> activePartitions) {
            this.blockingQueues = blockingQueues;
            this.fullPartitions = fullPartitions;
            this.activePartitions = activePartitions;
        }
    
        @Override
        public void run() {
            long start = System.currentTimeMillis();
            while(true) {
                blockingQueues.get(getNextID()).offer(new Random().nextInt(100000));
                try {
                    if(System.currentTimeMillis()-start<300000)
                        Thread.sleep(1);
                    else
                        break;
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
            }
        }
    
        private int getNextID() {
            if(System.currentTimeMillis()-timer>30000) {
                activePartitions = new ArrayList<>();
                long mean = 0l; 
                for(int i=0;i<fullPartitions.size();i++) 
                    mean += blockingQueues.get(i).size();
    
                mean  = mean/blockingQueues.size();
                for(int i=0;i<fullPartitions.size();i++) 
                    if(blockingQueues.get(i).size()-mean<THRESHOLD)
                        activePartitions.add(i);
    
                timer = System.currentTimeMillis();
            }
            int partitionID = activePartitions.get(currentQueue%activePartitions.size());
            currentQueue++;
            return partitionID;
        }
    }
    

    消费者:

    import java.util.concurrent.ArrayBlockingQueue;
    import java.util.concurrent.BlockingQueue;
    
    public class Consumer implements Runnable{
    
        private BlockingQueue<Integer> blockingQueue = new ArrayBlockingQueue<>(100000000);
        private int delayFactor;
        public Consumer(BlockingQueue<Integer> blockingQueue, int delayFactor, int consumerNo) {
            this.blockingQueue = blockingQueue;
            this.delayFactor = delayFactor;
        }
    
        @Override
        public void run() {
            long start = System.currentTimeMillis();
            while(true) {
                try {
                    blockingQueue.take();
                    if(blockingQueue.isEmpty())
                        System.out.println((System.currentTimeMillis()-start)/1000);
                    Thread.sleep(delayFactor);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
            }
        }
    
    }
    

    主线程:

    import java.util.ArrayList;
    import java.util.List;
    import java.util.concurrent.ArrayBlockingQueue;
    import java.util.concurrent.BlockingQueue;
    
    public class KafkaLoadBalancer {
    
        private static int MAX_PARTITION = 4;
    
        public static void main(String args[]) throws InterruptedException {
            List<BlockingQueue<Integer>> blockingQueues = new ArrayList<>();
            List<Integer> fullPartitions = new ArrayList<Integer>();
            List<Integer> activePartitions = new ArrayList<Integer>();
    
            System.out.println("Creating Queues");
            for(int i=0;i<MAX_PARTITION;i++) {
                blockingQueues.add(new ArrayBlockingQueue<>(1000000));
                fullPartitions.add(i);
                activePartitions.add(i);
            }
    
            System.out.println("Starting Producers");
            for(int i=0;i<MAX_PARTITION;i++) {
                Producer producer = new Producer(blockingQueues,fullPartitions,activePartitions);
                new Thread(producer).start();
            }
    
            System.out.println("Starting Consumers");
            for(int i=0;i<MAX_PARTITION;i++) {
                Consumer consumer = new Consumer(blockingQueues.get(i),i+1,i);
                new Thread(consumer).start();
            }
    
            System.out.println("Starting Display Thread");
            DisplayQueue dq = new DisplayQueue(blockingQueues);
            new Thread(dq).start();
        }
    }
    

    DispayQueue : 显示队列大小

    import java.util.List;
    import java.util.concurrent.BlockingQueue;
    
    public class DisplayQueue implements Runnable {
    
        private List<BlockingQueue<Integer>> blockingQueues;
    
        public DisplayQueue(List<BlockingQueue<Integer>> blockingQueues) {
            this.blockingQueues = blockingQueues;
        }
    
        @Override
        public void run() {
    
            long start = System.currentTimeMillis();
            while(true) {
                if(System.currentTimeMillis()-start>30000) {
                    for(int i=0;i<blockingQueues.size();i++)
                        System.out.println("Queue "+i+" size is=="+blockingQueues.get(i).size());
                    start = System.currentTimeMillis();
                }
            }
    
        }
    
    }
    

    【讨论】:

      猜你喜欢
      • 2016-07-15
      • 1970-01-01
      • 2019-04-30
      • 2020-02-28
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2012-11-27
      相关资源
      最近更新 更多