修复一个表创建失败的问题 https://github.com/AriaLyy/Aria/issues/570

修复一个非分块模式下导致下载失败的问题 https://github.com/AriaLyy/Aria/issues/571
修复一个服务器端无法创建socket连接,却没有返回码导致客户端卡住的问题 https://github.com/AriaLyy/Aria/issues/569
修复文件删除后,组合任务没有重新下载的问题 https://github.com/AriaLyy/Aria/issues/574
优化缓存队列和执行队列
This commit is contained in:
laoyuyu
2019-12-22 11:30:29 +08:00
parent 27c889e171
commit 8dc16ce302
38 changed files with 563 additions and 718 deletions

View File

@@ -95,8 +95,7 @@ public abstract class AbsTaskQueue<TASK extends AbsTask, TASK_WRAPPER extends Ab
* 停止所有任务
*/
@Override public void stopAllTask() {
for (String key : mExecutePool.getAllTask().keySet()) {
TASK task = mExecutePool.getAllTask().get(key);
for (TASK task : mExecutePool.getAllTask()) {
if (task != null) {
int state = task.getState();
if (task.isRunning() || (state != IEntity.STATE_COMPLETE
@@ -106,8 +105,13 @@ public abstract class AbsTaskQueue<TASK extends AbsTask, TASK_WRAPPER extends Ab
}
}
for (String key : mCachePool.getAllTask().keySet()) {
TASK task = mCachePool.getAllTask().get(key);
//for (String key : mCachePool.getAllTask().keySet()) {
// TASK task = mCachePool.getAllTask().get(key);
// if (task != null) {
// task.stop(TaskSchedulerType.TYPE_STOP_NOT_NEXT);
// }
//}
for (TASK task : mCachePool.getAllTask()) {
if (task != null) {
task.stop(TaskSchedulerType.TYPE_STOP_NOT_NEXT);
}

View File

@@ -26,7 +26,7 @@ import com.arialyy.aria.core.scheduler.TaskSchedulers;
import com.arialyy.aria.core.task.DownloadTask;
import com.arialyy.aria.util.ALog;
import java.util.LinkedHashSet;
import java.util.Map;
import java.util.List;
import java.util.Set;
/**
@@ -68,11 +68,10 @@ public class DTaskQueue extends AbsTaskQueue<DownloadTask, DTaskWrapper> {
*/
public void setTaskHighestPriority(DownloadTask task) {
task.setHighestPriority(true);
Map<String, DownloadTask> exeTasks = mExecutePool.getAllTask();
//Map<String, DownloadTask> exeTasks = mExecutePool.getAllTask();
List<DownloadTask> exeTasks = mExecutePool.getAllTask();
if (exeTasks != null && !exeTasks.isEmpty()) {
Set<String> keys = exeTasks.keySet();
for (String key : keys) {
DownloadTask temp = exeTasks.get(key);
for (DownloadTask temp : exeTasks) {
if (temp != null && temp.isRunning() && temp.isHighestPriorityTask() && !temp.getKey()
.equals(task.getKey())) {
ALog.e(TAG, "设置最高优先级任务失败,失败原因【任务中已经有最高优先级任务,请等待上一个最高优先级任务完成,或手动暂停该任务】");
@@ -80,28 +79,28 @@ public class DTaskQueue extends AbsTaskQueue<DownloadTask, DTaskWrapper> {
return;
}
}
}
int maxSize = AriaConfig.getInstance().getDConfig().getMaxTaskNum();
int currentSize = mExecutePool.size();
if (currentSize == 0 || currentSize < maxSize) {
startTask(task);
} else {
Set<DownloadTask> tempTasks = new LinkedHashSet<>();
for (int i = 0; i < maxSize; i++) {
DownloadTask oldTsk = mExecutePool.pollTask();
if (oldTsk != null && oldTsk.isRunning()) {
if (i == maxSize - 1) {
oldTsk.stop(TaskSchedulerType.TYPE_STOP_AND_WAIT);
mCachePool.putTaskToFirst(oldTsk);
break;
int maxSize = AriaConfig.getInstance().getDConfig().getMaxTaskNum();
int currentSize = mExecutePool.size();
if (currentSize == 0 || currentSize < maxSize) {
startTask(task);
} else {
Set<DownloadTask> tempTasks = new LinkedHashSet<>();
for (int i = 0; i < maxSize; i++) {
DownloadTask oldTsk = mExecutePool.pollTask();
if (oldTsk != null && oldTsk.isRunning()) {
if (i == maxSize - 1) {
oldTsk.stop(TaskSchedulerType.TYPE_STOP_AND_WAIT);
mCachePool.putTaskToFirst(oldTsk);
break;
}
tempTasks.add(oldTsk);
}
tempTasks.add(oldTsk);
}
}
startTask(task);
startTask(task);
for (DownloadTask temp : tempTasks) {
mExecutePool.putTask(temp);
for (DownloadTask temp : tempTasks) {
mExecutePool.putTask(temp);
}
}
}
}

View File

@@ -20,65 +20,43 @@ import android.text.TextUtils;
import com.arialyy.aria.core.task.AbsTask;
import com.arialyy.aria.util.ALog;
import com.arialyy.aria.util.CommonUtil;
import java.util.LinkedHashSet;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.TimeUnit;
import java.util.ArrayList;
import java.util.Deque;
import java.util.List;
import java.util.concurrent.LinkedBlockingDeque;
/**
* Created by lyy on 2016/8/14. 任务缓存池,所有下载任务最先缓存在这个池中
*/
public class BaseCachePool<TASK extends AbsTask> implements IPool<TASK> {
private static final String TAG = "BaseCachePool";
private final String TAG = CommonUtil.getClassName(this);
private static final int MAX_NUM = Integer.MAX_VALUE; //最大下载任务数
private static final long TIME_OUT = 1000;
private static final Object LOCK = new Object();
private Map<String, TASK> mCacheMap;
private LinkedBlockingQueue<TASK> mCacheQueue;
private Deque<TASK> mCacheQueue;
BaseCachePool() {
mCacheQueue = new LinkedBlockingQueue<>(MAX_NUM);
mCacheMap = new ConcurrentHashMap<>();
mCacheQueue = new LinkedBlockingDeque<>(MAX_NUM);
}
/**
* 获取被缓存的任务
*/
public Map<String, TASK> getAllTask() {
return mCacheMap;
public List<TASK> getAllTask() {
return new ArrayList<>(mCacheQueue);
}
/**
* 清除所有缓存的任务
*/
public void clear() {
for (String key : mCacheMap.keySet()) {
TASK task = mCacheMap.get(key);
mCacheQueue.remove(task);
mCacheMap.remove(key);
}
mCacheQueue.clear();
}
/**
* 将任务放在队首
*/
public boolean putTaskToFirst(TASK task) {
if (mCacheQueue.isEmpty()) {
return putTask(task);
} else {
Set<TASK> temps = new LinkedHashSet<>();
temps.add(task);
for (int i = 0, len = size(); i < len; i++) {
TASK temp = pollTask();
temps.add(temp);
}
for (TASK t : temps) {
putTask(t);
}
return true;
}
return mCacheQueue.offerFirst(task);
}
@Override public boolean putTask(TASK task) {
@@ -87,16 +65,12 @@ public class BaseCachePool<TASK extends AbsTask> implements IPool<TASK> {
ALog.e(TAG, "任务不能为空!!");
return false;
}
String key = task.getKey();
if (mCacheQueue.contains(task)) {
ALog.w(TAG, "任务【" + task.getTaskName() + "】进入缓存队列失败,原因:已经在缓存队列中");
return false;
} else {
boolean s = mCacheQueue.offer(task);
ALog.d(TAG, "任务【" + task.getTaskName() + "】进入缓存队列" + (s ? "成功" : "失败"));
if (s) {
mCacheMap.put(CommonUtil.keyToHashKey(key), task);
}
return s;
}
}
@@ -104,19 +78,8 @@ public class BaseCachePool<TASK extends AbsTask> implements IPool<TASK> {
@Override public TASK pollTask() {
synchronized (LOCK) {
try {
TASK task;
task = mCacheQueue.poll(TIME_OUT, TimeUnit.MICROSECONDS);
if (task != null) {
String url = task.getKey();
mCacheMap.remove(CommonUtil.keyToHashKey(url));
}
return task;
} catch (InterruptedException e) {
e.printStackTrace();
}
return mCacheQueue.pollFirst();
}
return null;
}
@Override public TASK getTask(String key) {
@@ -125,12 +88,17 @@ public class BaseCachePool<TASK extends AbsTask> implements IPool<TASK> {
ALog.e(TAG, "key 为null");
return null;
}
return mCacheMap.get(CommonUtil.keyToHashKey(key));
for (TASK task : mCacheQueue) {
if (task.getKey().equals(key)) {
return task;
}
}
}
return null;
}
@Override public boolean taskExits(String key) {
return mCacheMap.containsKey(CommonUtil.keyToHashKey(key));
return getTask(key) != null;
}
@Override public boolean removeTask(TASK task) {
@@ -139,8 +107,6 @@ public class BaseCachePool<TASK extends AbsTask> implements IPool<TASK> {
ALog.e(TAG, "任务不能为空");
return false;
} else {
String key = CommonUtil.keyToHashKey(task.getKey());
mCacheMap.remove(key);
return mCacheQueue.remove(task);
}
}
@@ -152,10 +118,7 @@ public class BaseCachePool<TASK extends AbsTask> implements IPool<TASK> {
ALog.e(TAG, "请传入有效的下载链接");
return false;
}
String temp = CommonUtil.keyToHashKey(key);
TASK task = mCacheMap.get(temp);
mCacheMap.remove(temp);
return mCacheQueue.remove(task);
return mCacheQueue.remove(getTask(key));
}
}

View File

@@ -21,26 +21,23 @@ import com.arialyy.aria.core.AriaManager;
import com.arialyy.aria.core.task.AbsTask;
import com.arialyy.aria.util.ALog;
import com.arialyy.aria.util.CommonUtil;
import java.util.Map;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.ArrayList;
import java.util.Deque;
import java.util.List;
import java.util.concurrent.LinkedBlockingDeque;
/**
* Created by lyy on 2016/8/15. 任务执行池所有当前下载任务都该任务池中默认下载大小为2
*/
public class BaseExecutePool<TASK extends AbsTask> implements IPool<TASK> {
private final String TAG = "BaseExecutePool";
private final String TAG = CommonUtil.getClassName(this);
private static final Object LOCK = new Object();
final long TIME_OUT = 1000;
ArrayBlockingQueue<TASK> mExecuteQueue;
Map<String, TASK> mExecuteMap;
Deque<TASK> mExecuteQueue;
int mSize;
BaseExecutePool() {
mSize = getMaxSize();
mExecuteQueue = new ArrayBlockingQueue<>(mSize);
mExecuteMap = new ConcurrentHashMap<>();
mExecuteQueue = new LinkedBlockingDeque<>(mSize);
}
/**
@@ -55,8 +52,8 @@ public class BaseExecutePool<TASK extends AbsTask> implements IPool<TASK> {
/**
* 获取所有正在执行的任务
*/
public Map<String, TASK> getAllTask() {
return mExecuteMap;
public List<TASK> getAllTask() {
return new ArrayList<>(mExecuteQueue);
}
@Override public boolean putTask(TASK task) {
@@ -88,17 +85,13 @@ public class BaseExecutePool<TASK extends AbsTask> implements IPool<TASK> {
*/
public void setMaxNum(int maxNum) {
synchronized (LOCK) {
try {
ArrayBlockingQueue<TASK> temp = new ArrayBlockingQueue<>(maxNum);
TASK task;
while ((task = mExecuteQueue.poll(TIME_OUT, TimeUnit.MICROSECONDS)) != null) {
temp.offer(task);
}
mExecuteQueue = temp;
mSize = maxNum;
} catch (InterruptedException e) {
e.printStackTrace();
Deque<TASK> temp = new LinkedBlockingDeque<>(maxNum);
TASK task;
while ((task = mExecuteQueue.poll()) != null) {
temp.offer(task);
}
mExecuteQueue = temp;
mSize = maxNum;
}
}
@@ -109,12 +102,8 @@ public class BaseExecutePool<TASK extends AbsTask> implements IPool<TASK> {
*/
boolean putNewTask(TASK newTask) {
synchronized (LOCK) {
String url = newTask.getKey();
boolean s = mExecuteQueue.offer(newTask);
ALog.d(TAG, "任务【" + newTask.getTaskName() + "】进入执行队列" + (s ? "成功" : "失败"));
if (s) {
mExecuteMap.put(CommonUtil.keyToHashKey(url), newTask);
}
return s;
}
}
@@ -124,37 +113,19 @@ public class BaseExecutePool<TASK extends AbsTask> implements IPool<TASK> {
*/
boolean pollFirstTask() {
synchronized (LOCK) {
try {
TASK oldTask = mExecuteQueue.poll(TIME_OUT, TimeUnit.MICROSECONDS);
if (oldTask == null) {
ALog.w(TAG, "移除任务失败原因任务为null");
return false;
}
oldTask.stop();
String key = CommonUtil.keyToHashKey(oldTask.getKey());
mExecuteMap.remove(key);
} catch (InterruptedException e) {
e.printStackTrace();
TASK oldTask = mExecuteQueue.pollFirst();
if (oldTask == null) {
ALog.w(TAG, "移除任务失败原因任务为null");
return false;
}
oldTask.stop();
return true;
}
}
@Override public TASK pollTask() {
synchronized (LOCK) {
try {
TASK task;
task = mExecuteQueue.poll(TIME_OUT, TimeUnit.MICROSECONDS);
if (task != null) {
String url = task.getKey();
mExecuteMap.remove(CommonUtil.keyToHashKey(url));
}
return task;
} catch (InterruptedException e) {
e.printStackTrace();
}
return null;
return mExecuteQueue.poll();
}
}
@@ -164,12 +135,18 @@ public class BaseExecutePool<TASK extends AbsTask> implements IPool<TASK> {
ALog.e(TAG, "key为null");
return null;
}
return mExecuteMap.get(CommonUtil.keyToHashKey(key));
for (TASK task : mExecuteQueue) {
if (task.getKey().equals(key)) {
return task;
}
}
return null;
}
}
@Override public boolean taskExits(String key) {
return mExecuteMap.containsKey(CommonUtil.keyToHashKey(key));
return getTask(key) != null;
}
@Override public boolean removeTask(TASK task) {
@@ -189,16 +166,8 @@ public class BaseExecutePool<TASK extends AbsTask> implements IPool<TASK> {
ALog.e(TAG, "key 为null");
return false;
}
String convertKey = CommonUtil.keyToHashKey(key);
TASK task = mExecuteMap.get(convertKey);
final int oldQueueSize = mExecuteQueue.size();
boolean isSuccess = mExecuteQueue.remove(task);
final int newQueueSize = mExecuteQueue.size();
if (isSuccess && newQueueSize != oldQueueSize) {
mExecuteMap.remove(convertKey);
return true;
}
return false;
return mExecuteQueue.remove(getTask(key));
}
}

View File

@@ -18,13 +18,14 @@ package com.arialyy.aria.core.queue.pool;
import com.arialyy.aria.core.AriaConfig;
import com.arialyy.aria.core.AriaManager;
import com.arialyy.aria.core.task.AbsTask;
import com.arialyy.aria.util.CommonUtil;
/**
* Created by AriaL on 2017/6/29.
* 单个下载任务的执行池
*/
class DGLoadExecutePool<TASK extends AbsTask> extends DLoadExecutePool<TASK> {
private final String TAG = "DGLoadExecutePool";
private final String TAG = CommonUtil.getClassName(this);
@Override protected int getMaxSize() {
return AriaConfig.getInstance().getDGConfig().getMaxTaskNum();

View File

@@ -18,9 +18,6 @@ package com.arialyy.aria.core.queue.pool;
import com.arialyy.aria.core.AriaConfig;
import com.arialyy.aria.core.task.AbsTask;
import com.arialyy.aria.util.ALog;
import com.arialyy.aria.util.CommonUtil;
import java.util.Set;
import java.util.concurrent.TimeUnit;
/**
* Created by AriaL on 2017/6/29.
@@ -45,9 +42,10 @@ class DLoadExecutePool<TASK extends AbsTask> extends BaseExecutePool<TASK> {
return false;
} else {
if (mExecuteQueue.size() >= mSize) {
Set<String> keys = mExecuteMap.keySet();
for (String key : keys) {
if (mExecuteMap.get(key).isHighestPriorityTask()) return false;
for (TASK temp : mExecuteQueue) {
if (temp.isHighestPriorityTask()) {
return false;
}
}
if (pollFirstTask()) {
return putNewTask(task);
@@ -61,22 +59,15 @@ class DLoadExecutePool<TASK extends AbsTask> extends BaseExecutePool<TASK> {
}
@Override boolean pollFirstTask() {
try {
TASK oldTask = mExecuteQueue.poll(TIME_OUT, TimeUnit.MICROSECONDS);
if (oldTask == null) {
ALog.w(TAG, "移除任务失败错误原因任务为null");
return false;
}
if (oldTask.isHighestPriorityTask()) {
return false;
}
oldTask.stop();
String key = CommonUtil.keyToHashKey(oldTask.getKey());
mExecuteMap.remove(key);
} catch (InterruptedException e) {
e.printStackTrace();
TASK oldTask = mExecuteQueue.pollFirst();
if (oldTask == null) {
ALog.w(TAG, "移除任务失败错误原因任务为null");
return false;
}
if (oldTask.isHighestPriorityTask()) {
return false;
}
oldTask.stop();
return true;
}
}