【问题标题】:implementing PriorityQueue on ThreadPoolExecutor在 ThreadPoolExecutor 中实现优先队列
【发布时间】:2015-06-01 13:25:54
【问题描述】:

已经为此苦苦挣扎了 2 多天。

实现了我在这里看到的答案 Specify task order execution in Java

public class PriorityExecutor extends ThreadPoolExecutor {

public PriorityExecutor(int corePoolSize, int maximumPoolSize,
                        long keepAliveTime, TimeUnit unit, BlockingQueue<Runnable> workQueue) {
    super(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue);
}
//Utitlity method to create thread pool easily
public static ExecutorService newFixedThreadPool(int nThreads) {
    return new PriorityExecutor(nThreads, nThreads, 0L,
            TimeUnit.MILLISECONDS, new PriorityBlockingQueue<Runnable>());
}
//Submit with New comparable task
public Future<?> submit(Runnable task, int priority) {
    return super.submit(new ComparableFutureTask(task, null, priority));
}
//execute with New comparable task
public void execute(Runnable command, int priority) {
    super.execute(new ComparableFutureTask(command, null, priority));
}
}

public class ComparableFutureTask<T> extends FutureTask<T>
    implements
    Comparable<ComparableFutureTask<T>> {

volatile int priority = 0;

public ComparableFutureTask(Runnable runnable, T result, int priority) {
    super(runnable, result);
    this.priority = priority;
}
public ComparableFutureTask(Callable<T> callable, int priority) {
    super(callable);
    this.priority = priority;
}

@Override
public int compareTo(ComparableFutureTask<T> o) {
    return Integer.valueOf(priority).compareTo(o.priority);
}
}

我使用的 Runnable:MyTask

public class MyTask implements Runnable{

 public MyTask(File file, Context context, int requestId) {
    this._file = file;
    this.context = context;
    this.requestId = requestId;
}

@Override
public void run() {
      // some work
    } catch (IOException e) {
        Log.e("Callable try", post.toString());

    }
}

我的服务:MediaDownloadService

public class MediaDownloadService extends Service {

private DBHelper helper;
Notification notification;
HashMap<Integer,Future> futureTasks = new HashMap<Integer, Future>();
final int _notificationId=1;
File file;

@Override
public IBinder onBind(Intent intent) {
    return sharonsBinder;
}


@Override
public int onStartCommand(Intent intent, int flags, int startId) {
    helper = new DBHelper(getApplicationContext());
    PriorityExecutor executor = (PriorityExecutor) PriorityExecutor.newFixedThreadPool(3);
    Log.e("requestsExists", helper.requestsExists() + "");
   if(helper.requestsExists()){
        // map of the index of the request and the string of the absolute path of the request
        Map<Integer,String> requestMap = helper.getRequestsToExcute(0);
        Set<Integer> keySet = requestMap.keySet();
        Iterator<Integer> iterator = keySet.iterator();
        Log.e("MAP",requestMap.toString());
        //checks if the DB requests exists
        if(!requestMap.isEmpty()){
            //execute them and delete the DB entry
            while(iterator.hasNext()){
                int iteratorNext = iterator.next();
                Log.e("ITREATOR", iteratorNext + "");
                file = new File(requestMap.get(iteratorNext));
                Log.e("file", file.toString());
                Log.e("thread Opened", "Thread" + iteratorNext);
                Future future = executor.submit(new MyTask(file, this, iteratorNext),10);
                futureTasks.put(iteratorNext, future);
                helper.requestTaken(iteratorNext);
            }
            Log.e("The priority queue",executor.getQueue().toString());
        }else{

            Log.e("stopself", "stop self after this");
            this.stopSelf();
        }
    }
    return START_STICKY;
}

在这一行不断出现错误: 未来的未来 = executor.submit(new MyTask(file, this, iteratorNext),10);

即使是 executor.submit();假设返回一个我不断得到的未来对象

Caused by: java.lang.ClassCastException: java.util.concurrent.FutureTask cannot be cast to java.lang.Comparable
        at java.util.concurrent.PriorityBlockingQueue.siftUpComparable(PriorityBlockingQueue.java:318)
        at java.util.concurrent.PriorityBlockingQueue.offer(PriorityBlockingQueue.java:450)
        at java.util.concurrent.ThreadPoolExecutor.execute(ThreadPoolExecutor.java:1331)
        at java.util.concurrent.AbstractExecutorService.submit(AbstractExecutorService.java:81)
        at com.vit.infibond.test.PriorityExecutor.submit(PriorityExecutor.java:26)
        at com.vit.infibond.test.MediaDownloadService.onStartCommand(MediaDownloadService.java:65)

谁能把我从噩梦中拯救出来?

我也尝试按照此答案的建议进行操作 Testing PriorityBlockingQueue in ThreadPoolExecutor

通过添加 forNewTask 覆盖只是为了再次获得强制转换执行,但这次是针对 RunnableFuture。

我了解我的理解中缺少一些基本内容,并希望得到深入的解释......

【问题讨论】:

  • 谢谢莎伦,我在下面扩展了你的 cmets。相当激烈的话题

标签: java priority-queue future threadpoolexecutor


【解决方案1】:

通过查看java.util.concurrent.ThreadPoolExecutor 的源代码,在提交期货时让它工作似乎真的很痛苦。您必须覆盖感觉是内部的受保护方法并进行一些讨厌的强制转换。

我建议您直接使用execute 方法。 Runnable 没有包装在那里,所以它应该可以工作。

如果您需要等待工作的结果,我建议您自己实施,以避免与ThreadPoolExecutor 内部结构混淆。

【讨论】:

  • 我的问题不是等待我的工作完成,而是我想管理它们,因为它们在队列中。例如,我在队列中有项目,我现在需要改进他们的优先级之一,所以我需要引用该未来对象,以便我可以取消该任务并以更高的优先级重新插入它。如果你能解释一下我自己实现它的更多信息,你的意思是我自己创建一个线程池?
  • 要更改优先级,您需要访问ComparableFutureTask,因为这是队列中的内容,而不是Future。因此,为此,您仍然可以使用executor.execute() 而不是submit()。
  • executor.execute() 会返回一个 ComparableFutureTask 吗?如果不是,我将如何达到我想取消的任务..?我认为要更改优先级,我必须取消并重新插入任务,并且取消是通过未来对象完成的......或者是吗?
  • @sharongur 不,它没有,但是由于您已经实现了 PriorityExecutor,您可以创建一个自定义方法在执行时返回它。
  • 太棒了!绝对精彩!我更改了执行方法以返回一个 ComparableFutureTask 并添加了一个方法来增加包含 Comparable 对象的哈希图的 id 的优先级,然后只需更改其优先级非常感谢!
【解决方案2】:

sharon gur 在最底层的建议是改变

//execute with New comparable task
public void execute(Runnable command, int priority) {
    super.execute(new ComparableFutureTask(command, null, priority));
}

到

//execute with New comparable task
public ComparableFutureTask  execute(Runnable command, int priority) {
    ComparableFutureTask task = new ComparableFutureTask(command, null, priority);
    super.execute(task);
    return task;
}

然后在你的调用者中:

CurrentTask currentTask = new CurrentTask(priority,queue)
RunnableFuture task = enhancedExecutor.execute(currentTask,priority.value)
task?.get()

我有一个问题

RunnableFuture task = myExecutor.submit(currentTask)
task?.get()

导致currentTask 然后转换为FutureTask 并且不了解我在 CurrentTask 中的对象的问题。就像.execute 一样,一切都很好。这个 hack 似乎已经够半 / 接近足够的工作了。

所以它工作得很好,但没有生成文件

RunnableFuture task = myExecutor.execuute(currentTask)
    task?.get()

所以这就是我让它工作的方式(优先级被处理两次)感觉不对,但工作......

当前任务::

class CurrentTask implements Runnable {
    private Priority priority
    private MyQueue queue

    public int getPriority() {
        return priority.value
    }

    public CurrentTask(Priority priority,ReportsQueue queue){
        this.priority = priority
        this.queue=queue
    }

    @Override
    public void run() {
...
}
}

优先级:

public enum Priority {

    HIGHEST(0),
    HIGH(1),
    MEDIUM(2),
    LOW(3),
    LOWEST(4)

    int value

    Priority(int val) {
        this.value = val
    }

    public int getValue(){
        return value
    }
}

然后你的执行者调用

public YourExecutor() {

    public YourExecutor() {
        super(maxPoolSize,maxPoolSize,timeout,TimeUnit.SECONDS, new PriorityBlockingQueue<Runnable>(1000,new ReverseComparator()))
    }

因此,在更改为新方法之前,提交点击下面的比较器,并且作为 TaskExecutor 不会理解 .priority?.value ,默认情况下 .execute currentTask 是什么命中了这个并且一切正常

public int compare(final Runnable lhs, final Runnable rhs) {

    if(lhs instanceof Runnable && rhs instanceof Runnable){
      // Favour a higher priority
        println "${lhs} vs ${lhs.getClass()}"
      if(((Runnable)lhs)?.priority?.value<((Runnable)rhs)?.priority?.value){
 ...
}

}

因此,上面的 hack 和下面的更改似乎正在工作

class  ReverseComparator implements Comparator<ComparableFutureTask>{

  @Override
  public int compare(final ComparableFutureTask lhs, final ComparableFutureTask rhs) {

    if(lhs instanceof ComparableFutureTask && rhs instanceof ComparableFutureTask){

        // run higher priority (lower numbers before higher numbers)
        println "${lhs} vs ${lhs.getClass()} ::: ${lhs.priority}"
      if(((Runnable)lhs)?.priority<((Runnable)rhs)?.priority){
          println "-returning -1"
        return -1;
      } else if (((Runnable)lhs)?.priority>((Runnable)rhs)?.priority){
      println "-returning @@@1"
        return 1;
      } 


    }
    println "-returning ==0 "
    return 0;
  }  

仅仅是因为我们传入了覆盖ComparableFutureTask,它具有扩展FutureTask的优先级

希望现在转了一整天是有意义的:)

【讨论】:

猜你喜欢
  • 1970-01-01
  • 2011-11-10
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2012-11-26
  • 1970-01-01
相关资源
最近更新 更多