You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.25 07:29:08