@@ -947,9 +947,7 @@ public void fusedBoundary() {
947947 Flowable .range (1 , 10000 )
948948 .switchMap ((Function <Integer , Flowable <Object >>) _ -> Flowable .just (2 ).hide ()
949949 .observeOn (Schedulers .single ())
950- .map ((Function <Integer , Object >) _ -> {
951- return Thread .currentThread ().getName ();
952- }))
950+ .map ((Function <Integer , Object >) _ -> Thread .currentThread ().getName ()))
953951 .to (TestHelper .<Object >testConsumer ())
954952 .awaitDone (5 , TimeUnit .SECONDS )
955953 .assertNever (thread )
@@ -1110,9 +1108,7 @@ public void cancellationShouldTriggerInnerCancellationRace() throws Throwable {
11101108
11111109 int n = 10_000 ;
11121110 for (int i = 0 ; i < n ; i ++) {
1113- Flowable .<Integer >create (it -> {
1114- it .onNext (0 );
1115- }, BackpressureStrategy .MISSING )
1111+ Flowable .<Integer >create (it -> it .onNext (0 ), BackpressureStrategy .MISSING )
11161112 .switchMap (_ -> createFlowable (inner ))
11171113 .observeOn (Schedulers .computation ())
11181114 .doFinally (outer ::incrementAndGet )
@@ -1129,12 +1125,8 @@ Flowable<Integer> createFlowable(AtomicInteger inner) {
11291125 return Flowable .<Integer >unsafeCreate (s -> {
11301126 SerializedSubscriber <Integer > it = new SerializedSubscriber <>(s );
11311127 it .onSubscribe (new BooleanSubscription ());
1132- Schedulers .cached ().scheduleDirect (() -> {
1133- it .onNext (1 );
1134- }, 0 , TimeUnit .MILLISECONDS );
1135- Schedulers .cached ().scheduleDirect (() -> {
1136- it .onNext (2 );
1137- }, 0 , TimeUnit .MILLISECONDS );
1128+ Schedulers .cached ().scheduleDirect (() -> it .onNext (1 ), 0 , TimeUnit .MILLISECONDS );
1129+ Schedulers .cached ().scheduleDirect (() -> it .onNext (2 ), 0 , TimeUnit .MILLISECONDS );
11381130 })
11391131 .doFinally (inner ::incrementAndGet );
11401132 }
0 commit comments