使用PublishSubject与flatMap实现并发:Android服务动态任务处理问题
嘿,我明白你现在的困扰——要在Android服务里搞一个支持动态调整、并发控制的任务队列确实有点棘手,尤其是还要兼顾任务的新增和排序。我之前也帮不少开发者解决过类似问题,给你一套实用的方案试试:
核心思路:带优先级的阻塞队列+线程池组合实现
我们需要一个能动态添加、自动排序的线程安全任务容器,再配合线程池严格控制并发数。这里推荐用PriorityBlockingQueue作为任务队列,它天然支持优先级排序,而且能完美适配线程池的调度逻辑。
第一步:定义可排序的任务模型
先把每个ITEM抽象成任务类,实现Runnable用于执行任务,同时实现Comparable接口定义排序规则,让队列能自动帮我们调整任务顺序。
public class TransferTask implements Runnable, Comparable<TransferTask> { // 任务类型枚举 public enum TaskType { UPLOAD, DOWNLOAD } private String taskId; // 唯一标识任务 private TaskType type; // 上传/下载 private int priority; // 优先级(数值越大优先级越高) private String taskData; // 任务数据:比如文件路径、目标URL private volatile boolean isCancelled; // 任务取消标记 public TransferTask(String taskId, TaskType type, int priority, String taskData) { this.taskId = taskId; this.type = type; this.priority = priority; this.taskData = taskData; } @Override public void run() { // 先检查是否被取消 if (isCancelled) return; try { if (type == TaskType.UPLOAD) { doUpload(taskData); } else { doDownload(taskData); } } catch (Exception e) { // 这里可以添加任务失败后的处理逻辑:比如重试、标记状态 e.printStackTrace(); } } // 模拟上传逻辑 private void doUpload(String data) throws InterruptedException { System.out.println("开始上传:" + data); Thread.sleep(3000); // 模拟耗时操作 System.out.println("上传完成:" + data); } // 模拟下载逻辑 private void doDownload(String data) throws InterruptedException { System.out.println("开始下载:" + data); Thread.sleep(3000); System.out.println("下载完成:" + data); } // 实现排序规则:优先级高的先执行;优先级相同则按任务ID排序 @Override public int compareTo(TransferTask other) { // PriorityBlockingQueue是小顶堆,反转优先级让高优任务排在队首 int priorityCompare = Integer.compare(other.priority, this.priority); return priorityCompare != 0 ? priorityCompare : this.taskId.compareTo(other.taskId); } // 任务取消方法 public void cancelTask() { this.isCancelled = true; } // Getter方法 public String getTaskId() { return taskId; } public int getPriority() { return priority; } }
第二步:构建任务管理核心类
这个类负责维护任务队列、线程池,以及提供动态添加、重新排序、取消任务的方法,保证全局单例避免重复创建。
public class TaskManager { private static final int CONCURRENT_TASK_COUNT = 3; // 并发任务数,可根据设备性能调整 private ExecutorService taskExecutor; private final Map<String, TransferTask> allTasks = new ConcurrentHashMap<>(); // 存储所有任务 // 单例模式 private TaskManager() { // 初始化线程池:核心/最大线程数=并发数,用优先级队列作为任务容器 taskExecutor = new ThreadPoolExecutor( CONCURRENT_TASK_COUNT, CONCURRENT_TASK_COUNT, 0L, TimeUnit.MILLISECONDS, new PriorityBlockingQueue<>() ); } public static TaskManager getInstance() { return SingletonHolder.INSTANCE; } private static class SingletonHolder { private static final TaskManager INSTANCE = new TaskManager(); } // 添加新任务 public void addTask(TransferTask task) { allTasks.put(task.getTaskId(), task); taskExecutor.submit(task); } // 重新排序任务:正在执行的任务无法强制中断,只能调整待执行任务的顺序 public void reorderTasks(List<TransferTask> orderedTasks) { // 先停止接受新任务,等待正在执行的任务完成(可设置超时) taskExecutor.shutdown(); try { taskExecutor.awaitTermination(1, TimeUnit.SECONDS); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } // 更新任务列表,清空旧队列 allTasks.clear(); orderedTasks.forEach(task -> allTasks.put(task.getTaskId(), task)); // 重新初始化线程池,提交排序后的任务 taskExecutor = new ThreadPoolExecutor( CONCURRENT_TASK_COUNT, CONCURRENT_TASK_COUNT, 0L, TimeUnit.MILLISECONDS, new PriorityBlockingQueue<>() ); allTasks.values().forEach(taskExecutor::submit); } // 取消指定任务 public void cancelTask(String taskId) { TransferTask task = allTasks.get(taskId); if (task != null) { task.cancelTask(); allTasks.remove(taskId); } } // 销毁任务管理器,在Service销毁时调用 public void shutdown() { taskExecutor.shutdownNow(); allTasks.clear(); } }
第三步:在Android Service中集成
把任务管理器和Service绑定,建议使用前台Service避免被系统后台回收。
public class TransferService extends Service { private TaskManager taskManager; @Override public void onCreate() { super.onCreate(); taskManager = TaskManager.getInstance(); // 启动前台服务,防止被系统杀死 startForeground(1, new NotificationCompat.Builder(this, "transfer_channel") .setContentTitle("传输服务运行中") .setContentText("正在处理上传/下载任务") .setSmallIcon(R.drawable.ic_notification) .build()); } @Override public int onStartCommand(Intent intent, int flags, int startId) { // 从Intent获取任务参数,创建并添加任务 if (intent != null) { String taskId = intent.getStringExtra("task_id"); String type = intent.getStringExtra("task_type"); int priority = intent.getIntExtra("task_priority", 0); String data = intent.getStringExtra("task_data"); TransferTask.TaskType taskType = "UPLOAD".equals(type) ? TransferTask.TaskType.UPLOAD : TransferTask.TaskType.DOWNLOAD; TransferTask task = new TransferTask(taskId, taskType, priority, data); taskManager.addTask(task); } return START_STICKY; } // 提供给外部调用的重新排序方法(可通过Binder或广播触发) public void reorderTasks(List<TransferTask> tasks) { taskManager.reorderTasks(tasks); } @Nullable @Override public IBinder onBind(Intent intent) { return new TransferBinder(); } // Binder类,用于和Activity交互 public class TransferBinder extends Binder { public TransferService getService() { return TransferService.this; } } @Override public void onDestroy() { super.onDestroy(); taskManager.shutdown(); stopForeground(true); } }
关键注意事项
- 任务取消:不要直接中断正在执行的线程(可能导致数据损坏),用
volatile标记位让任务自行终止更安全。 - 排序局限:
PriorityBlockingQueue的排序是在任务入队时完成的,已经入队的任务无法动态调整顺序,所以重新排序时需要重启线程池重新提交任务。 - Service生命周期:务必在Service销毁时调用
TaskManager.shutdown(),避免内存泄漏。
内容的提问来源于stack exchange,提问作者Vino
相关产品推荐
相关产品推荐

