|
|
! M1 G- Q. a3 }$ h- F# o
<h4 id="flink系列文章">Flink系列文章</h4>
3 u! R! B0 R$ q<ol>. K Q! v* `0 M, w
<li><a href="https://www.ikeguang.com/?p=1976">第01讲:Flink 的应用场景和架构模型</a></li>! ^5 h% r' b" x+ l d
<li><a href="https://www.ikeguang.com/?p=1977">第02讲:Flink 入门程序 WordCount 和 SQL 实现</a></li>
& m1 X4 Y7 x, [& Y* J2 k4 F; x3 m<li><a href="https://www.ikeguang.com/?p=1978">第03讲:Flink 的编程模型与其他框架比较</a></li>5 c' n7 D- x& u
<li><a href="https://www.ikeguang.com/?p=1982">第04讲:Flink 常用的 DataSet 和 DataStream API</a></li>% K) m( [# R7 }% X5 V; [
<li><a href="https://www.ikeguang.com/?p=1983">第05讲:Flink SQL & Table 编程和案例</a></li>
) ^: A. {6 h! j& [; l% ~+ p: v. Z<li><a href="https://www.ikeguang.com/?p=1985">第06讲:Flink 集群安装部署和 HA 配置</a></li>& E3 Q+ A( W v5 q( D' a
<li><a href="https://www.ikeguang.com/?p=1986">第07讲:Flink 常见核心概念分析</a></li>4 W" @5 E! n9 S" _3 o
<li><a href="https://www.ikeguang.com/?p=1987">第08讲:Flink 窗口、时间和水印</a></li>7 G2 K: Y7 y% w* f0 r/ A
<li><a href="https://www.ikeguang.com/?p=1988">第09讲:Flink 状态与容错</a></li> s' H8 T" w& V7 `7 Q* a$ o
</ol>
; n/ M( g* g- {<blockquote>
' l& u" b. Z, _2 p+ U5 w' I<p>关注公众号:<code>大数据技术派</code>,回复<code>资料</code>,领取<code>1024G</code>资料。</p>+ W5 p; @, S/ Z& Y/ |$ \
</blockquote>
5 c* T/ f/ r: g. l; b<p>这一课时将介绍 Flink 中提供的一个很重要的功能:旁路分流器。</p>7 o9 N5 y- J( e2 j3 F0 Y
<h3 id="分流场景">分流场景</h3>5 o/ G. D1 ]2 y& Z$ Y1 H
<p>我们在生产实践中经常会遇到这样的场景,需把输入源按照需要进行拆分,比如我期望把订单流按照金额大小进行拆分,或者把用户访问日志按照访问者的地理位置进行拆分等。面对这样的需求该如何操作呢?</p> b" w: _' v' {( W. X* s
<h3 id="分流的方法">分流的方法</h3>4 p* E F7 ?% P Q% c; U
<p>通常来说针对不同的场景,有以下三种办法进行流的拆分。</p>
5 l1 K6 a3 L Z" f<h4 id="filter-分流">Filter 分流</h4>1 g2 b$ J' H/ ^. F) R
<p><img src="https://kingcall.oss-cn-hangzhou.aliyuncs.com/blog/img/CgqCHl7CAy6ADUaXAACSFUbdpuA911-20210223084827182.png" ></p>& S9 t% |/ l* n
<p>Filter 方法我们在第 04 课时中(Flink 常用的 DataSet 和 DataStream API)讲过,这个算子用来根据用户输入的条件进行过滤,每个元素都会被 filter() 函数处理,如果 filter() 函数返回 true 则保留,否则丢弃。那么用在分流的场景,我们可以做多次 filter,把我们需要的不同数据生成不同的流。</p>
2 K! L! D; p/ t/ c# J) x<p>来看下面的例子:</p>
* l) a3 ]" O W) @<p>复制代码</p>
/ q# S' `: }0 B) S( g<pre><code class="language-java">public static void main(String[] args) throws Exception {
& ~! i+ K' v- A' j: T StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();, D. C4 ?. x) ~; ~! s4 k6 S
//获取数据源7 o/ `6 Y& Z3 ?) q
List data = new ArrayList<Tuple3<Integer,Integer,Integer>>();+ |0 Y. K# q, ]
data.add(new Tuple3<>(0,1,0));, i9 @& _5 q2 Z2 O3 H4 _
data.add(new Tuple3<>(0,1,1));
% c0 B+ m/ G) J) d, g( ] data.add(new Tuple3<>(0,2,2));" j5 J; ^: `5 k \ [, I
data.add(new Tuple3<>(0,1,3));, f& [+ ?+ Z" w7 T
data.add(new Tuple3<>(1,2,5));2 A/ Q* j2 I* L7 h' k% u- i4 A
data.add(new Tuple3<>(1,2,9));! g1 Z% N5 _4 i3 G9 E
data.add(new Tuple3<>(1,2,11));
( M! m4 `' [3 H0 c9 }5 t data.add(new Tuple3<>(1,2,13));* Y2 `; k* d# K8 |) V. P' Z1 i
; e# k+ `* D- D. l DataStreamSource<Tuple3<Integer,Integer,Integer>> items = env.fromCollection(data);( E8 p; g7 U' u% B
) I+ `/ x$ r# n
9 S* D' D5 k8 D; }4 y! z6 ^. v0 p' q4 `9 l: [3 @3 N
SingleOutputStreamOperator<Tuple3<Integer, Integer, Integer>> zeroStream = items.filter((FilterFunction<Tuple3<Integer, Integer, Integer>>) value -> value.f0 == 0);
4 ^8 X1 B- g4 y# x( ^1 ~* q3 S$ ~$ r3 G4 a. a' J
SingleOutputStreamOperator<Tuple3<Integer, Integer, Integer>> oneStream = items.filter((FilterFunction<Tuple3<Integer, Integer, Integer>>) value -> value.f0 == 1);2 Q6 O1 f; v9 B. z2 ?5 b$ S, n1 { [
; V$ z, r$ |3 t7 i1 o0 m" L
0 f" B% P7 E G g) B0 u$ y3 q- v- N# N) o1 t6 A1 L- k
zeroStream.print();9 [$ J5 d# K5 C$ A6 a
/ b; X) [/ P; Z' E2 |# @2 p oneStream.printToErr();
: c8 G7 {" _5 N; S
" c. Q5 w1 C! V( x* I2 |; k9 Q) G3 S5 v9 s+ I
: L9 G1 w) n! c C1 ^
4 p6 r6 \3 g! d% \" V( d' G# _ |9 e0 t7 g) [) ~* y
//打印结果5 t( V& }8 q& i- R9 L. f# G$ \
/ E7 {: k* e6 ~+ H- J6 D String jobName = "user defined streaming source";/ f. D9 e6 e9 K2 G
" `/ Q1 J4 r1 f# d# u5 F env.execute(jobName);6 u, W l+ D/ Z' [9 l; | F
) e' o" w' Z- L$ k}
3 X4 ~1 c$ ?/ ?2 h0 ^</code></pre>
! i$ ]" R7 F; L# J5 ]<p>在上面的例子中我们使用 filter 算子将原始流进行了拆分,输入数据第一个元素为 0 的数据和第一个元素为 1 的数据分别被写入到了 zeroStream 和 oneStream 中,然后把两个流进行了打印。</p>4 V' p' R/ o. B7 E- t% ?' }9 l) S! I
<p><img src="https://kingcall.oss-cn-hangzhou.aliyuncs.com/blog/img/Ciqc1F7CA2WAYbshAAKj494h86s723-20210223084827296.png" ></p>/ i3 j7 U, x% y
<p>可以看到 zeroStream 和 oneStream 分别被打印出来。</p>
) W+ _6 A" |$ Q, C# p5 b<p>Filter 的弊端是显而易见的,为了得到我们需要的流数据,需要多次遍历原始流,这样无形中浪费了我们集群的资源。</p>
8 ~1 f" y! H1 x! z7 x8 Y* v# q<h4 id="split-分流">Split 分流</h4>' P! K% l! N/ e
<p>Split 也是 Flink 提供给我们将流进行切分的方法,需要在 split 算子中定义 OutputSelector,然后重写其中的 select 方法,将不同类型的数据进行标记,最后对返回的 SplitStream 使用 select 方法将对应的数据选择出来。</p>
$ h! [- i& J: _! U6 L<p>我们来看下面的例子:</p>
) [0 G+ S8 R8 e4 w<p>复制代码</p>
5 |. A8 _8 d8 [<pre><code class="language-java">public static void main(String[] args) throws Exception {
; |' h8 G* T6 o: m5 ~, @& a4 V5 }' C3 \0 F7 P& I
! k/ P4 ]2 w7 h( |
3 ~ d ^$ V, k* ]
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();" J- Y# e: ] X3 n6 c+ ~- K
" c! I0 B( X3 V. t //获取数据源
0 F* g, g9 V Y0 C
+ c% ~$ l5 X" d0 c List data = new ArrayList<Tuple3<Integer,Integer,Integer>>();1 N6 m6 O& E" c2 U2 v8 G
l0 _4 I( \/ \ data.add(new Tuple3<>(0,1,0));1 t6 k8 D9 L3 N5 C
1 L1 m5 S, R5 w$ [% j! i, K data.add(new Tuple3<>(0,1,1));" ~6 B& T. Y" C9 ]+ q
, ?; H1 t- t) W7 Z data.add(new Tuple3<>(0,2,2));# c- s! t( \* {. ?/ e
$ D- d9 C2 h* \6 ^8 ^ data.add(new Tuple3<>(0,1,3));
) X8 G% S [1 L, f- f1 F7 b0 [; t; w$ n/ D# G+ E
data.add(new Tuple3<>(1,2,5));
7 n$ y( {2 Y# e, r# r: y" R% d1 ^4 y& G. v; @ I
data.add(new Tuple3<>(1,2,9));- n5 q9 o& G: j9 S: x9 T
0 k, a( `) k6 X* K$ Q( ^
data.add(new Tuple3<>(1,2,11));5 b% P7 |8 A$ B _- M9 y' | h. v
: y0 ?1 i7 l9 l: ]7 m) Z& H
data.add(new Tuple3<>(1,2,13));6 w: j' J" ^# m! c' E7 R
8 i7 F% w' M1 a9 Y1 P B* F4 x1 C% y
9 K& X \; g: s. N1 O1 m r8 f
( r7 N- Q, B: P- q0 h& a
, C/ e/ W# p" J. C3 E! Y/ w DataStreamSource<Tuple3<Integer,Integer,Integer>> items = env.fromCollection(data);
) s7 g' q2 ?# `5 }2 S; _9 U p- @& J, R" n$ ~* F# l
! _, z3 E4 y) h+ W$ l2 k, ~: F8 m: D* m4 f7 G' Q
- k/ D/ \( Q1 }9 ~) T- O4 Y4 p$ i
$ P1 w G; G8 A- A2 G
SplitStream<Tuple3<Integer, Integer, Integer>> splitStream = items.split(new OutputSelector<Tuple3<Integer, Integer, Integer>>() {# X& ?4 s6 l: v6 n% l
. f% |0 n+ {: E1 h1 r @Override* r6 c7 i3 j0 Z- j$ w# ?
/ _' Q& _8 @& O# \2 t public Iterable<String> select(Tuple3<Integer, Integer, Integer> value) {
% C$ `+ `' A; w! Z5 S/ N! @$ G# Z5 C6 z& ?
List<String> tags = new ArrayList<>();
( C0 `* c% t6 y# R8 }9 g9 J; ^5 c
7 y$ }% Q j3 k# Y% M5 s7 E& u, X if (value.f0 == 0) {
: l( s u- z S7 e) l; D0 a! U0 V% i, X- \
tags.add("zeroStream");
" [5 B' z& N2 d' Z, A4 p5 {3 H2 t
( h8 g, F/ v0 i3 P G _0 R( c; x3 f } else if (value.f0 == 1) {1 T G- O: f. n- M3 I7 V
7 C# F! w9 ]7 O3 C9 s! ?
tags.add("oneStream");* `6 n1 B6 W3 X) C2 F5 z& d/ A
- C0 r4 E, v% G* U }
9 p4 p& t6 z4 ^. m1 X) ]: j M' H' M* m# k" b' D- X; a: B
return tags;9 l3 a3 I% ]0 A/ v9 R7 T3 h
, W2 {' X- n2 ~/ D+ @% Q' q# e) d }: _" Y8 X5 [- v I# C+ m% M/ y
5 w/ C3 H$ O) Y% X: q
});
* L0 R4 @$ \ h/ L. u' }# k. g; ~+ q R
; }+ W1 j: N) q, C- x' }* J* `' F! G* [ d
splitStream.select("zeroStream").print();
$ c- p7 l! K; o8 m& L2 W% d1 W( @7 V# v' l! A3 H* w* R
splitStream.select("oneStream").printToErr();6 ?2 w- a. e5 R0 t6 @. n5 ]
; y9 [8 u7 `) p- } H
# z7 g6 X' G$ D; F
4 d; l! j/ c# Z //打印结果0 m: I8 ~6 [% }! g. D: b& Z
0 t( P5 f+ E1 Y8 [ g. G String jobName = "user defined streaming source";
- c7 r2 r u; }$ |* Y9 K2 H( [8 v
env.execute(jobName);
& w; m6 q- m; F& }0 H9 @/ r6 z* l; z& k1 x1 N- L3 V j' J0 K
}4 \9 L/ v; v( r; q, a T
</code></pre>
% y& G/ I/ x6 L3 H2 |% Q<p>同样,我们把来源的数据使用 split 算子进行了切分,并且打印出结果。</p>! l4 L3 D( V& O# }! l9 F2 o; ~. t
<p><img src="https://kingcall.oss-cn-hangzhou.aliyuncs.com/blog/img/CgqCHl7CA4aAbUSJAAG1LWNB3qw627-20210223084827377.png" ></p>
6 a6 ?2 v+ t" B) N7 n<p>但是要注意,使用 split 算子切分过的流,是不能进行二次切分的,假如把上述切分出来的 zeroStream 和 oneStream 流再次调用 split 切分,控制台会抛出以下异常。</p># ^. Y+ t& V3 R* U# r0 e
<p>复制代码</p>
8 ]7 P# ?% B+ j9 D7 r<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.. I/ Z. B4 c. A
</code></pre>3 `5 @' o. o4 h: }5 U. M! [
<p>这是什么原因呢?我们在源码中可以看到注释,该方式已经废弃并且建议使用最新的 SideOutPut 进行分流操作。</p>2 W! p5 [6 _" c* b) J! u
<p><img src="https://kingcall.oss-cn-hangzhou.aliyuncs.com/blog/img/CgqCHl7CA6OAJ-JDAAIrh1JSAEo033-20210223084827602.png" ></p>( t4 m! j* T- A8 l$ Y/ k
<h4 id="sideoutput-分流">SideOutPut 分流</h4>$ G9 Q" N* H4 ]! I- r
<p>SideOutPut 是 Flink 框架为我们提供的最新的也是最为推荐的分流方法,在使用 SideOutPut 时,需要按照以下步骤进行:</p>
; _: V) s7 k7 c( Z8 h<ul>
0 W& R8 l) I) o6 T7 }* Y: v! N<li>定义 OutputTag</li>' |4 p' l/ M9 S4 n
<li>调用特定函数进行数据拆分: D) @" Y( t" T8 t% ]6 `
<ul>
V+ Z1 E7 \& k<li>ProcessFunction</li>
6 k3 Q+ z9 ?; \/ W5 |- V3 y<li>KeyedProcessFunction</li>
. d7 J3 i( D7 h' j" q9 Q0 ?* R0 h<li>CoProcessFunction</li>1 j% c7 \& S2 d7 ]: \2 N
<li>KeyedCoProcessFunction</li>
$ i% e& s( j/ l1 G; o& L, Y- m; Q<li>ProcessWindowFunction</li>( j& q6 _( G1 b O
<li>ProcessAllWindowFunction</li>
3 i5 ?8 t( T! p1 E# L</ul>& C5 a. }% f3 f, k
</li>, z* a: \; w; H/ \
</ul>' s; w3 w2 C8 H% ]' ^
<p>在这里我们使用 ProcessFunction 来讲解如何使用 SideOutPut:</p>
3 q) b) b( I; L p$ K5 k<p>复制代码</p>8 H4 t' b2 r6 ]7 _( b
<pre><code class="language-java">public static void main(String[] args) throws Exception {
: t. B1 y+ f# g1 G$ C3 n4 ^( F& x2 O4 ~+ ~* [! T) M* f
& W B& e- }& j, r# R4 }* H7 N1 m u
; E- i( P! d. ` StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();! S) P6 w% v3 D; E4 F+ y
) o! h- \4 W8 n) X //获取数据源. v/ z. t# U z+ j( T
9 p6 \# w- A7 ~! u8 H. \% N List data = new ArrayList<Tuple3<Integer,Integer,Integer>>();
; r( v( ~5 D! w$ |
4 x" a' L. p$ j data.add(new Tuple3<>(0,1,0));
/ z, Q9 G# J6 _9 A( y) ?0 E$ J5 [3 r7 n) a% v) e
data.add(new Tuple3<>(0,1,1));: Q1 B% v# }! ^. _+ u9 Z$ H
9 j& F8 u$ k4 U/ J/ y# E0 e, e
data.add(new Tuple3<>(0,2,2));- c3 l: c8 }- j. d# \' n& E
. {* u# w& _ d" V! a2 F' y% T/ T2 l; T
data.add(new Tuple3<>(0,1,3));: H) l+ i7 I* F
) H1 r8 ~% z1 d6 B3 E2 d data.add(new Tuple3<>(1,2,5));9 }4 Q( \& w. d0 h! s" L
; |* B# A( g& o1 p7 W! H" X7 s4 j data.add(new Tuple3<>(1,2,9));
7 z6 O) s" T! M% i' P% a* Y5 Y# r* R5 v& U, V, B3 @
data.add(new Tuple3<>(1,2,11));
9 j- b. M" {. Q
# Z6 N) h7 r5 {% t5 }/ J data.add(new Tuple3<>(1,2,13));/ B% {) _$ L( ]/ @$ z( n. B' {
# X* D0 S% x, D# [7 j- s, n; D3 L& O6 r
3 m' u. r3 X1 A) e/ O# D% Y6 z
" Z* J9 t& U1 ?# a3 L2 _ @9 [0 h/ b$ q
DataStreamSource<Tuple3<Integer,Integer,Integer>> items = env.fromCollection(data);
$ t! l$ m) J) c* ]6 b* ]& O8 ]) t
! Y( s# V" R. |9 t
7 S" i q% l; W% g3 C8 J" m OutputTag<Tuple3<Integer,Integer,Integer>> zeroStream = new OutputTag<Tuple3<Integer,Integer,Integer>>("zeroStream") {};/ J A. P8 L& I; ~ x+ \
2 j1 G0 N: `3 S6 t+ X; H OutputTag<Tuple3<Integer,Integer,Integer>> oneStream = new OutputTag<Tuple3<Integer,Integer,Integer>>("oneStream") {};. b9 M# f! m# x7 y
. h* |, S+ j0 ~0 X; S8 t
7 S; Z2 D' w+ V7 V! s) F5 c8 M# j8 e
7 m. l# U8 I4 ~$ g& ~3 S0 d2 i* I! L8 p* w0 W4 c
8 p' H3 ?, V$ u% m' S" R
SingleOutputStreamOperator<Tuple3<Integer, Integer, Integer>> processStream= items.process(new ProcessFunction<Tuple3<Integer, Integer, Integer>, Tuple3<Integer, Integer, Integer>>() {
, A" J5 D4 [9 Y" c2 J5 H" J: M+ j- o) x4 E, r
@Override
9 S% R; Q3 P, V3 t1 g
$ I* m' ^1 l3 w$ ]' u4 K+ k public void processElement(Tuple3<Integer, Integer, Integer> value, Context ctx, Collector<Tuple3<Integer, Integer, Integer>> out) throws Exception {
) f0 P7 |0 t* w' X* R+ m
. A5 C* W! y2 f+ \9 l) R8 N8 |2 t; o2 w0 {# C1 k8 l) Z) ^) C
4 x% J7 F+ V# ?! {8 q if (value.f0 == 0) {& ^6 [* C4 D( C: b! Z1 ?* i2 I
- L# m2 y$ N# @ ctx.output(zeroStream, value);* f) |* d0 e; q7 ~
. B' d1 v3 S( E& ?$ n) I% m } else if (value.f0 == 1) {2 w% R. W& \% D8 O) C5 x
; A& A7 N }* J2 t, P9 h% r" b ctx.output(oneStream, value);
3 I% u( k3 F! `4 b- L& K2 T' n
7 C2 ]( z! p0 s2 P2 Y9 m }7 m1 A2 o0 P& l. R! B" O
1 r, L/ Z! Q1 X! l. ?3 l7 _1 S
}
( c+ |5 f& O# Y
8 t; q, _: P) z$ N2 g# a });2 |5 V* p. T0 x" _) Q0 C x& _
) f x3 a" x- f4 M& k3 I: U
6 z. L) T. |. i$ H6 g) {/ k$ j, S2 N
DataStream<Tuple3<Integer, Integer, Integer>> zeroSideOutput = processStream.getSideOutput(zeroStream);. D0 J% i9 m2 O8 j2 a& g( w3 m- x0 e
/ P5 ]: }2 @5 ?% g+ G6 ?/ C
DataStream<Tuple3<Integer, Integer, Integer>> oneSideOutput = processStream.getSideOutput(oneStream);7 |) t6 z/ j9 d& j) p; O/ @) S
9 p5 D- s! P" O! a$ j7 F$ a* ^4 |$ a f2 s+ ^
6 P, `- [2 Z$ d) V/ r zeroSideOutput.print();
% I( P- w6 r0 F, v& M7 }
7 Z. Q w4 u9 D( d) p' k oneSideOutput.printToErr();
6 S# h" {& U6 r0 F( c O
( J+ ], ^$ b/ B* A( a; O- H, l6 S' X5 m+ x/ u# |
0 i! i3 G- u" y# Y. `2 ~$ O
1 S2 b4 r% C3 ]8 j, \
7 ^4 d/ D$ m4 ]- q
//打印结果8 R9 O; y, `# `7 I
/ o7 [& e8 |! w) ^" A String jobName = "user defined streaming source";: b+ B8 x9 x! F' N$ y7 E0 [
* r6 r7 u3 G q
env.execute(jobName);1 o2 W% T2 W9 N
, U% s; a# L6 J}
5 E# F7 @* z" x</code></pre> p2 t8 {, r% [ }
<p>可以看到,我们将流进行了拆分,并且成功打印出了结果。这里要注意,Flink 最新提供的 SideOutPut 方式拆分流是<strong>可以多次进行拆分</strong>的,无需担心会爆出异常。</p>+ T$ u P; j Z- u+ W) R3 e* h/ T
<p><img src="https://kingcall.oss-cn-hangzhou.aliyuncs.com/blog/img/CgqCHl7CBMKAGHoUAAM-5UL5geg132-20210223084827698.png" ></p>& M) }) P; s2 x% C5 u& n# |
<h3 id="总结">总结</h3>6 S3 k6 [" r% E! x) L I% z
<p>这一课时我们讲解了 Flink 的一个小的知识点,是我们生产实践中经常遇到的场景,Flink 在最新的版本中也推荐我们使用 SideOutPut 进行流的拆分。</p>
2 T2 W* x* Z. ]3 K1 h. K6 S! D: Q<blockquote>
9 d4 q, D0 D2 d% x W6 d<p>关注公众号:<code>大数据技术派</code>,回复<code>资料</code>,领取<code>1024G</code>资料。</p>
, `: D5 D# [. M- t( D1 B C</blockquote>3 Q- g! }( \' `' O2 B9 E' w
' B, F# ~2 E, E8 f. B% D
|
|