|
|
* v* f) M3 S' L3 B
<h4 id="flink系列文章">Flink系列文章</h4>- a2 f9 x7 M7 ?+ U2 e U: F1 m: V
<ol>
9 B | _! e- |1 ?/ R0 d3 J' W<li><a href="https://www.ikeguang.com/?p=1976">第01讲:Flink 的应用场景和架构模型</a></li>1 x7 g L5 r! d' x3 ?1 q6 A
<li><a href="https://www.ikeguang.com/?p=1977">第02讲:Flink 入门程序 WordCount 和 SQL 实现</a></li>1 }' L% ~! r+ p g0 y5 Z1 _2 c
<li><a href="https://www.ikeguang.com/?p=1978">第03讲:Flink 的编程模型与其他框架比较</a></li>! ~" c0 c1 G1 k+ S& g7 i! Q
<li><a href="https://www.ikeguang.com/?p=1982">第04讲:Flink 常用的 DataSet 和 DataStream API</a></li>8 M- L4 O4 i3 I [, [
<li><a href="https://www.ikeguang.com/?p=1983">第05讲:Flink SQL & Table 编程和案例</a></li>
; S e) `- [, t% N<li><a href="https://www.ikeguang.com/?p=1985">第06讲:Flink 集群安装部署和 HA 配置</a></li>
4 L- K8 u% |! h/ q! j; E/ G. {<li><a href="https://www.ikeguang.com/?p=1986">第07讲:Flink 常见核心概念分析</a></li>
/ ?. [, K8 e, P+ B- I1 \1 h7 j<li><a href="https://www.ikeguang.com/?p=1987">第08讲:Flink 窗口、时间和水印</a></li>/ c0 h+ e2 E$ B- ~7 m; A
<li><a href="https://www.ikeguang.com/?p=1988">第09讲:Flink 状态与容错</a></li>1 n6 ]& ?7 m9 u$ Q, M9 k
</ol>. B7 a3 U& L. o7 O0 ~1 \
<blockquote>
8 u4 E. ~6 ^* j9 f' z! N. ~<p>关注公众号:<code>大数据技术派</code>,回复<code>资料</code>,领取<code>1024G</code>资料。</p>6 l) L: \2 ?+ U) z! s+ {( Q( M5 s
</blockquote>
; E4 V1 J( m5 O<p>这一课时将介绍 Flink 中提供的一个很重要的功能:旁路分流器。</p>
4 Y n8 \" p4 _9 c' ?. K<h3 id="分流场景">分流场景</h3>
* w5 r: G# _" j1 e# k+ K<p>我们在生产实践中经常会遇到这样的场景,需把输入源按照需要进行拆分,比如我期望把订单流按照金额大小进行拆分,或者把用户访问日志按照访问者的地理位置进行拆分等。面对这样的需求该如何操作呢?</p>5 N9 |9 L& \" p3 F" R! I( S4 {3 z
<h3 id="分流的方法">分流的方法</h3>
, p9 [6 n: R5 [0 o2 H<p>通常来说针对不同的场景,有以下三种办法进行流的拆分。</p>+ p; I; Y' i! I
<h4 id="filter-分流">Filter 分流</h4>
% ^- P* _, [$ v<p><img src="https://kingcall.oss-cn-hangzhou.aliyuncs.com/blog/img/CgqCHl7CAy6ADUaXAACSFUbdpuA911-20210223084827182.png" ></p>5 z* t- J* e& n# T' B" r8 e
<p>Filter 方法我们在第 04 课时中(Flink 常用的 DataSet 和 DataStream API)讲过,这个算子用来根据用户输入的条件进行过滤,每个元素都会被 filter() 函数处理,如果 filter() 函数返回 true 则保留,否则丢弃。那么用在分流的场景,我们可以做多次 filter,把我们需要的不同数据生成不同的流。</p>6 a5 U. w" b! w8 u- s
<p>来看下面的例子:</p>
\ q4 H/ l% u h7 G<p>复制代码</p>
% \ Z) Y( j5 @) v8 c# L' ~0 V<pre><code class="language-java">public static void main(String[] args) throws Exception {8 U( p* Y- W- m8 [* C2 q0 ^/ Z" l
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();* ^; {) ~$ \/ A1 D5 P; K7 y* |7 E
//获取数据源
1 t, v E# q' W& o List data = new ArrayList<Tuple3<Integer,Integer,Integer>>();# i# }; w4 M8 e3 L; o$ Y, T" \
data.add(new Tuple3<>(0,1,0));2 c: ]7 Y9 ` K5 D% n- J0 ], ^3 N H
data.add(new Tuple3<>(0,1,1));- u9 S( _& y; c9 y' G; q7 A
data.add(new Tuple3<>(0,2,2));: h: r, ~+ U1 Z! |- M
data.add(new Tuple3<>(0,1,3));
. d% }1 E( m3 z, a- U+ \# \ data.add(new Tuple3<>(1,2,5));3 x0 u4 ~/ ~# c0 J' l! s2 o. a; [
data.add(new Tuple3<>(1,2,9));
, ?' u: {2 o5 H. v) O- _ data.add(new Tuple3<>(1,2,11));7 {6 e( z( R, r L3 ?3 K) N( S
data.add(new Tuple3<>(1,2,13));7 `9 j, a5 f5 ~
) E4 K) M9 T. k$ `0 ] DataStreamSource<Tuple3<Integer,Integer,Integer>> items = env.fromCollection(data);3 V' u( M# r8 ~$ `( X! k; O
1 o: }" R2 Z+ D
& @8 e& H/ |! c" ]7 }; t' @, Q
SingleOutputStreamOperator<Tuple3<Integer, Integer, Integer>> zeroStream = items.filter((FilterFunction<Tuple3<Integer, Integer, Integer>>) value -> value.f0 == 0);
* G" |1 a! o; V( X
5 @9 X# J0 W# D' u' a SingleOutputStreamOperator<Tuple3<Integer, Integer, Integer>> oneStream = items.filter((FilterFunction<Tuple3<Integer, Integer, Integer>>) value -> value.f0 == 1);3 r/ V6 H- r4 ^4 q
3 D8 ~. ?, ~- B8 Y- ^
W/ w% X5 I) j X
* s. L x, x# c4 e4 h zeroStream.print();
9 U2 N. @, \0 z! ^
) ~, w- f- N5 E7 e# | oneStream.printToErr();
0 A" {4 P2 w' i% ~/ A2 I) R, {0 L. |. ?) F% L7 J
5 e1 [! q$ z" T3 y1 E' U
% `2 H1 e5 j; V/ {# d! p
1 h. \/ d# z8 M, h0 {4 o
2 y) J s7 w$ I. g( l //打印结果! y' r. k7 l: f# J7 m0 E( m
: Y" L1 A9 h9 {' h/ Y/ D
String jobName = "user defined streaming source";
3 s- N( y; ^. ]: i
* {! I" Z) {8 v' p env.execute(jobName);* C( }6 r. x8 {0 W
) u7 Y0 s9 a3 H G. c( k1 h: |
}
) Y. N5 v8 [1 f$ c' p4 G</code></pre> r E* c0 }& }) I f' m3 u
<p>在上面的例子中我们使用 filter 算子将原始流进行了拆分,输入数据第一个元素为 0 的数据和第一个元素为 1 的数据分别被写入到了 zeroStream 和 oneStream 中,然后把两个流进行了打印。</p>
% C$ [3 A2 W& F1 h* F7 b$ r<p><img src="https://kingcall.oss-cn-hangzhou.aliyuncs.com/blog/img/Ciqc1F7CA2WAYbshAAKj494h86s723-20210223084827296.png" ></p>+ n. H u! u/ w8 L/ s
<p>可以看到 zeroStream 和 oneStream 分别被打印出来。</p>+ t. \( s* {' o
<p>Filter 的弊端是显而易见的,为了得到我们需要的流数据,需要多次遍历原始流,这样无形中浪费了我们集群的资源。</p>+ F3 O+ M1 A0 H+ B1 g& M) H
<h4 id="split-分流">Split 分流</h4>
+ f4 l- d7 |5 [7 f: H- a<p>Split 也是 Flink 提供给我们将流进行切分的方法,需要在 split 算子中定义 OutputSelector,然后重写其中的 select 方法,将不同类型的数据进行标记,最后对返回的 SplitStream 使用 select 方法将对应的数据选择出来。</p>
7 T/ i4 o' i6 w/ \0 i( S<p>我们来看下面的例子:</p>2 D% k0 r9 v4 }. ^
<p>复制代码</p>
7 h4 d' U" f; f/ w1 s3 l<pre><code class="language-java">public static void main(String[] args) throws Exception {) [6 [* W7 O* A& `8 B
" }. @; m5 B* C7 G
( I: ^ Q O, Y
6 S F9 k8 x0 n, ]3 g
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
4 R' m8 u5 {* z( u1 a
* A1 a( s: d( c1 ~ //获取数据源
% V9 s/ C. [0 U7 {" i/ K! B& k' j
List data = new ArrayList<Tuple3<Integer,Integer,Integer>>();
|) z+ b! z1 z+ ?1 r/ X
( a+ V! @) n [) l5 _5 `& N data.add(new Tuple3<>(0,1,0));
! X# i5 Z( \: K+ o6 e) g8 R6 }7 [. j
0 `! i) [% k0 [) i6 o( ? data.add(new Tuple3<>(0,1,1));
" }1 x3 O: ?5 s9 b% r+ W5 z. D0 T H- ?9 O3 G0 @) @
data.add(new Tuple3<>(0,2,2));5 R4 R) s) ]: I) r# Y5 _
( K6 G+ }2 R" u' q
data.add(new Tuple3<>(0,1,3));
9 W: E, [( Q1 k0 i% s. s0 ^' X9 h1 }" _/ L1 P! s0 c8 d
data.add(new Tuple3<>(1,2,5));
/ r7 j) \1 |9 m# Q. O/ T
$ U$ T" t5 I1 r; e7 Q) n* T data.add(new Tuple3<>(1,2,9));
1 B) [7 ^" w9 K7 h
. C. ~. I V0 s1 `1 L data.add(new Tuple3<>(1,2,11));
2 A8 J0 V/ G) n1 u7 r/ `2 W! A9 ~3 @1 V
data.add(new Tuple3<>(1,2,13));7 j* I+ g4 e; m- v8 M% b& j0 ` a
, T) I2 c" d% M2 w) s4 i4 N/ s' r# Z i; O- s
2 O6 y9 a+ {0 O1 x3 Q0 T, ~6 [! K v% D% X: O& W2 S2 w7 Y
& h- e0 I* t4 W0 O, C+ z* N- L" t
DataStreamSource<Tuple3<Integer,Integer,Integer>> items = env.fromCollection(data);7 Q' Q. b0 L" h8 x
8 _# E6 D; I9 Q4 D% O4 l3 `
! a% j2 V# j0 A) k2 E
, q6 ~; p" R9 v: j2 m: i: k1 j5 ?- c) f* i9 ?" m' {1 P) X
# i! C0 k z/ Q, j @6 Y
SplitStream<Tuple3<Integer, Integer, Integer>> splitStream = items.split(new OutputSelector<Tuple3<Integer, Integer, Integer>>() {
+ W& ?3 ^! g" K7 Y; ^0 ?1 D4 B" j, ?# ]
@Override3 f- Q8 M0 k( F5 ?
- ]3 @6 I# A+ \1 V1 {* _/ I
public Iterable<String> select(Tuple3<Integer, Integer, Integer> value) {
' L, K! U8 \! K7 M$ U# \, G5 P, o
. j! Q, y d* O9 f/ _/ L List<String> tags = new ArrayList<>();' ?* r: s; M( E/ W' j5 R
; t$ {1 Y. I, |0 o1 _. y# C# T- } if (value.f0 == 0) {
6 [4 i! M: X: @ G K0 l u
6 _) C( g5 M5 p7 g) i) \: c1 U! @' n: S tags.add("zeroStream");- M' u5 h! x0 b$ j! g
9 V. L5 j1 L5 @7 v; F, x } else if (value.f0 == 1) {
. R z! b3 m3 Z. z
# g7 N. P+ I0 A; [8 j9 R tags.add("oneStream");" L k( q" Y% J& z+ q2 W( Z
- f4 C; |$ L! o2 F6 ]* k
}
- O( z- D, m6 e7 q- D" o6 \& f# k
( Q$ q5 M4 i3 H/ D. N ] return tags;
0 G/ k+ C4 e7 R/ } l% ]/ J4 u5 V, W) V- S; W
} j1 q7 F' ~( P ^9 a) M0 P* y' J2 Y
; d; p0 h/ [7 |: S) a& ?3 b
});/ g) @) e# D# A6 j6 K
3 c6 O5 V r0 x9 T$ m
2 f' k5 ~" H% Y/ l7 @
6 D x7 n) G( Z3 l splitStream.select("zeroStream").print();+ h8 {( L j1 j
/ q6 w3 B; p* \# X8 N" e9 b
splitStream.select("oneStream").printToErr();* O4 k/ S: {' T! k/ F. V
, ?* P' U- y: X# C
% }8 q8 v! x0 {6 ?0 y& |& ~
* L+ \1 T" c% ]- ? //打印结果
0 N' M- C" ]1 t4 J! r* e1 T" {4 z0 z6 V( K
String jobName = "user defined streaming source";
4 w( M- `1 _; R, m5 H# n% W9 X
/ e G K$ q% P0 O& _! L8 r env.execute(jobName);8 J4 U0 H! E5 |& h
- u! r5 r3 N1 z& ?7 s/ K" J" | u8 `}2 N. X$ g! K0 _
</code></pre>
% v7 q! D6 u) `& I( g<p>同样,我们把来源的数据使用 split 算子进行了切分,并且打印出结果。</p>
8 ]; v m+ W- v- W3 x" k<p><img src="https://kingcall.oss-cn-hangzhou.aliyuncs.com/blog/img/CgqCHl7CA4aAbUSJAAG1LWNB3qw627-20210223084827377.png" ></p>
) M" S6 W" ^, c# G! I7 B @<p>但是要注意,使用 split 算子切分过的流,是不能进行二次切分的,假如把上述切分出来的 zeroStream 和 oneStream 流再次调用 split 切分,控制台会抛出以下异常。</p>/ n6 g$ t8 D3 i' C3 }' f
<p>复制代码</p>4 b% |( s" E* G; t9 n2 x
<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.
$ k: S2 t1 c3 Z6 z</code></pre>1 S2 F* u# ]6 C9 d" N1 T5 v$ t* `( k
<p>这是什么原因呢?我们在源码中可以看到注释,该方式已经废弃并且建议使用最新的 SideOutPut 进行分流操作。</p>; f {6 m- p5 h9 `) R2 X- z( N
<p><img src="https://kingcall.oss-cn-hangzhou.aliyuncs.com/blog/img/CgqCHl7CA6OAJ-JDAAIrh1JSAEo033-20210223084827602.png" ></p>
0 w3 [- b% M0 k$ m6 ~& Z) u<h4 id="sideoutput-分流">SideOutPut 分流</h4>. S9 `0 z0 F y( B9 j+ q
<p>SideOutPut 是 Flink 框架为我们提供的最新的也是最为推荐的分流方法,在使用 SideOutPut 时,需要按照以下步骤进行:</p>
5 e0 t4 T: R9 o7 p9 M3 F<ul>6 O* ~- A9 u% R3 h0 F
<li>定义 OutputTag</li>
/ ?* h* R# l' F! C/ N; a1 s8 c<li>调用特定函数进行数据拆分' c/ @" F2 R! K0 i
<ul># \! y# M) ^/ E2 v. r# h6 F
<li>ProcessFunction</li>
4 X( Z- c2 K# W<li>KeyedProcessFunction</li>
# O1 y: b/ o- g- t% }<li>CoProcessFunction</li>
) {) M4 p$ Y( y) X& s4 `' D; L<li>KeyedCoProcessFunction</li>8 w- L( I" |) Z6 {+ f) X
<li>ProcessWindowFunction</li>% o% ^: `; b. O1 O6 m9 Q
<li>ProcessAllWindowFunction</li>+ o6 S& C1 k/ s4 y( e) \
</ul>7 v( O+ u# z: v1 e" k+ N
</li>
" u; j& C7 H# a</ul>
, ?3 `8 _1 C3 g# n: @0 ~$ C<p>在这里我们使用 ProcessFunction 来讲解如何使用 SideOutPut:</p> e& \9 w" S' Z) `# W9 e
<p>复制代码</p>2 Z- E- y1 v9 L2 p# r
<pre><code class="language-java">public static void main(String[] args) throws Exception {8 m& \4 V, Z, s' C2 p
' v0 y- _ h7 l% P2 V7 E( e
( g2 M H6 @; y7 X) _8 e- b" x0 G& R. U) c! x! ~
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
& T, ?0 I8 _+ U- \ U7 b5 z, h6 R& l# Z( ^1 |
//获取数据源
3 ^/ j) Z7 j/ I: [# l( s3 s) n0 M3 x: k6 i2 z6 c) T* s
List data = new ArrayList<Tuple3<Integer,Integer,Integer>>();
8 c/ E; X& |/ C/ i2 r( E2 Z% ~9 i
1 |7 Q5 C, D# p4 b data.add(new Tuple3<>(0,1,0));, G9 A0 J8 ^/ W# Z6 g
% X& p1 ^7 Y _7 Q% _7 ?8 ?9 ?# W- B
data.add(new Tuple3<>(0,1,1));
) o$ c* {7 Q) B l
0 n8 O% ~# g5 z data.add(new Tuple3<>(0,2,2));
1 t( `2 v3 i( r7 Y
* R) ?* k+ j% f" C data.add(new Tuple3<>(0,1,3));* h/ x, L7 V9 s) {4 o3 F l
7 u- ?, o7 w( ~$ U$ D
data.add(new Tuple3<>(1,2,5));
/ @; }5 R: A, |( c4 d3 V6 @
) s, `1 m9 z2 I4 |- Z& t7 Q data.add(new Tuple3<>(1,2,9));3 v) S6 A/ e) ` D' P
, W* w4 \% M$ {: [4 c5 H" n, M data.add(new Tuple3<>(1,2,11));
5 C6 x- E: N8 q U
|$ B7 r: r M0 B7 f data.add(new Tuple3<>(1,2,13));
) T+ L' o9 p1 ?. m; L3 [5 [9 i7 d! V; G, h4 n
+ r0 o3 E1 ?: }- `7 O! ~3 ]% f+ \) J; G
' O: ]: m& J/ e$ H- }3 p- x% c
6 P+ w; {, B0 o& u' n- H DataStreamSource<Tuple3<Integer,Integer,Integer>> items = env.fromCollection(data);3 t2 D$ J* U. [ J9 F9 Y
, Y4 j+ O% d7 R* [7 I, V3 w
1 F z5 }# {6 S! a5 N1 R
3 T) T" Z: S; _ OutputTag<Tuple3<Integer,Integer,Integer>> zeroStream = new OutputTag<Tuple3<Integer,Integer,Integer>>("zeroStream") {};2 Z, u5 Q. J! Z$ v2 S/ E0 s, Z% L/ N
1 R9 z$ ^. }: Z* _; ~# v4 x
OutputTag<Tuple3<Integer,Integer,Integer>> oneStream = new OutputTag<Tuple3<Integer,Integer,Integer>>("oneStream") {};
) D- y7 o) A) d6 [5 r6 v7 P1 Y8 w9 t _# k/ f
- Z0 u5 p, `1 \4 ~: p. L
" ~$ h- Z1 t6 ?* W0 i2 c% l5 h
+ o- h2 W4 J% P/ z9 ~# S2 I
& c) u, e: y2 F: N) r" J U3 y* N0 q5 f SingleOutputStreamOperator<Tuple3<Integer, Integer, Integer>> processStream= items.process(new ProcessFunction<Tuple3<Integer, Integer, Integer>, Tuple3<Integer, Integer, Integer>>() {9 V- n+ `" V5 d
! _5 P5 O" u/ V0 y9 w9 l @Override
; a8 S* U& y \) W* u$ D( |% @
& d' j( s( N0 e public void processElement(Tuple3<Integer, Integer, Integer> value, Context ctx, Collector<Tuple3<Integer, Integer, Integer>> out) throws Exception {
. f4 Y; ?6 v. q2 F$ J. p1 F4 |
. y+ S; w p Z' m+ C, K# l
6 f) m: U) p0 A% M* O8 z% `
3 y, i- h7 P, p# y if (value.f0 == 0) {
/ y* z( p" f0 [. r8 E+ t* e4 v% D! G2 m$ x& `$ B
ctx.output(zeroStream, value);( p# W3 U/ V9 L) W0 v. C" l0 d
/ A! Z* U( R" W
} else if (value.f0 == 1) {- d. o U, e0 `! j% L
* q7 d8 |% K5 T& H# l# [+ [
ctx.output(oneStream, value);
5 N8 d3 E6 a# A8 p. ~- C" K1 w5 W* F
6 `# g5 R6 a6 V, J+ Y }
- _& ^% I6 E% A3 Z! Z4 U& \) j. }0 T" }
}
% ]; y" c& L% q! s; Z) u' M1 E& [8 N6 L1 V
});
0 M! J2 U7 v8 c1 q9 C
( z4 Q$ k n& S) {4 V9 r, M5 s0 n8 ^- d9 |
" w9 ?4 z1 h9 ~# G$ X p4 h DataStream<Tuple3<Integer, Integer, Integer>> zeroSideOutput = processStream.getSideOutput(zeroStream);* p) J( b! N' \/ F9 ~
2 L2 ^- Z$ n+ h* A, N: a) W DataStream<Tuple3<Integer, Integer, Integer>> oneSideOutput = processStream.getSideOutput(oneStream);, F6 h$ x% T1 _" W- d, }
+ q2 c! y: n: v. b! `
4 X3 k: r1 l7 L: s( l4 O9 L0 w) T: g: A7 e
zeroSideOutput.print();! O$ v7 |/ F6 s7 {% M3 Z# h) z! M
0 M' J: ~) x/ b& ] oneSideOutput.printToErr();
4 j% A4 ~9 ~* ?" J$ V# I
$ R, T4 F$ r( B9 Z' Y! Z' B# k3 Z4 ^0 v' a1 _ j9 ?
) A9 L% f% L1 G! z; z, t" _
! e# p' }' X q! ~5 _5 ]- K: d1 d$ x/ E! }( I: P
//打印结果
. R6 Z- Z! o$ k3 v3 i
$ \% c! }& v ~- S8 R+ Z4 b" T String jobName = "user defined streaming source";
2 `8 p1 L6 Q- O% h, L" `
, r- [3 ]- j2 a2 `* _4 C# U env.execute(jobName);; \* H; ~+ Z+ i: C% m0 V8 ~$ ?
: I7 F( X1 C" A' k( K* M: m/ n
}8 f8 H1 y8 F+ C# W, a
</code></pre> e1 v6 g3 H# S4 T U' @ |
<p>可以看到,我们将流进行了拆分,并且成功打印出了结果。这里要注意,Flink 最新提供的 SideOutPut 方式拆分流是<strong>可以多次进行拆分</strong>的,无需担心会爆出异常。</p>9 o- ]2 A) z1 j
<p><img src="https://kingcall.oss-cn-hangzhou.aliyuncs.com/blog/img/CgqCHl7CBMKAGHoUAAM-5UL5geg132-20210223084827698.png" ></p>
/ m( N' G' t# f+ G6 a! r<h3 id="总结">总结</h3>; t: S6 g, {# `4 U' t; O0 e
<p>这一课时我们讲解了 Flink 的一个小的知识点,是我们生产实践中经常遇到的场景,Flink 在最新的版本中也推荐我们使用 SideOutPut 进行流的拆分。</p>
9 \) D+ J% h7 d! v0 P2 @<blockquote>
9 M8 w: K0 h# m9 B) D f7 K<p>关注公众号:<code>大数据技术派</code>,回复<code>资料</code>,领取<code>1024G</code>资料。</p>7 H9 {( S1 w" Q! d1 k4 M3 C! s
</blockquote>9 G) o3 U9 a8 v& c
& E* Q# Z% d* A) f4 @8 |) ` |
|