2016年6月9日星期四

SQL Server的WHERE条件"短路评价"大杀器CASE WHEN

WHERE 快速条件 or 慢速条件,执行时,当"快速条件"已经true时,后面的"慢速条件"还会被继续执行吗?理想的就是不被执行,即遵循"短路评价”方式, short circuit evaluation。

N年来我都回避了这个问题,最近不得不直面这个问题,于是在SQL Server 2012上作了调查,结论是和条件是否利用了表函数有关系。

顺便还证实了CASE WHEN语法里的条件绝对是短路评价,这在MSDN文档里明确有说明,这个东西可以用来提高速度。

可惜证据文档取不出来,而且冗长,就大致做个备忘录,有心的人看看就明白了。


1. 条件里没有子查询 -> 短路评价

SELECT * FROM test WHERE some_fast_check()=1 OR some_slow_check(c) = 1


2. 条件里有子查询,但子查询是针对真实表的 -> 准“短路评价”

SELECT * FROM test WHERE some_fast_check()=1 OR c IN (SELECT * FROM normal_or_temp_table)
需要注意的是:子查询一定会最先执行,相当于先做成一个内存表,勉强说得过去,但是万一这个子查询很慢,偏偏第一个条件已经满足了的时候,那就悲剧了,拖累了整个速度
的确,如果第一个条件满足的话,这之后对于形成的内存表进行Scan的次数=0,所以说这是准 “短路评价”。


3. 条件里有子查询,且使用了“表函数” -> 开始失去控制了。

SELECT * FROM test WHERE some_fast_check()=1 OR c IN (SELECT * FROM some_slow_data())
尤其是上述例子里“标量函数”和“表函数”混合的,是最超乎想象的,这时,
some_slow_data()表函数一定最先被执行一次以便其产生内存表,这个只好忍了。
但是这之后就不可忍了!对于这个内存表,总会被进行N次Scan,而N显然取决于主表里符合其他条件的件数。
就是说就算所有行的some_fast_check()=1已经为true了,这后面针对内存表的Scan还是继续做!
傻啊。没办法,从Plan看就是这样的,SQL Server就是这么任性。
当然,结果速度快不快,还要取决于优化器否决定并发执行,有时也不慢。


4. 短路评价大杀器"CASE WHEN"

例如把下面这句整体当做一个表达式来评价时,当"快速条件"满足了时,整个表达式就出结果(1)了,"慢速条件"都不会被执行,
CASE WHEN 快速条件 THEN 1 WHEN 慢速条件 THEN 1 ELSE 0 END
于是,整个SQL改成:
SELECT * FROM test WHERE 
    CASE WHEN some_fast_check()=1 THEN 1
         WHEN c IN (SELECT * FROM some_slow_data()) THEN 1
         ELSE 0
    END = 1
就可以变快了。


还有,不知道DB为什么没有默认开启RECOMPILE选项,这个选项会减少哪些不必要的子查询。
测试时,必须注意用dbcc命令清楚缓存,具体的命令一查就行。
另外一个确定的经验是,IN换成EXISTS或者TABLE JOIN的方式在第一次执行时(没有缓存),Plan是一样的,后来有了Cache之后, TABLE JOIN方式就快些。

2016年6月8日星期三

RxJava(ReactiveX,Observable)的一些大白话

RxJava, ReactiveX, Observable....种种称呼,这东西怎么说呢,优美统一,绝对值得用。但是,她的一些名称取得实在太操蛋了,光这个Observable就与java.util里的那个重复了。另外,关于线程的地方介绍地太暧昧了,非得仔细的琢磨文档甚至得动手试一试,光靠哪些花哨的文档真不靠谱。 于是我做了些练习,最终做了一个挑战,看看如何实现:假设有3个不同类型的Rx, 说不准谁快,其中有一个还有可能出错误,现在让这三个Rx并行执行,但是要等待得到所有的结果,包括错误。

一部分人压根就没想过Publisher和Subscriber的代码分别在什么线程里执行,
做Android的人天然会意识到这个,因为很多例子里都会写subScribeOn(Schedulers.io())和observeOn(AndroidSchedulers.mainThread()),所以没问题。
做了个测试工程RxJava_Test。代码都是在RxJava_Test.java里。
得到几个大白话的结论,可不容易在reactivex的文档里明确找到。
文档里那一堆花花绿绿的图真的能够给人信心吗?虽说画的很有道理,连thread的颜色都有区别,可总不容大白话好啊。


1. 线程指定

如果subScribeOn和observeOn都不指定,那么所有的动作都在当前线程里顺序执行。他们的作用分别是指定了Publisher和Subscriber的动作在什么thread里执行。
subScribeOn(scheduler) 这个函数名字取得真不合适,让人误以为subscriber的回调函数执行在指定的thread里,实际上是指定了Publisher代码执行在指定的thread里。
非要抬杠的话,由于后续的subscriber默认就的确被调用在该thread里,而且publisher的动作实际是在OnSubscibe callback里才被执行的,所以这个函数名马马虎虎说得过去,但是后续的observeOn才是真正指定了subscriber callback执行thread啊
明明是Publisher/Subscriber对称的名词,却冒出了observerOn和subscribeOn两个类似的函数名称,没有人觉得不爽?
测试代码摘要:
    private Observable<String> slowPublisher = Observable.create(subscriber -> {
                sleep(1); //sleep 1ms to just let other thread run so can get predictable output
                println("[SLOW publisher] begin");
                println("[SLOW publisher] do some work");
                sleep(3000);
                println("[SLOW publisher] publish");
                subscriber.onNext("SLOW result");
                subscriber.onCompleted();
                println("[SLOW publisher] end");
            }
    );

    @Test
    public void test_publisher_subscriber_in_current_thread() throws Exception {
        new Thread(() -> {
            println("Enter test function");

            slowPublisher
                    .subscribe(result -> {
                        println("---- subscriber got " + result);
                    });

            println("Leave test function");
        }, "CurrentThread" /*threadName*/).start();

        assertOut("17:02:56.571 @CurrentThread Enter test function");
        assertOut("17:02:56.667 @CurrentThread [SLOW publisher] begin");
        assertOut("17:02:56.667 @CurrentThread [SLOW publisher] do some work");
        assertOut("17:02:59.668 @CurrentThread [SLOW publisher] publish");
        assertOut("17:02:59.668 @CurrentThread ---- subscriber got SLOW result");
        assertOut("17:02:59.668 @CurrentThread [SLOW publisher] end");
        assertOut("17:02:59.669 @CurrentThread Leave test function");
    }

    @Test
    public void test_publisher_in_a_thread_and_subscriber_also_in_same() throws Exception {
        new Thread(() -> {
            println("Enter test function");

            slowPublisher
                    .subscribeOn(schedulerForWorkThread1) //cause publisher run in new thread
                    .subscribe(result -> {
                        println("---- subscriber got " + result);
                    });

            println("Leave test function");
        }, "CurrentThread" /*threadName*/).start();

        assertOut("11:49:51.169 @CurrentThread Enter test function");
        assertOut("11:49:51.217 @CurrentThread Leave test function");
        assertOut("11:49:51.218 @RxWorkThread1 [SLOW publisher] begin");
        assertOut("11:49:51.218 @RxWorkThread1 [SLOW publisher] do some work");
        assertOut("11:49:54.221 @RxWorkThread1 [SLOW publisher] publish");
        assertOut("11:49:54.221 @RxWorkThread1 ---- subscriber got SLOW result");
        assertOut("11:49:54.222 @RxWorkThread1 [SLOW publisher] end");
    }

    @Test
    public void test_publisher_in_a_thread_and_subscriber_in_another() throws Exception {
        new Thread(() -> {
            println("Enter test function");

            slowPublisher
                    .subscribeOn(schedulerForWorkThread1) //cause publisher run in new thread
                    .observeOn(schedulerForWorkThread2) //cause subscriber run in another new thread
                    .subscribe(result -> {
                        sleep(1); //sleep 1ms to just let other thread run so can get predictable output
                        println("---- subscriber got " + result);
                    });

            println("Leave test function");
        }, "CurrentThread" /*threadName*/).start();

        assertOut("12:41:10.523 @CurrentThread Enter test function");
        assertOut("12:41:10.676 @CurrentThread Leave test function");
        assertOut("12:41:10.678 @RxWorkThread1 [SLOW publisher] begin");
        assertOut("12:41:10.678 @RxWorkThread1 [SLOW publisher] do some work");
        assertOut("12:41:13.682 @RxWorkThread1 [SLOW publisher] publish");
        assertOut("12:41:13.682 @RxWorkThread1 [SLOW publisher] end");
        assertOut("12:41:13.683 @RxWorkThread2 ---- subscriber got SLOW result");
    }

追踪了一大坨rxjava-1.1.5-sources.jar里的代码后,大致的动作归纳起来就是:
[Publisher]
in some thread(如果没指定subscribeOn(...)那就是当前thread):
        Do some long time work
        result_queue.putAndSetSignal(result)
[Subscriber]
in some thread(如果没指定observeOn(...)那就和Publisher的一样,否则就是别的thread):
       result = result_queue.getOrWait()
       callback.onNext(result)
       callback.onComplete()


2. 结果硬取

如果有个Rx(例如Rx<String>),要想在当前thread里就取出结果(例如String),那么唯一的方法就是toBlocking()然后.first()或者.iterable()之类的。他就这么设计的,几乎所有的方法的返回值依然是个Rx以便链式调用。
以toBlocking().first()为例,看看 rxjava-1.1.5-sources.jar!/rx/observables/BlockingObservable.java 里这部分的源码:
    private T blockForSingle(final Observable<? extends T> observable) {
        final AtomicReference<T> returnItem = new AtomicReference<T>();
        final AtomicReference<Throwable> returnException = new AtomicReference<Throwable>();
        final CountDownLatch latch = new CountDownLatch(1);

        Subscription subscription = ((Observable<T>)observable).subscribe(new Subscriber<T>() {
            @Override
            public void onCompleted() {
                latch.countDown();
            }

            @Override
            public void onError(final Throwable e) {
                returnException.set(e);
                latch.countDown();
            }

            @Override
            public void onNext(final T item) {
                returnItem.set(item);
            }
        });
        BlockingUtils.awaitForComplete(latch, subscription);

        if (returnException.get() != null) {
            if (returnException.get() instanceof RuntimeException) {
                throw (RuntimeException) returnException.get();
            } else {
                throw new RuntimeException(returnException.get());
            }
        }

        return returnItem.get();
    }
真心不可怕,看完了完全掌控了。大致意思就是:
当前thread:
    out_result = null
    fire publisher_event to start Publisher
    wait complete_event
Subscriber在自己的thread里:
    result = result_queue.getOrWait()
    callback.onNext(result)
    out_result = result
    fire complete_event
破除了暧昧和神秘感。自己做一个也比较简单,无非就是搞个信号等待一下,只是做得粗糙,没考虑到错误时的信号通知,代码就不贴了。


3. Error处理

对于阻塞式Rx(经过toBlocking()转换之后的Rx),必须在try/catch里做。而一般的Rx,必须在subscribe(...onError handler)做。
相比这个error handler,还有一个令人混淆的doOnError(handler),它的效果和前者不一样,并不被看作真正的error handler,就是说出了错误的时候如果没有找到前者,Subscriber thread里会爆OnErrorNotImplementedException。
    @Test
    public void test_error_in_blocking_mode_will_always_throw_out_in_current_thread() throws Exception {
        new Thread(() -> {
            println("Enter test function");

            try {
                errorPublisher
                        .toBlocking().first();
            } catch (Exception e) {
                println("---- test1: " + e);
            }

            try {
                errorPublisher
                        .subscribeOn(schedulerForWorkThread1) //cause publisher run in new thread
                        .toBlocking().first();
            } catch (Exception e) {
                println("---- test2: " + e);
            }

            try {
                errorPublisher
                        .subscribeOn(schedulerForWorkThread1) //cause publisher run in new thread
                        .doOnError(e -> {/*do nothing*/})
                        .toBlocking().first();
            } catch (Exception e) {
                println("---- test3: " + e);
            }

            println("Leave test function");
        }, "CurrentThread" /*threadName*/).start();

        assertOut("14:38:43.078 @CurrentThread Enter test function");
        assertOut("14:38:43.113 @CurrentThread ---- test1: java.lang.NullPointerException");
        assertOut("14:38:43.126 @CurrentThread ---- test2: java.lang.NullPointerException");
        assertOut("14:38:43.152 @CurrentThread ---- test3: java.lang.NullPointerException");
        assertOut("14:38:43.152 @CurrentThread Leave test function");
    }

    @Test
    public void test_error_subscribe_without_error_handler() throws Exception {
        new Thread(() -> {
            println("Enter test function");

            try {
                errorPublisher
                        .subscribeOn(schedulerForWorkThread1) //cause publisher run in new thread
                        .subscribe();
            } catch (Exception e) {
                println("---- should not come here : " + e);
            }

            println("Leave test function");
        }, "CurrentThread" /*threadName*/).start();

        assertOut("21:06:34.342 @CurrentThread Enter test function");
        assertOut("21:06:34.397 @CurrentThread Leave test function");
    }

    @Test
    public void test_error_subscribe_without_error_handler_cause_OnErrorNotImplementedException() throws Exception {
        new Thread(() -> {
            println("Enter test function");

            try {
                errorPublisher
                        .subscribe();
            } catch (Exception e) {
                println("---- test1 error: " + e);
            }

            try {
                errorPublisher
                        .doOnError(e -> {/*do nothing*/})
                        .subscribe();
            } catch (Exception e) {
                println("---- test2 error: " + e);
            }

            println("Leave test function");
        }, "CurrentThread" /*threadName*/).start();

        assertOut("14:15:57.268 @CurrentThread Enter test function");
        assertOut("14:15:57.304 @CurrentThread ---- test1 error: rx.exceptions.OnErrorNotImplementedException");
        assertOut("14:15:57.308 @CurrentThread ---- test2 error: rx.exceptions.OnErrorNotImplementedException");
        assertOut("14:15:57.308 @CurrentThread Leave test function");
    }

    @Test
    public void test_error_subscribe_with_error_handler() throws Exception {
        new Thread(() -> {
            println("Enter test function");

            errorPublisher
                    .subscribe(result -> {
                        //nothing
                    }, e -> { //onError
                        sleep(1); //sleep 1ms to just let other thread run so can get predictable output
                        println("-------- error1: " + e);
                    }, () -> { //onComplete
                        //nothing
                    });

            errorPublisher
                    .subscribeOn(schedulerForWorkThread1) //cause publisher run in new thread
                    .subscribe(result -> {
                        //nothing
                    }, e -> { //onError
                        sleep(1); //sleep 1ms to just let other thread run so can get predictable output
                        println("-------- error2: " + e);
                    }, () -> { //onComplete
                        //nothing
                    });

            println("Leave test function");
        }, "CurrentThread" /*threadName*/).start();

        assertOut("16:19:20.813 @CurrentThread Enter test function");
        assertOut("16:19:20.830 @CurrentThread -------- error1: java.lang.NullPointerException");
        assertOut("16:19:20.840 @CurrentThread Leave test function");
        assertOut("16:19:20.843 @RxWorkThread1 -------- error2: java.lang.NullPointerException");
    }


4. Merge,Concat顺序

Merge并不保证结果的顺序。Concat操作是确保结果顺序的。
    @Test
    public void test_merge_in_parallel_but_result_order_not_predictable() throws Exception {
        new Thread(() -> {
            println("Enter test function");

            slowPublisher
                    .subscribeOn(schedulerForWorkThread1) //cause publisher run in new thread
                    .mergeWith(
                            fastPublisher
                                    .subscribeOn(schedulerForWorkThread2)
                    )
                    .subscribe(result -> {
                        println("---- subscriber got " + result);

                    });

            println("Leave test function");
        }, "CurrentThread" /*threadName*/).start();

        assertOut("10:50:56.507 @CurrentThread Enter test function");
        assertOut("10:50:56.637 @CurrentThread Leave test function");
        assertOut("10:50:56.638 @RxWorkThread1 [SLOW publisher] begin");
        assertOut("10:50:56.638 @RxWorkThread1 [SLOW publisher] do some work");
        assertOut("10:50:56.690 @RxWorkThread2 [FAST publisher] begin");
        assertOut("10:50:56.690 @RxWorkThread2 ---- subscriber got FAST result");
        assertOut("10:50:56.691 @RxWorkThread2 [FAST publisher] end");
        assertOut("10:50:59.640 @RxWorkThread1 [SLOW publisher] publish");
        assertOut("10:50:59.640 @RxWorkThread1 ---- subscriber got SLOW result");
        assertOut("10:50:59.641 @RxWorkThread1 [SLOW publisher] end");
    }

    @Test
    public void test_merge_in_serial_if_not_specify_schedule_thread() throws Exception {
        new Thread(() -> {
            println("Enter test function");

            slowPublisher
                    .mergeWith(
                            fastPublisher
                    )
                    .subscribe(result -> {
                        println("---- subscriber got " + result);

                    });

            println("Leave test function");
        }, "CurrentThread" /*threadName*/).start();

        assertOut("10:52:59.765 @CurrentThread Enter test function");
        assertOut("10:52:59.846 @CurrentThread [SLOW publisher] begin");
        assertOut("10:52:59.846 @CurrentThread [SLOW publisher] do some work");
        assertOut("10:53:02.849 @CurrentThread [SLOW publisher] publish");
        assertOut("10:53:02.850 @CurrentThread ---- subscriber got SLOW result");
        assertOut("10:53:02.850 @CurrentThread [SLOW publisher] end");
        assertOut("10:53:02.904 @CurrentThread [FAST publisher] begin");
        assertOut("10:53:02.904 @CurrentThread ---- subscriber got FAST result");
        assertOut("10:53:02.904 @CurrentThread [FAST publisher] end");
        assertOut("10:53:02.905 @CurrentThread Leave test function");
    }

    @Test
    public void test_concat_will_always_serially_so_predictable_order() throws Exception {
        new Thread(() -> {
            println("Enter test function");

            slowPublisher
                    .subscribeOn(schedulerForWorkThread1) //cause publisher run in new thread
                    .concatWith(
                            fastPublisher
                                    .subscribeOn(schedulerForWorkThread2)
                    )
                    .subscribe(result -> {
                        println("---- subscriber got " + result);

                    });

            println("Leave test function");
        }, "CurrentThread" /*threadName*/).start();

        assertOut("10:53:55.314 @CurrentThread Enter test function");
        assertOut("10:53:55.454 @CurrentThread Leave test function");
        assertOut("10:53:55.456 @RxWorkThread1 [SLOW publisher] begin");
        assertOut("10:53:55.456 @RxWorkThread1 [SLOW publisher] do some work");
        assertOut("10:53:58.461 @RxWorkThread1 [SLOW publisher] publish");
        assertOut("10:53:58.461 @RxWorkThread1 ---- subscriber got SLOW result");
        assertOut("10:53:58.463 @RxWorkThread1 [SLOW publisher] end");
        assertOut("10:53:58.516 @RxWorkThread2 [FAST publisher] begin");
        assertOut("10:53:58.516 @RxWorkThread2 ---- subscriber got FAST result");
        assertOut("10:53:58.517 @RxWorkThread2 [FAST publisher] end");
    }


练的差不多了,

5. 挑战:异种Rx并发执行+统一掌控

假设有3个不同类型的Rx, 说不准谁快,其中有一个还有可能出错误,现在想让这三个Rx并行执行,但是要等待得到所有的结果,包括错误。
说真的,如果不是硬要用RxJava里的方法,那就各自并行publish,subscribe,在callback里用CountDownLatch信号量来搞算了,没啥难的,还好懂。
    private static class RxResult<T> {
        T result;
        Throwable error;
    }

    private static class RxResultAndSignal<T> extends RxResult<T> {
        CountDownLatch latch = new CountDownLatch(1);
    }
    @Test
    public void test_merge_3_different_type_rx_by_self_idea() throws Exception {
        new Thread(() -> {
            println("Enter test function");

            RxResultAndSignal<String> r1 = new RxResultAndSignal<>();
            RxResultAndSignal<StringBuilder> r2 = new RxResultAndSignal<>();
            RxResultAndSignal<Number> r3 = new RxResultAndSignal<>();

            slowPublisher
                    .subscribeOn(schedulerForWorkThread1) //cause publisher run in new thread
                    .subscribe(result -> {
                        r1.result = result;
                    }, e -> { //onError
                        r1.error = e;
                        r1.latch.countDown();
                    }, () -> { //onComplete
                        r1.latch.countDown();
                    });

            incompatiblePublisher
                    .subscribeOn(schedulerForWorkThread2)
                    .subscribe(result -> {
                        r2.result = result;
                    }, e -> { //onError
                        r2.error = e;
                        r2.latch.countDown();
                    }, () -> { //onComplete
                        r2.latch.countDown();
                    });

            buggyPublisher
                    .subscribeOn(schedulerForWorkThread3)
                    .subscribe(result -> {
                        r3.result = result;
                    }, e -> { //onError
                        r3.error = e;
                        r3.latch.countDown();
                    }, () -> { //onComplete
                        r3.latch.countDown();
                    });

            try {
                r1.latch.await();
                r2.latch.await();
                r3.latch.await();
            } catch (InterruptedException e) {
                e.printStackTrace();
            }

            sleep(1); //sleep 1ms to let other thread run so can get predictable output
            println("---- got all result: {" + r1.result + "}, {" + r2.result + "}, {" + r3.result + "}");
            println("---- got all error: {" + r1.error + "}, {" + r2.error + "}, {" + r3.error + "}");

            println("Leave test function");
        }, "CurrentThread" /*threadName*/).start();

        assertOut("13:27:55.033 @CurrentThread Enter test function");
        assertOut("13:27:55.065 @RxWorkThread1 [SLOW publisher] begin");
        assertOut("13:27:55.065 @RxWorkThread1 [SLOW publisher] do some work");
        assertOut("13:27:55.082 @RxWorkThread2 [Incompatible publisher] begin");
        assertOut("13:27:55.082 @RxWorkThread2 [Incompatible publisher] do some work");
        assertOut("13:27:55.086 @RxWorkThread3 [Buggy publisher] begin");
        assertOut("13:27:57.084 @RxWorkThread2 [Incompatible publisher] end");
        assertOut("13:27:58.069 @RxWorkThread1 [SLOW publisher] publish");
        assertOut("13:27:58.070 @RxWorkThread1 [SLOW publisher] end");
        assertOut("13:27:58.071 @CurrentThread ---- got all result: {SLOW result}, {Incompatible result}, {null}");
        assertOut("13:27:58.071 @CurrentThread ---- got all error: {null}, {null}, {java.lang.NullPointerException}");
        assertOut("13:27:58.071 @CurrentThread Leave test function");
    }

如果非要全用RxJava里的方法实现的话,折腾了一阵子,发现这个方法比较健壮:
    @Test
    public void test_merge_3_different_type_rx() throws Exception {
        new Thread(() -> {
            println("Enter test function");

            RxResult<String> r1 = new RxResult<>();
            RxResult<StringBuilder> r2 = new RxResult<>();
            RxResult<Number> r3 = new RxResult<>();

            Completable.merge(
                    slowPublisher
                            .subscribeOn(schedulerForWorkThread1) //cause publisher run in new thread
                            .doOnNext(result -> r1.result = result)
                            .doOnError(e -> r1.error = e)
                            .toCompletable() //so can mergeWith other type rx
                            .onErrorComplete() //auto call complete when error
                    ,
                    incompatiblePublisher
                            .subscribeOn(schedulerForWorkThread2)
                            .doOnNext(result -> r2.result = result)
                            .doOnError(e -> r2.error = e)
                            .toCompletable() //so can mergeWith other type rx
                            .onErrorComplete() //auto call complete when error
                    ,
                    buggyPublisher
                            .subscribeOn(schedulerForWorkThread3)
                            .doOnNext(result -> r3.result = result)
                            .doOnError(e -> r3.error = e)
                            .toCompletable() //so can mergeWith other type rx
                            .onErrorComplete() //auto call complete when error
            )
                    .await(/*can specify total timeout*/);

            sleep(1); //sleep 1ms to let other thread run so can get predictable output
            println("---- got all result: {" + r1.result + "}, {" + r2.result + "}, {" + r3.result + "}");
            println("---- got all error: {" + r1.error + "}, {" + r2.error + "}, {" + r3.error + "}");

            println("Leave test function");
        }, "CurrentThread" /*threadName*/).start();

        assertOut("13:47:50.789 @CurrentThread Enter test function");
        assertOut("13:47:50.852 @RxWorkThread1 [SLOW publisher] begin");
        assertOut("13:47:50.852 @RxWorkThread1 [SLOW publisher] do some work");
        assertOut("13:47:50.862 @RxWorkThread2 [Incompatible publisher] begin");
        assertOut("13:47:50.863 @RxWorkThread2 [Incompatible publisher] do some work");
        assertOut("13:47:50.869 @RxWorkThread3 [Buggy publisher] begin");
        assertOut("13:47:52.868 @RxWorkThread2 [Incompatible publisher] end");
        assertOut("13:47:53.854 @RxWorkThread1 [SLOW publisher] publish");
        assertOut("13:47:53.854 @RxWorkThread1 [SLOW publisher] end");
        assertOut("13:47:53.856 @CurrentThread ---- got all result: {SLOW result}, {Incompatible result}, {null}");
        assertOut("13:47:53.856 @CurrentThread ---- got all error: {null}, {null}, {java.lang.NullPointerException}");
        assertOut("13:47:53.856 @CurrentThread Leave test function");
    }

绝对还有其他方法。RxJava构建了一个自己的世界,要折腾什么样子的都不稀奇,不过些复杂了就不好懂,光看暧昧的方法调用其实无法让人确信如自己希望的那样动作。

2016年6月7日星期二

Windows/Mac下的Docker,如何把Windows的目录映射到容器里,有个小坑

如果你的目录没有放在C:¥Users目录(Mac就是/Users)以及其子目录下,那恭喜你,你会掉到坑里。 docker-compose工具同理收到这个影响,docker-compose.yml里的目录映射(volumns)设定里的Host目录,都要注意。

Windows/Mac OS X下的docker是运行在一个Linux虚拟机的, 在这个虚拟机里,运行多个容器。三层啊,真不爽,什么事儿都是隔山跨水的,现在事儿就来了,容器里会报错说找不到被映射过来的文件。
拿Windows举例,有个目录叫做 C:¥WinTestDir,放了几个文件,现在运行一个容器例如busybox,试图把这个目录映射到容器里的 /xxx/yyy目录下,无论以C:¥WinTestDir还是/c/WindowsDir还是/WindowsDir都没有效果。
docker run -it -v C:¥WinTestDir:/xxx/yyy busybox
结果会提示说"和:非法字符。

docker run -it -v /WinTestDir:/xxx/yyy busybox
docker run -it -v /c/WinTestDir:/xxx/yyy busybox

结果容器里的/xxx/yyy目录下是空的。
这时一般都会意识到,这个冒号左边的目录实际上得是那个虚拟机里存在的目录才行,那么哪些目录被映射到虚拟机里了呢?Docker的文档里写了,
Windows系统:C:¥Users ->  /c/Users
Mac OS X系统:/Users  -> /Users
其它的目录压根没被映射进去,跟别提之后往容器里映射了。
docker-machine ssh default
docker@rethink:~$ ls /c
Users
docker@rethink:~$ ls /c/Users/Administrator
... Documents Downloads ...
docker@rethink:~$ ls /WinTestDir
No such file or directory
总之,除非用VirtualBox修改这个虚拟机的共享目录设定,否则在虚拟机里只能看到C:¥Users以下的文件。

所以,省事儿的方法就是把WinTestDir挪到C:¥Users下随便一层目录,例如
C:¥Users¥q¥Documents¥WinTestDir,然后就可以用
/c/Users/q/Documents/WinTestDir来做映射源目录了。例如
docker run -it -v /c/Users/q/Documents/WinTestDir:/xxx/yyy busybox

/ # ls /xxx/yyy
  ... some files ...
这就OK了。

同理,docker-composer所使用的docker-composer.yml文件里关于目录映射(volumns)的地方就得小心。例如这里写的someHostDir。
test:
  image: busybox
  command: /bin/find /xxx/yyy
  volumes:
    - ./someHostDir:/xxx/yyy
看起来./someHostDir用的是相对目录,挺优美的样子,可是如果这个目录不属于C:¥Users底下的,例如C:/work,那就实际上无法映射了。可以用这个config子命令来看看实际目录。Windows下就变成了/c/work/someHostDir了,实际在虚拟机里这个目录是和Windows机器里的目录没有关联起来,就是说是空的。
暂时离开了Windows,就用Mac OS X做实验:
例如当前位于 /private/tmp/docker-compose-test目录下,结果是这样
$ docker-compose config

networks: {}
services:
  test:
    command: /bin/find /xxx/yyy
    image: busybox
    network_mode: bridge
    volumes:
    - /private/tmp/docker-compose-test/someHostDir:/xxx/yyy:rw
version: '2.0'
volumes: {}
那么运行docker-compose up会显示出/xxx/yyy里的内容,什么都没有。
$ docker-compose up
Recreating dockercomposetest_test_1
Attaching to dockercomposetest_test_1
test_1  | /xxx/yyy
dockercomposetest_test_1 exited with code 0

2016年5月16日星期一

IntelliJ IDEA Community版独特的一个bug:某些目录下的文件能显示但不被编译

IntelliJ IDEA Ultimate不存在的问题。没办法只好自己动手折腾了。
版本是2016.1.1 和 2016.1.2。Windows上和Mac上都一样的现象。
问题的详细都提交给JetBrain公司了,但是还没有修复:
https://youtrack.jetbrains.com/oauth?state=%2Fissue%2FIDEA-155471
问题的起源是,有个工程里有个目录叫做rcs,好死不死的正好和古老的CVS版本管理系统所使用的隐藏信息文件后缀同名,默认就不被显示在工程里,看都看不到。这个在Ultimate版里也存在。
这个不是什么大问题,到设定里到File Types设定里,找到"Ignore files and folders"设定,从中把rcs字眼给删除。这就可以显示了,
但是之后,Ultimate版可以正常编译rcs目录下的java文件,Community版不行。
于是想出了一个怪招,见一个rrr符号连接,指向rcs目录,然后在IDEA里把rrr加到"Ignore files and folders"设定里以便隐藏。这的确可以编译了,但是,每次改变一点东西都要手动到菜单里之行“编译 ...当前文件”命令或者全体重新编译,不然代码就执行老的。
于是研究起来,发现是lib/jps-model.jar里又个org/jetbrains/jps/model/impl/JpsFileTypesConfigurationImpl.class 里的固定字符串里改掉就行了。
于是用二进制编辑器替换掉rcs -> rrr,在用jar或者7zip什么的替换掉jps-model.jar里的class为改好的class文件就好了。

这就可以结束了。



具体的研究过程就是把IDEA里所有的jar都展开,用gstrings工具找字符串再grep,最后确定了三个文件可疑,挨个试一下就成了。

cd "/Applications/IntelliJ IDEA CE.app/Contents"
find . -name '*.jar' | while read f; do rm -fr ~/tmp/IntelliJ_IDEA_Community/$f; mkdir -p ~/tmp/IntelliJ_IDEA_Community/$(dirname $f); unzip $f -d ~/tmp/IntelliJ_IDEA_Community/$f > /dev/null; done
cd /Users/q/tmp/IntelliJ_IDEA_Community/
find . -name '*.class' -or -name '*.xml' -or -name '*.json' -or -name '*.properties' | while read f; do gstrings -n 3 $f | grep -w rcs && echo ---- $f ----;done

        NG:
            ".dependency-info;CVS;RCS;SCCS;rcs;
            ---- ./lib/idea.jar/com/intellij/openapi/fileTypes/impl/FileTypeConfigurable$FileTypePanel.class ----
        NG:
            u*.hprof;*.pyc;*.pyo;*.rbc;*.yarb;*~;.DS_Store;.git;.hg;.svn;CVS;RCS;SCCS;__pycache__;_svn;rcs;vssver.scc;vssver2.scc;
            ---- ./lib/idea.jar/com/intellij/openapi/fileTypes/impl/FileTypeManagerImpl.class ----
        OK:
            CVS;SCCS;RCS;rcs;.DS_Store;.svn;.pyc;.pyo;*.pyc;*.pyo;.git;*.hprof;_svn;.hg;*.lib;*~;__pycache__;.bundle;vssver.scc;vssver2.scc;*.rbc;
            ---- ./lib/jps-model.jar/org/jetbrains/jps/model/impl/JpsFileTypesConfigurationImpl.class ----
        NG:
              <ignoreFiles list="CVS;SCCS;RCS;rcs;.DS_Store;.svn;.pyc;.pyo"/>
            ---- ./lib/resources.jar/CommunityFileTypes.xml ----