Skip to content

Navigation Menu

Sign in
Appearance settings

Search code, repositories, users, issues, pull requests...

Provide feedback

We read every piece of feedback, and take your input very seriously.

Saved searches

Use saved searches to filter your results more quickly

Appearance settings

Latest commit

 

History

History
History
114 lines (95 loc) · 4.1 KB

File metadata and controls

114 lines (95 loc) · 4.1 KB
Copy raw file
Download raw file
Open symbols panel
Edit and raw actions
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
/**
* Copyright 2013 Netflix, Inc.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package rx;
import java.util.HashMap;
import java.util.Map;
import org.junit.Test;
import rx.CovarianceTest.CoolRating;
import rx.CovarianceTest.ExtendedResult;
import rx.CovarianceTest.HorrorMovie;
import rx.CovarianceTest.Media;
import rx.CovarianceTest.Movie;
import rx.CovarianceTest.Rating;
import rx.CovarianceTest.Result;
import rx.EventStream.Event;
import rx.observables.GroupedObservable;
import rx.util.functions.Action1;
import rx.util.functions.Func1;
import rx.util.functions.Func2;
public class ZipTests {
@Test
public void testZipObservableOfObservables() {
EventStream.getEventStream("HTTP-ClusterB", 20)
.groupBy(new Func1<Event, String>() {
@Override
public String call(Event e) {
return e.instanceId;
}
// now we have streams of cluster+instanceId
}).flatMap(new Func1<GroupedObservable<String, Event>, Observable<Map<String, String>>>() {
@Override
public Observable<Map<String, String>> call(final GroupedObservable<String, Event> ge) {
return ge.scan(new HashMap<String, String>(), new Func2<Map<String, String>, Event, Map<String, String>>() {
@Override
public Map<String, String> call(Map<String, String> accum, Event perInstanceEvent) {
accum.put("instance", ge.getKey());
return accum;
}
});
}
})
.take(10)
.toBlockingObservable().forEach(new Action1<Map<String, String>>() {
@Override
public void call(Map<String, String> v) {
System.out.println(v);
}
});
System.out.println("**** finished");
}
/**
* This won't compile if super/extends isn't done correctly on generics
*/
@Test
public void testCovarianceOfZip() {
Observable<HorrorMovie> horrors = Observable.from(new HorrorMovie());
Observable<CoolRating> ratings = Observable.from(new CoolRating());
Observable.<Movie, CoolRating, Result> zip(horrors, ratings, combine).toBlockingObservable().forEach(action);
Observable.<Movie, CoolRating, Result> zip(horrors, ratings, combine).toBlockingObservable().forEach(action);
Observable.<Media, Rating, ExtendedResult> zip(horrors, ratings, combine).toBlockingObservable().forEach(extendedAction);
Observable.<Media, Rating, Result> zip(horrors, ratings, combine).toBlockingObservable().forEach(action);
Observable.<Media, Rating, ExtendedResult> zip(horrors, ratings, combine).toBlockingObservable().forEach(action);
Observable.<Movie, CoolRating, Result> zip(horrors, ratings, combine);
}
Func2<Media, Rating, ExtendedResult> combine = new Func2<Media, Rating, ExtendedResult>() {
@Override
public ExtendedResult call(Media m, Rating r) {
return new ExtendedResult();
}
};
Action1<Result> action = new Action1<Result>() {
@Override
public void call(Result t1) {
System.out.println("Result: " + t1);
}
};
Action1<ExtendedResult> extendedAction = new Action1<ExtendedResult>() {
@Override
public void call(ExtendedResult t1) {
System.out.println("Result: " + t1);
}
};
}
Morty Proxy This is a proxified and sanitized view of the page, visit original site.