【问题标题】:Realm - implementing asynchronous queueRealm - 实现异步队列
【发布时间】:2017-04-18 16:17:06
【问题描述】:

我有一个 dagger-singleton-wrapper 来处理我的基本 Realm 请求。其中一个看起来像这样:

public void insertOrUpdateAsync(final List<RealmMessage> messages, @Nullable final OnInsertListener listener) {
    Realm instance = getRealmInstance();
    instance.executeTransactionAsync(realm -> {
                List<RealmMessage> newMessages = insertOrUpdateMessages(realm, messages);
            },
            () -> success(listener, instance),
            error -> error(listener, error, instance));
}

private List<RealmMessage> insertOrUpdateMessages(@NonNull Realm realm, @NonNull final List<RealmMessage> messages) {
    ...
    return realm.copyToRealmOrUpdate(unattendedMessages);
}

效果很好。

但是有一个极端情况——长话短说——我多次启动 insertOrUpdateAsynch()。经过一些请求,我得到了这个:

Caused by: java.util.concurrent.RejectedExecutionException: Task java.util.concurrent.FutureTask@b7b848 rejected from io.realm.internal.async.RealmThreadPoolExecutor@80f96e1[Running, pool size = 17, active threads = 17, queued tasks = 100, completed tasks = 81]

我的问题是:我应该如何在不重建整个应用程序流程的情况下处理这个问题。 我的想法是通过 RxJava 对传入的请求进行排队。我对吗?我应该考虑和自学哪些运营商?

还是我以完全错误的方式处理这个问题? 从我的大部分谷歌搜索中,我注意到问题主要在于像我这样的循环启动方法。我没有使用任何。在我的情况下,问题是这种方法是由多个响应启动的,并且由于当前的后端实现而改变它是不可能的。

【问题讨论】:

    标签: java android asynchronous realm rx-java


    【解决方案1】:

    如果您不想重新设计您的应用程序,您可以使用计数信号量。您将看到两个线程将立即获得锁。另一个线程将阻塞,直到某个调用释放一个锁。不建议在没有 Timeout 的情况下使用 acquire()。

    为了使用 RxJava,您必须更改应用程序的设计,而 RxJava 中的速率限制并不是那么容易,因为它完全与输出有关。

    private final Semaphore semaphore = new Semaphore(2);
    
    @Test
    public void name() throws Exception {
        Thread t1 = new Thread(() -> {
            doNetworkStuff();
        });
        Thread t2 = new Thread(() -> {
            doNetworkStuff();
        });
        Thread t3 = new Thread(() -> {
            doNetworkStuff();
        });
    
        t1.start();
        t2.start();
        t3.start();
    
        Thread.sleep(1500);
    }
    
    private void doNetworkStuff() {
        try {
            System.out.println("enter doNetworkStuff");
    
            semaphore.acquire();
    
            System.out.println("acquired");
    
            Thread.sleep(1000);
    
        } catch (InterruptedException e) {
            e.printStackTrace(); // Don't do this!!
        } finally {
            semaphore.release();
        }
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2011-05-19
      • 2021-04-01
      • 2013-08-25
      • 1970-01-01
      • 2018-04-19
      • 1970-01-01
      • 2012-01-23
      • 1970-01-01
      相关资源
      最近更新 更多