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

周期性更新与增删操作同步问题:Flowable用户任务API调用实现

解决Flowable周期性任务与手动操作的同步问题

看起来你现在遇到的是周期性批量任务和手动业务操作之间的并发同步问题——也就是每隔60秒自动更新用户任务的流程,和手动的创建/更新/删除任务操作可能会同时修改数据库,导致数据不一致或者冲突对吧?我结合RxJava和数据库操作的实践,给你几个可行的解决方案:

1. 用串行调度器统一所有数据库操作

最直接的方式是给所有涉及用户任务的数据库操作(包括周期性更新和手动增删改)指定同一个单线程调度器,确保所有操作串行执行,从根源上避免并发冲突。

你可以用RxJava自带的Schedulers.single(),或者自定义一个单线程的Executor调度器:

// 定义全局的串行数据库调度器
private final Scheduler dbSerialScheduler = Schedulers.single();

// 调整周期性任务的调度逻辑
Flowable.interval(60, TimeUnit.SECONDS)
    .concatMapCompletable(_tick -> mUserRepository.getAllUsers()
        .subscribeOn(dbSerialScheduler) // 数据库查询也走串行调度
        .flatMapObservable(Observable::fromIterable)
        .flatMapCompletable(user -> 
            // 下载操作可以在IO线程并行,但数据库操作切换到串行调度
            downloadAndPersistTasks(user)
                .subscribeOn(Schedulers.io())
                .observeOn(dbSerialScheduler)
        )
    )
    .subscribeOn(dbSerialScheduler)
    .subscribe(
        () -> {}, 
        error -> Log.e("TaskSync", "Sync failed", error)
    );

// 手动的增删改操作也统一使用这个串行调度器
public Completable updateUserTask(Task task) {
    return mTaskRepository.updateTask(task)
        .subscribeOn(dbSerialScheduler);
}

2. 引入显式锁机制控制并发

如果不想完全串行化所有操作,可以用Java的ReentrantLock或者RxJava的信号量来做显式的并发控制,确保同一时间只有一个任务相关的修改操作在执行。

示例代码:

private final ReentrantLock taskOperationLock = new ReentrantLock();

// 给周期性任务的核心逻辑加锁
private Completable downloadAndPersistTasksWithLock(User user) {
    return Completable.fromRunnable(taskOperationLock::lock)
        .andThen(downloadAndPersistTasks(user))
        .doFinally(taskOperationLock::unlock);
}

// 手动操作同样加锁
public Completable deleteUserTask(String taskId) {
    return Completable.fromRunnable(taskOperationLock::lock)
        .andThen(mTaskRepository.deleteTask(taskId))
        .doFinally(taskOperationLock::unlock);
}

// 调整周期性任务调用加锁后的方法
Flowable.interval(60, TimeUnit.SECONDS)
    .flatMapCompletable(_tick -> mUserRepository.getAllUsers()
        .flatMapObservable(Observable::fromIterable)
        .flatMapCompletable(this::downloadAndPersistTasksWithLock)
        .subscribeOn(Schedulers.io())
    , false, 1)
    .subscribe(...);

注意:锁需要包裹整个任务修改流程(比如查询旧数据、删除、插入新数据的完整步骤),避免只锁部分操作导致中间状态冲突。

3. 利用数据库事务保证原子性

如果你的数据库支持事务,把downloadAndPersistTasks里的“删除旧任务+插入新任务”逻辑封装成一个原子事务,同时手动的增删改操作也使用事务,这样可以避免出现“旧任务删了但新任务没插入”的中间不一致状态,也能减少冲突概率。

以Room数据库为例:

@Dao
public interface TaskDao {
    // 用@Transaction标记原子操作
    @Transaction
    default void replaceUserTasks(String userId, List<Task> newTasks) {
        deleteOldTasks(userId);
        insertNewTasks(newTasks);
    }

    @Query("DELETE FROM tasks WHERE user_id = :userId")
    void deleteOldTasks(String userId);

    @Insert
    void insertNewTasks(List<Task> tasks);
}

这样每次更新用户任务都是一个原子操作,要么全部完成,要么全部回滚,不会留下脏数据。

4. 避免周期性任务自身重叠执行

你的现有代码用了flatMapCompletable(..., false, 1),但Flowable.interval会按时发射tick,不管上一次任务有没有完成,可能导致前一次的同步还在进行,下一次又开始了。换成concatMapCompletable可以解决这个问题——它会等待上一个Completable执行完成后,再处理下一个tick:

Flowable.interval(60, TimeUnit.SECONDS)
    .concatMapCompletable(_tick -> mUserRepository.getAllUsers()
        .flatMapObservable(Observable::fromIterable)
        .flatMapCompletable(user -> downloadAndPersistTasks(user)
            .subscribeOn(Schedulers.io())
        )
        .subscribeOn(Schedulers.io())
    )
    .subscribe(...);

推荐组合方案

我个人推荐方案1 + 方案3 + 方案4的组合:

  • 用串行调度器统一所有数据库操作,避免跨操作并发冲突;
  • 用事务保证单个任务更新的原子性;
  • 用concatMapCompletable避免周期性任务自身的重叠执行。

这样既能保证数据一致性,又能维持合理的执行效率。

内容的提问来源于stack exchange,提问作者elnino

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 08:20:30