0%

Flink join

flink join

Cogroup

CoGroupFunction中会返回所有数据,不管有没有匹配上

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();

上面的例子,左流有三个元素 Tuple3<String,String,String>,右流有两个元素Tuple2<String,String>
两个流第一个元素相互关联。分别指定两个流的事件时间字段。
两个流关联后,按照EventTime划分窗口。与单流类似。
不管元素是否可以关联上,都会输出

用户可以定义CoGroupFunction函数, 可以实现在窗口内,任意组合,如笛卡尔积

举例说明:

window Join

interval join

broadcast join