【问题标题】:Prove me that PublishSubject in RxJava is not thread safe证明 RxJava 中的 PublishSubject 不是线程安全的
【发布时间】:2017-04-06 19:42:24
【问题描述】:

声明 PublishSubject 在 RxJava 中不是线程安全的。好的。

我正在尝试找到任何示例,我正在尝试构建任何示例来模拟竞争条件,这会导致不需要的结果。但我不能:(

谁能提供一个例子来证明 PublishSubject 不是线程安全的?

【问题讨论】:

  • 证明是文档says so.
  • 哦,谢谢你的链接!!!你太棒了!!!
  • 您不需要证明。您需要证明或断言它 线程安全的。否则它应该被认为是不安全的。
  • @EJP,哇!专家的另一个非常有用的评论已经仔细阅读了这个问题!!!!你也很厉害!!!
  • 我投票决定将此问题作为题外话结束,因为这不是一个问题,它是对代码的模糊请求,以重现在未指定版本的库中可能存在或不存在的错误。

标签: java multithreading thread-safety rx-java


【解决方案1】:

通常,人们会问为什么他们的设置会出现意外和/或崩溃,答案是:因为他们同时调用 Subject 上的 onXXX 方法:

import java.util.concurrent.CountDownLatch;
import java.util.concurrent.atomic.AtomicInteger;

import org.junit.Test;

import rx.Scheduler.Worker;
import rx.exceptions.MissingBackpressureException;
import rx.observers.AssertableSubscriber;
import rx.schedulers.Schedulers;
import rx.subjects.*;

public class PublishSubjectRaceTest {

    @Test
    public void racy() throws Exception {
        Worker worker = Schedulers.computation().createWorker();
        try {
            for (int i = 0; i < 1000; i++) {
                AtomicInteger wip = new AtomicInteger(2);

                PublishSubject<Integer> ps = PublishSubject.create();

                AssertableSubscriber<Integer> as = ps.test(1);

                CountDownLatch cdl = new CountDownLatch(1);

                worker.schedule(() -> {
                    if (wip.decrementAndGet() != 0) {
                        while (wip.get() != 0) ;
                    }
                    ps.onNext(1);

                    cdl.countDown();
                });
                if (wip.decrementAndGet() != 0) {
                    while (wip.get() != 0) ;
                }
                ps.onNext(1);

                cdl.await();

                as.assertFailure(MissingBackpressureException.class, 1);
            }
        } finally {
            worker.unsubscribe();
        }
    }

    @Test
    public void nonRacy() throws Exception {
        Worker worker = Schedulers.computation().createWorker();
        try {
            for (int i = 0; i < 1000; i++) {
                AtomicInteger wip = new AtomicInteger(2);

                Subject<Integer, Integer> ps = PublishSubject.<Integer>create()
                    .toSerialized();

                AssertableSubscriber<Integer> as = ps.test(1);

                CountDownLatch cdl = new CountDownLatch(1);

                worker.schedule(() -> {
                    if (wip.decrementAndGet() != 0) {
                        while (wip.get() != 0) ;
                    }
                    ps.onNext(1);

                    cdl.countDown();
                });
                if (wip.decrementAndGet() != 0) {
                    while (wip.get() != 0) ;
                }
                ps.onNext(1);

                cdl.await();

                as.assertFailure(MissingBackpressureException.class, 1);
            }
        } finally {
            worker.unsubscribe();
        }
    }
}

【讨论】:

  • 我知道 toSerialized() 可以使主题线程安全,但是 toSerialized() 的折衷是什么?表现 ?因为我看到我们在每个 onXX() 方法中都有仔细检查同步块。谢谢!
  • toSerialized 增加了一些开销,除非您对延迟有严格的限制,否则您的应用几乎不会注意到差异。
  • 我不明白你的例子。在您的每种情况下,多线程在哪里。一切都发生在一个线程中,主题未订阅。很不清楚...
  • 多线程:Schedulers.computation().createWorker() & worker.schedule(...);订阅发生在ps.test(1)
【解决方案2】:

我找到了证据。我认为这个例子比@akarnokd 提供的更明显。

    AtomicInteger counter = new AtomicInteger();

    // Thread-safe
    // SerializedSubject<Object, Object> subject = PublishSubject.create().toSerialized();

    // Not Thread Safe
    PublishSubject<Object> subject = PublishSubject.create();

    Action1<Object> print = (x) -> System.out.println(Thread.currentThread().getName() + " " + counter);

    Consumer<Integer> sleep = (s) -> {
        try {
            Thread.sleep(s);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
    };

    subject
            .doOnNext(i -> counter.incrementAndGet())
            .doOnNext(i -> counter.decrementAndGet())
            .doOnNext(print)
            .filter(i -> counter.get() != 0)
            .doOnNext(i -> {
                        throw new NullPointerException("Concurrency detected");
                    }
            )
            .subscribe();

    Runnable r = () -> {
        for (int i = 0; i < 100000; i++) {
            sleep.accept(1);
            subject.onNext(i);
        }
    };

    ExecutorService pool = Executors.newFixedThreadPool(2);
    pool.execute(r);
    pool.execute(r);

【讨论】:

    猜你喜欢
    • 2017-03-07
    • 1970-01-01
    • 2011-01-25
    • 2021-12-18
    • 2019-01-13
    • 2018-07-11
    • 1970-01-01
    相关资源
    最近更新 更多