|
|
0 c6 v+ T0 z; |/ p. z<h4 id="flink系列文章">Flink系列文章</h4>' S7 t- L! R7 {( E* H
<ol>% O, {7 R- y- o
<li><a href="https://www.ikeguang.com/?p=1976">第01讲:Flink 的应用场景和架构模型</a></li>4 v; {$ c4 v8 R
<li><a href="https://www.ikeguang.com/?p=1977">第02讲:Flink 入门程序 WordCount 和 SQL 实现</a></li>
" B& s* [& M- B# i<li><a href="https://www.ikeguang.com/?p=1978">第03讲:Flink 的编程模型与其他框架比较</a></li>
. ~) r$ o3 E" G; m& ^9 n0 l<li><a href="https://www.ikeguang.com/?p=1982">第04讲:Flink 常用的 DataSet 和 DataStream API</a></li>, c1 U: F7 _5 u0 \
<li><a href="https://www.ikeguang.com/?p=1983">第05讲:Flink SQL & Table 编程和案例</a></li>
! b1 n# D7 ~2 Q5 ~<li><a href="https://www.ikeguang.com/?p=1985">第06讲:Flink 集群安装部署和 HA 配置</a></li>
. W& S* H$ p0 b<li><a href="https://www.ikeguang.com/?p=1986">第07讲:Flink 常见核心概念分析</a></li>% S( a# ~$ ]" q4 v1 E
<li><a href="https://www.ikeguang.com/?p=1987">第08讲:Flink 窗口、时间和水印</a></li>
3 p' V1 u g3 |! S) l1 ^, F<li><a href="https://www.ikeguang.com/?p=1988">第09讲:Flink 状态与容错</a></li>
0 O0 z3 D; x$ {+ W+ F) _+ h& n( g1 W</ol>/ c" ?+ d8 E- m& t2 F& e
<blockquote>
" T) Q& K5 \3 Q n<p>关注公众号:<code>大数据技术派</code>,回复<code>资料</code>,领取<code>1024G</code>资料。</p>) j& K) Y+ o/ X( h. U
</blockquote>. v* t/ e* w, O8 H
<p>这一课时将介绍 Flink 中提供的一个很重要的功能:旁路分流器。</p>
7 U- q: s2 N7 [- q$ l" I: g, p7 E% O<h3 id="分流场景">分流场景</h3>
* T0 U- z3 P5 d<p>我们在生产实践中经常会遇到这样的场景,需把输入源按照需要进行拆分,比如我期望把订单流按照金额大小进行拆分,或者把用户访问日志按照访问者的地理位置进行拆分等。面对这样的需求该如何操作呢?</p>
" R9 [/ \4 j! b* j$ R4 t<h3 id="分流的方法">分流的方法</h3>
5 Y. O2 u; O7 f& e<p>通常来说针对不同的场景,有以下三种办法进行流的拆分。</p>
& U/ @( a" W' o" J+ m. H3 r<h4 id="filter-分流">Filter 分流</h4>1 A: e7 p& p- j" Q" p
<p><img src="https://kingcall.oss-cn-hangzhou.aliyuncs.com/blog/img/CgqCHl7CAy6ADUaXAACSFUbdpuA911-20210223084827182.png" ></p>8 @8 |) Z2 Q3 O0 ~0 V' _
<p>Filter 方法我们在第 04 课时中(Flink 常用的 DataSet 和 DataStream API)讲过,这个算子用来根据用户输入的条件进行过滤,每个元素都会被 filter() 函数处理,如果 filter() 函数返回 true 则保留,否则丢弃。那么用在分流的场景,我们可以做多次 filter,把我们需要的不同数据生成不同的流。</p>
1 ^# b3 Q* d3 r2 I) l! r<p>来看下面的例子:</p>
; Q5 @, V$ {. c- e+ X4 w |- s, P<p>复制代码</p>
& \, \* q' s4 f- H" Z5 n- ^<pre><code class="language-java">public static void main(String[] args) throws Exception {
9 s, F7 T) p# a9 ^2 u& f2 g StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
) p3 o; x7 Q# z, h4 \5 y //获取数据源
, b5 [1 s" \- _ List data = new ArrayList<Tuple3<Integer,Integer,Integer>>();
7 ^# S: K2 K# [: n2 c data.add(new Tuple3<>(0,1,0));
6 k/ n# H# P) {0 ]0 K data.add(new Tuple3<>(0,1,1));
+ H8 e8 G' w7 V1 o' `3 W data.add(new Tuple3<>(0,2,2));
: R. j3 y; S6 Q: q1 q4 M2 f; w data.add(new Tuple3<>(0,1,3));
- f. `, J" w5 E% Q data.add(new Tuple3<>(1,2,5));! \; X W8 ?3 v6 m H' K6 S! T+ `
data.add(new Tuple3<>(1,2,9));
# S! h: v/ C: B. G' B1 {% o8 B data.add(new Tuple3<>(1,2,11));$ i! K7 q4 ~- ^, b4 ^
data.add(new Tuple3<>(1,2,13));
- E9 Q9 L+ m2 w; [; I! |3 b: `6 j4 A# T# z2 _5 B4 y+ X9 ^
DataStreamSource<Tuple3<Integer,Integer,Integer>> items = env.fromCollection(data);0 a4 x. n4 I( r
4 B& S% g9 C- D
+ X: y% A4 _8 O" u# @4 ?+ s
; w$ g5 }4 T1 O8 C* I
SingleOutputStreamOperator<Tuple3<Integer, Integer, Integer>> zeroStream = items.filter((FilterFunction<Tuple3<Integer, Integer, Integer>>) value -> value.f0 == 0);/ \' p1 I% K; r- H" ^, i
1 }9 q1 n7 h/ p# R1 F, U
SingleOutputStreamOperator<Tuple3<Integer, Integer, Integer>> oneStream = items.filter((FilterFunction<Tuple3<Integer, Integer, Integer>>) value -> value.f0 == 1);+ ?' a# O/ j+ W u; Q$ M$ m1 L
2 w5 ]; Q5 {$ x- a; Q, c3 O+ u0 A: I
9 }- _5 [, N& e6 |# \+ `/ W h! O0 ?8 }7 N
zeroStream.print();
& ]; `+ U( f! k: w; r( j" O
! ?% |8 A, e& I( @0 z. H+ |1 g a oneStream.printToErr();4 A* K8 @- N0 g
: @, \; N) g/ o( r0 }% T. h8 t4 D: _4 k0 l, i* |- O
2 b/ B/ Y2 U" Q0 b5 H/ s! \7 P
8 c, i" M- H- P; f+ e- @0 d8 y4 L3 h8 t; W3 _
//打印结果
9 X) a: u7 f2 e7 ]9 O% }" G
2 H: ^ T/ y0 i* s5 ^ String jobName = "user defined streaming source";
( ~% X# A) C& f: c+ O( i% ^3 [: J! ?: ~
env.execute(jobName);# _# ^. Z( ]0 e% K
. K9 a9 r2 E1 {& V* u6 l( S- c
}
5 }1 d6 O5 q) s1 W</code></pre>
4 v/ o/ ]4 Z+ j2 u7 o+ l<p>在上面的例子中我们使用 filter 算子将原始流进行了拆分,输入数据第一个元素为 0 的数据和第一个元素为 1 的数据分别被写入到了 zeroStream 和 oneStream 中,然后把两个流进行了打印。</p>
0 @: Z- x- i! r" V<p><img src="https://kingcall.oss-cn-hangzhou.aliyuncs.com/blog/img/Ciqc1F7CA2WAYbshAAKj494h86s723-20210223084827296.png" ></p>
/ r0 g: R& w2 I8 T* M! J+ V7 b<p>可以看到 zeroStream 和 oneStream 分别被打印出来。</p>$ L/ K9 O T; R. s
<p>Filter 的弊端是显而易见的,为了得到我们需要的流数据,需要多次遍历原始流,这样无形中浪费了我们集群的资源。</p>& e" Z" Z/ i1 {2 Y2 W1 G
<h4 id="split-分流">Split 分流</h4>2 R0 N. C2 R x( @" v
<p>Split 也是 Flink 提供给我们将流进行切分的方法,需要在 split 算子中定义 OutputSelector,然后重写其中的 select 方法,将不同类型的数据进行标记,最后对返回的 SplitStream 使用 select 方法将对应的数据选择出来。</p>
9 e) y: S s- F9 L<p>我们来看下面的例子:</p># |: [( |; ^9 n# C" D. B6 s
<p>复制代码</p>
* R' ~3 c$ h; Y: u6 C; g<pre><code class="language-java">public static void main(String[] args) throws Exception {
7 G/ v" r/ B6 u0 E5 ?/ E: i2 x: _
5 U7 d7 [! d) R: m7 L" H! c' }' H
* `8 B3 S1 m, {4 Y
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
) e7 m: n* u% G% N
; X1 b' x3 t& ^1 }4 d" n( E //获取数据源5 U( _2 V4 a& \1 h8 a
% f0 d4 f9 c/ `0 W0 K List data = new ArrayList<Tuple3<Integer,Integer,Integer>>();* n8 X5 s r; J9 h9 v+ ]+ j4 a
$ f% W$ w$ V2 \, r f- |2 B
data.add(new Tuple3<>(0,1,0));1 j8 j* ]* b' m
+ H# p6 q0 P" L" [0 |# L
data.add(new Tuple3<>(0,1,1));8 z8 ~; |2 [' m; @/ I; i/ y& t0 R7 i
( m2 F4 i& m3 t# Y# q }' d
data.add(new Tuple3<>(0,2,2));1 V2 @, ]3 D9 C5 L
% a% M, E" _1 q9 W& u data.add(new Tuple3<>(0,1,3));( O k) r7 L# F7 \$ p3 z
& l! q& B% ^* R% i; i; J data.add(new Tuple3<>(1,2,5));$ L% P1 j& U: e
7 T# d6 Y* J, a4 l+ L/ [# Q; x0 W data.add(new Tuple3<>(1,2,9));. y, k# g) a+ x- b! F
. X5 S& _( Y) N$ x# f% R
data.add(new Tuple3<>(1,2,11));
6 {8 K, b: \* { K
2 e# U: Y* E/ S P' F* ~ data.add(new Tuple3<>(1,2,13));; Z7 {: Q c0 m- u6 M# V3 k
% F- h# i' V) F; C8 V* J5 x: a: ~' `' ^& g3 J5 K
" m; k- Z# m- W6 R, r
, a* H* Q* b# c8 ?1 ]0 D0 _9 {/ L7 P6 G/ ]% w, H/ d
DataStreamSource<Tuple3<Integer,Integer,Integer>> items = env.fromCollection(data);1 J8 o* [( d. v& J3 }
) k" w- b& H; Q( g7 a
6 V. v; q% `- }. x' Q5 \
; k- q2 ~- h+ e0 @0 a3 h
( A2 R$ V2 Y- J2 E9 F7 y( m, m
- [: `$ U. J/ j3 E* @1 h7 ~ SplitStream<Tuple3<Integer, Integer, Integer>> splitStream = items.split(new OutputSelector<Tuple3<Integer, Integer, Integer>>() {" v7 X2 h2 s5 S0 e8 Z
" ^" M5 l4 [: W n9 ~
@Override; U8 O; z# N( h0 P) {
+ P: \# ]6 h4 j, P8 t public Iterable<String> select(Tuple3<Integer, Integer, Integer> value) {7 U3 h) `, r/ O' {! e5 j
% u3 F/ N5 @0 r& T List<String> tags = new ArrayList<>();" }( P V- V3 ?( D1 W
3 o( u4 O% J6 E m. `& ~! \
if (value.f0 == 0) {8 u; r; G$ [+ s3 _. a2 `. q/ E
+ M8 a7 N7 G5 X tags.add("zeroStream");
# {- T% g" o, O% M2 n4 |% {( u G, `' b% A( q& I
} else if (value.f0 == 1) {1 x' q$ t" M/ ~
, ^ o$ t$ N h1 C: q tags.add("oneStream");! }+ q7 Y* | E5 M- K+ `. `" N) \
O! v1 Q3 D7 c% h! [+ G
}; X( W* S# p6 g3 R0 F4 w( r. ^5 A
* y, X3 p; ^7 X
return tags;
4 V8 ^, s7 J2 T! _
- E& s& d0 a" B' l% o' | }
/ X+ R' L8 D8 z: S' A: r
$ G( {' g" Q: v# ~: w$ _ });% x/ x5 [) T+ \9 m9 B
# l9 \. W7 |- X1 ^0 i) z! `9 t- u& O8 A- |$ p8 P4 a
. y2 t3 h. ?8 G" W4 S/ W0 b {1 [1 y splitStream.select("zeroStream").print();
1 _& `( u& F' a/ b9 D7 f. ?" a9 p7 r# m; v4 a
splitStream.select("oneStream").printToErr();- K4 N! b& Y) `7 h
. C; s5 h+ U& X1 W+ S1 V# r) _) x1 F3 L3 K
. C# b+ r6 y+ q# s! t7 t
//打印结果& u# q! e- a; _
3 x& G( j V$ w( z( v: _6 { String jobName = "user defined streaming source";# I4 f1 Y) C4 G, `; j7 r
, }8 c' y7 `6 | env.execute(jobName);3 f7 t% P1 n) a8 H# ?# Z: T6 k
3 e9 o# A {( ?6 L2 j5 ^
}
, c/ z" `6 b6 e4 c5 {</code></pre>9 C, G+ H6 W7 t6 i4 [8 Q8 Z# ^9 P
<p>同样,我们把来源的数据使用 split 算子进行了切分,并且打印出结果。</p>
; w% V! z0 t" Q# n# o8 j n/ N& u<p><img src="https://kingcall.oss-cn-hangzhou.aliyuncs.com/blog/img/CgqCHl7CA4aAbUSJAAG1LWNB3qw627-20210223084827377.png" ></p>
$ l% G$ ~8 z* B; d<p>但是要注意,使用 split 算子切分过的流,是不能进行二次切分的,假如把上述切分出来的 zeroStream 和 oneStream 流再次调用 split 切分,控制台会抛出以下异常。</p>
, j: } m& @% N/ U<p>复制代码</p>
" a3 z. A1 s- p$ g+ a) w/ G, T# N<pre><code class="language-java">Exception in thread "main" java.lang.IllegalStateException: Consecutive multiple splits are not supported. Splits are deprecated. Please use side-outputs.
) P$ R4 Z6 M" f" |4 P! p- }7 U+ p</code></pre>2 E. W; R- T) r7 Z: K
<p>这是什么原因呢?我们在源码中可以看到注释,该方式已经废弃并且建议使用最新的 SideOutPut 进行分流操作。</p># p+ G, v9 P4 w* U8 h. @: k- ^& ~
<p><img src="https://kingcall.oss-cn-hangzhou.aliyuncs.com/blog/img/CgqCHl7CA6OAJ-JDAAIrh1JSAEo033-20210223084827602.png" ></p>
: C2 Q/ v6 r* q; v7 V<h4 id="sideoutput-分流">SideOutPut 分流</h4>$ n5 E0 _- E$ D
<p>SideOutPut 是 Flink 框架为我们提供的最新的也是最为推荐的分流方法,在使用 SideOutPut 时,需要按照以下步骤进行:</p>% U L) w* d" w
<ul>
9 y9 D$ X% n# [! s% ^0 U<li>定义 OutputTag</li>
- Y6 e' @8 {9 F" f h8 I<li>调用特定函数进行数据拆分( C/ j. `% m2 w7 {4 M
<ul>
* |& a( P2 ^2 s5 B! |. Z4 Q8 Y P<li>ProcessFunction</li>
/ C7 ?6 V; O/ g8 \<li>KeyedProcessFunction</li>
7 h) Z5 G: G/ P& p! M9 B<li>CoProcessFunction</li>. N6 R* l( S; e! y0 k7 K; R1 }
<li>KeyedCoProcessFunction</li>
- M T/ p) N* G% S% I; a3 a<li>ProcessWindowFunction</li>
2 A0 E0 B: @' u6 n<li>ProcessAllWindowFunction</li>& i5 V3 ]; y" u% F
</ul>
+ e D' a* t" j& j</li>
' T" P4 Q$ J6 F5 ?</ul>
# R5 q& i0 Z0 m# A$ W<p>在这里我们使用 ProcessFunction 来讲解如何使用 SideOutPut:</p>4 z# c# P/ H4 `
<p>复制代码</p>- \" r: [4 T$ g+ w& L# q5 }" y
<pre><code class="language-java">public static void main(String[] args) throws Exception {
' v- k- j" X. c( w* V' n8 K
/ N4 O; E% ~8 H% u
' w1 p1 |: ?& E1 E" l3 Z$ N) z: ~. d! i' ?: ~7 c0 c. U& i
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();; S9 y* Y& }' a* _; O/ m* k: ^, `
# O* C- f9 K1 p) d3 X
//获取数据源9 I7 ~ U3 S H% h* Y( I
. Q7 q/ y: e" ~2 D
List data = new ArrayList<Tuple3<Integer,Integer,Integer>>();1 ^& _3 Y" d9 R4 o
2 U2 M. }0 t. m2 T. p. ]8 Y data.add(new Tuple3<>(0,1,0));
5 B. W* a% l* s3 C. }$ ]' L3 R2 X9 c
data.add(new Tuple3<>(0,1,1));
! v2 S$ l) N. R2 Y' M) R
: Q5 ~; p: _ y* S data.add(new Tuple3<>(0,2,2));9 M E& A2 r \$ i2 [- M
. J6 Y/ s+ S7 |& } data.add(new Tuple3<>(0,1,3));: D4 I9 r- |/ [+ {
0 S( M# [( I% V$ `" b9 j! Y( E" D data.add(new Tuple3<>(1,2,5));0 ^+ [3 |; l e- r
: c3 e% c) d3 z, @7 ^! K data.add(new Tuple3<>(1,2,9));% ~( U7 e/ c" z
+ H" l: V7 U& J: E data.add(new Tuple3<>(1,2,11));
& Q' n4 U: o `, `1 l/ E+ h: w4 o3 N% T E
data.add(new Tuple3<>(1,2,13));, I( A7 O0 h' p% g+ I; j. W
( I8 Z- w+ `: t; o1 x% `2 g
* y% C2 Q9 N9 h6 O/ H6 }/ R+ D& N9 j
/ l% u% Z0 s1 ]3 K0 O5 n& r& @, ?
' ^: G% |9 a: T" E DataStreamSource<Tuple3<Integer,Integer,Integer>> items = env.fromCollection(data);. ~7 R( l6 q3 n1 ^( Q% [" {* o
w. {% g% T6 x! ~2 S* m6 X7 A3 i* a) [( S) C7 ~8 N, M1 d
$ |: D1 |8 I. l3 `( e9 R OutputTag<Tuple3<Integer,Integer,Integer>> zeroStream = new OutputTag<Tuple3<Integer,Integer,Integer>>("zeroStream") {};( Z! h7 A* B6 ]* e* {8 |7 I/ L
0 q T0 j! z8 e1 W OutputTag<Tuple3<Integer,Integer,Integer>> oneStream = new OutputTag<Tuple3<Integer,Integer,Integer>>("oneStream") {};
* j) O3 M; E/ M0 m
2 \$ W$ n8 m) t8 f) K, b& ]. \! M' c" U
/ T1 Y, N# w# g& k$ i3 L% @* ~3 s: G9 z y
3 x: X, F# v. t
SingleOutputStreamOperator<Tuple3<Integer, Integer, Integer>> processStream= items.process(new ProcessFunction<Tuple3<Integer, Integer, Integer>, Tuple3<Integer, Integer, Integer>>() {, X; D! Q. p+ K2 U; [) K
0 n: z" D! z5 m- p1 i# l
@Override$ V2 n8 } L) S& p& \4 h7 L
7 K8 L+ y+ a- K, Q
public void processElement(Tuple3<Integer, Integer, Integer> value, Context ctx, Collector<Tuple3<Integer, Integer, Integer>> out) throws Exception {
i5 `$ q# H7 c, x/ u
/ W9 P8 X2 Y$ C: U( K/ P" F% e7 V( S& k* ~% j/ U* k, e# n
% |6 A3 O# N" z7 V& S1 X4 \8 ]
if (value.f0 == 0) {2 e: h7 e6 f7 ^ N! J
9 |# q Q$ m) W* z$ V0 L
ctx.output(zeroStream, value);) b q0 ~5 i5 K+ _
$ P& x; w7 G+ y! w+ |7 X. e( u- `
} else if (value.f0 == 1) {
% h+ d5 n- b) L/ O% X6 r( l7 Z l) t5 b7 |) K. U
ctx.output(oneStream, value);) a9 C) V8 e( i% r4 g' j+ D2 }. O
% p" M" J' |/ x, |: u7 V
}
. z4 I' u" W# j! X/ R/ q( D9 m, n, R8 m- ?2 }7 B i' d; d
}6 ], x' `2 @9 c4 j+ M
0 o: }5 Z7 Y4 V/ |! ^0 S, E B$ K
});
6 t: N9 [3 F e# a
' T0 P$ [* B' Y" M1 p
8 B) F. K3 z, j" j) z
" ?+ ^% _- f: M9 ~- Y( M DataStream<Tuple3<Integer, Integer, Integer>> zeroSideOutput = processStream.getSideOutput(zeroStream);
! n8 [$ v" |3 o- b# ?4 }' D
& B7 Z+ n/ i$ j9 h/ Y, Y( [ DataStream<Tuple3<Integer, Integer, Integer>> oneSideOutput = processStream.getSideOutput(oneStream);1 G0 {8 E$ x0 B/ f) Q
1 ?, X0 K5 t+ \- d' S
9 k* s4 [' U/ }5 O9 c3 k3 c# L' V" } P5 b( K' w+ F, E, _" e
zeroSideOutput.print();
4 W; {8 f0 A o: R r1 H5 \7 ? J* o, w# g
oneSideOutput.printToErr();
3 ?; x1 |* ]: q% o m: }/ Y5 `# F" v; E6 w! K
6 r) P" z6 \6 g
6 \8 v$ F. j8 f7 y' c8 q2 a) g) b# U4 i
- T5 Q) M5 A9 M# P( e
//打印结果
; F+ B- N6 E" C2 m+ b/ f. N2 H5 x
String jobName = "user defined streaming source";0 S- b3 u0 D- `: H% S6 \' j7 w5 Q, p
4 Q$ n- k" Q" Q! B* _3 y5 M0 ?/ [
env.execute(jobName);
3 Q) u, r+ {3 G4 J/ _ E1 ]/ ?6 W9 ^ e* r$ w
}
1 A* g; o \- t+ t</code></pre>5 k! H( m# U5 H8 B4 W' f. i
<p>可以看到,我们将流进行了拆分,并且成功打印出了结果。这里要注意,Flink 最新提供的 SideOutPut 方式拆分流是<strong>可以多次进行拆分</strong>的,无需担心会爆出异常。</p>) l$ n2 b1 W# D; v
<p><img src="https://kingcall.oss-cn-hangzhou.aliyuncs.com/blog/img/CgqCHl7CBMKAGHoUAAM-5UL5geg132-20210223084827698.png" ></p>/ i3 U E4 w) k6 Z, i& k9 T0 d( {) j
<h3 id="总结">总结</h3>
; l. _1 b2 d, u$ w<p>这一课时我们讲解了 Flink 的一个小的知识点,是我们生产实践中经常遇到的场景,Flink 在最新的版本中也推荐我们使用 SideOutPut 进行流的拆分。</p>( j- C, P8 i' L0 v! R/ `2 x
<blockquote>
& l: T# D- ~5 @" z9 j<p>关注公众号:<code>大数据技术派</code>,回复<code>资料</code>,领取<code>1024G</code>资料。</p>
3 P, |- H5 x: C: o, H6 r</blockquote>
4 y- k: U8 u- s! G2 G Z
4 R$ x3 Z2 N1 Z6 N4 [0 J |
|