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
| DataStream<Tuple3<Long, String, String>> input1 = ...; input1 = input1.assignTimestampsAndWatermarks(new AscendingTimestampExtractor<Tuple3<Long, String, String>>() {
@Override public long extractAscendingTimestamp(Tuple3<Long, String, String> arg0) { return arg0.f0; }
});
DataStream<Tuple2<Long, String>> input2 = ...; input2 = input2.assignTimestampsAndWatermarks(new AscendingTimestampExtractor<Tuple2<Long, String>>() { @Override public long extractAscendingTimestamp(Tuple2<Long, String> stringStringTuple2) { return stringStringTuple2.f0; } });
input1.coGroup(input2).where(new KeySelector<Tuple3<Long, String, String>, String>() { @Override public String getKey(Tuple3<Long, String, String> itemEntity) throws Exception { return itemEntity.f1; } }) .equalTo(new KeySelector<Tuple2<Long, String>, String>() { @Override public String getKey(Tuple2<Long, String> value) throws Exception { return value.f1; } }) .window(TumblingEventTimeWindows.of(Time.minutes(1))) .apply(new CoGroupFunction<Tuple3<Long, String, String>, Tuple2<Long, String>, String>() { @Override public void coGroup(Iterable<Tuple3<Long, String, String>> first, Iterable<Tuple2<Long, String>> second, Collector<String> collector) throws Exception { StringBuilder buffer = new StringBuilder(); buffer.append("DataStream first:\n"); for (Tuple3<Long, String, String> value : first) { buffer.append(value).append("\n"); } buffer.append("DataStream second:\n"); for (Tuple2<Long, String> value : second) { buffer.append(value.f0).append("=>").append(value.f1).append("\n"); } collector.collect(buffer.toString()); } }) .print();
|