【问题标题】:Can I parallelise an operation driven by a flatMap?我可以并行化由 flatMap 驱动的操作吗?
【发布时间】:2017-03-14 16:45:30
【问题描述】:

假设我有一个类似下面代码的方法,其中 List 被 flatMapped 到单个字符串,每个字符串都应用了一些昂贵的操作。有什么方法可以并行化昂贵的操作,就像我在 Java 8 中使用 parallelStream() 一样?

final List<String> names = new ArrayList<String>() {{
        add("Ringo");
        add("John");
        add("Paul");
        add("George");
    }};

    Observable.just(names).subscribeOn(Schedulers.io())
            .flatMap(new Func1<List<String>, Observable<String>>() {
                @Override
                public Observable<String> call(final List<String> names) {
                    return Observable.create(new Observable.OnSubscribe<String>() {
                        @Override
                        public void call(Subscriber<? super String> subscriber) {
                            for (String name : names) {
                                subscriber.onNext(name);
                            }
                        }
                    });
                }
            })
            .map(new Func1<String, String>() {
                @Override
                public String call(String s) {
                    //Simulate expensive operation
                    try {
                        Thread.sleep(6000);
                    } catch (InterruptedException e) {
                        e.printStackTrace();
                    }
                    return s.toUpperCase();
                }
            }).subscribe(new Subscriber<String>() {
        @Override
        public void onCompleted() {

        }

        @Override
        public void onError(Throwable e) {

        }

        @Override
        public void onNext(String s) {
            Log.v("RXExample", s + " on " + Thread.currentThread().getName());
        }
    });

为了完成,应用答案中推荐的更改如下所示并且效果很好!

final List<String> names = new ArrayList<String>() {{
        add("Ringo");
        add("John");
        add("Paul");
        add("George");
    }};

    Observable.just(names).subscribeOn(Schedulers.io())
            .flatMap(new Func1<List<String>, Observable<String>>() {
                @Override
                public Observable<String> call(final List<String> names) {
                    return Observable.create(new Observable.OnSubscribe<String>() {
                        @Override
                        public void call(final Subscriber<? super String> subscriber) {
                            for (final String name : names) {
                                Observable
                                        .just(name)
                                        .subscribeOn(Schedulers.from(Executors.newFixedThreadPool(5)))
                                        .map(new Func1<String, String>() {
                                            @Override
                                            public String call(String s) {
                                                //Simulate expensive operation
                                                try {
                                                    Thread.sleep(6000);
                                                } catch (InterruptedException e) {
                                                    e.printStackTrace();
                                                }
                                                return s.toUpperCase();
                                            }
                                        }).subscribe(new Observer<String>() {
                                    @Override
                                    public void onCompleted() {

                                    }

                                    @Override
                                    public void onError(Throwable e) {

                                    }

                                    @Override
                                    public void onNext(String s) {
                                        subscriber.onNext(name);
                                    }
                                });
                            }
                        }
                    });
                }
            })
            .subscribe(new Subscriber<String>() {
        @Override
        public void onCompleted() {

        }

        @Override
        public void onError(Throwable e) {

        }

        @Override
        public void onNext(String s) {
            Log.v("RXExample", s + " on " + Thread.currentThread().getName());
        }
    });

【问题讨论】:

    标签: java parallel-processing functional-programming rx-java rx-android


    【解决方案1】:

    您可以使用 flatMap 并行工作,如下例所示。我正在使用 RxJava2 进行测试。

    如需进一步解释,请阅读此处的 flatMap 用法:http://tomstechnicalblog.blogspot.de/2015/11/rxjava-achieving-parallelization.html

    @Test
    public void name() throws Exception {
        final List<String> names = new ArrayList<String>() {{
            add("Ringo");
            add("John");
            add("Paul");
            add("George");
        }};
    
        Observable<String> stringObservable = Observable.fromIterable(names)
                .flatMap(s -> {
                    return longWork(s).doOnNext(s1 -> {
                        printCurrentThread(s1);
                    }).subscribeOn(Schedulers.newThread());
                });
    
        TestObserver<String> test = stringObservable.test();
    
        test.awaitDone(2_000, TimeUnit.MILLISECONDS).assertValueCount(4);
    }
    
    private Observable<String> longWork(String s) throws InterruptedException {
        return Observable.fromCallable(() -> {
            Thread.sleep(1_000);
    
            return s;
        });
    }
    
    private void printCurrentThread(String additional) {
        System.out.println(additional + "_" + Thread.currentThread());
    }
    

    【讨论】:

    • 酷,谢谢!我已将它应用到我的代码 sn-p 以显示它在 AndroidRX 中的外观。
    猜你喜欢
    • 1970-01-01
    • 2019-11-14
    • 2012-03-28
    • 1970-01-01
    • 1970-01-01
    • 2021-08-28
    • 1970-01-01
    • 2015-08-16
    • 2017-10-20
    相关资源
    最近更新 更多