ICode9

精准搜索请尝试: 精确搜索
首页 > 其他分享> 文章详细

《从0到1学习Flink》—— Data Sink 介绍

2021-04-03 15:02:29  阅读:212  来源: 互联网

标签:target Flink Source Sink Data public


图片

前言

再上一篇文章中 《从0到1学习Flink》—— Data Source 介绍 讲解了 Flink Data Source ,那么这里就来讲讲 Flink Data Sink 吧。

首先 Sink 的意思是:

image.png

大概可以猜到了吧!Data sink 有点把数据存储下来(落库)的意思。

image.png

如上图,Source 就是数据的来源,中间的 Compute 其实就是 Flink 干的事情,可以做一系列的操作,操作完后就把计算后的数据结果 Sink 到某个地方。(可以是 MySQL、ElasticSearch、Kafka、Cassandra 等)。这里我说下自己目前做告警这块就是把 Compute 计算后的结果 Sink 直接告警出来了(发送告警消息到钉钉群、邮件、短信等),这个 sink 的意思也不一定非得说成要把数据存储到某个地方去。其实官网用的 Connector 来形容要去的地方更合适,这个 Connector 可以有 MySQL、ElasticSearch、Kafka、Cassandra RabbitMQ 等。

Flink Data Sink

前面文章 《从0到1学习Flink》—— Data Source 介绍 介绍了 Flink Data Source 有哪些,这里也看看 Flink Data Sink 支持的有哪些。

图片

看下源码有哪些呢?

图片

可以看到有 Kafka、ElasticSearch、Socket、RabbitMQ、JDBC、Cassandra POJO、File、Print 等 Sink 的方式。

SinkFunction

图片

从上图可以看到 SinkFunction 接口有 invoke 方法,它有一个 RichSinkFunction 抽象类。

上面的那些自带的 Sink 可以看到都是继承了 RichSinkFunction 抽象类,实现了其中的方法,那么我们要是自己定义自己的 Sink 的话其实也是要按照这个套路来做的。

这里就拿个较为简单的 PrintSinkFunction 源码来讲下:

 1@PublicEvolving
2public class PrintSinkFunction<IN> extends RichSinkFunction<IN> {
3    private static final long serialVersionUID = 1L;
4
5    private static final boolean STD_OUT = false;
6    private static final boolean STD_ERR = true;
7
8    private boolean target;
9    private transient PrintStream stream;
10    private transient String prefix;
11
12    /**
13     * Instantiates a print sink function that prints to standard out.
14     */
15    public PrintSinkFunction() {}
16
17    /**
18     * Instantiates a print sink function that prints to standard out.
19     *
20     * @param stdErr True, if the format should print to standard error instead of standard out.
21     */
22    public PrintSinkFunction(boolean stdErr) {
23        target = stdErr;
24    }
25
26    public void setTargetToStandardOut() {
27        target = STD_OUT;
28    }
29
30    public void setTargetToStandardErr() {
31        target = STD_ERR;
32    }
33
34    @Override
35    public void open(Configuration parameters) throws Exception {
36        super.open(parameters);
37        StreamingRuntimeContext context = (StreamingRuntimeContext) getRuntimeContext();
38        // get the target stream
39        stream = target == STD_OUT ? System.out : System.err;
40
41        // set the prefix if we have a >1 parallelism
42        prefix = (context.getNumberOfParallelSubtasks() > 1) ?
43                ((context.getIndexOfThisSubtask() + 1) + "> ") : null;
44    }
45
46    @Override
47    public void invoke(IN record) {
48        if (prefix != null) {
49            stream.println(prefix + record.toString());
50        }
51        else {
52            stream.println(record.toString());
53        }
54    }
55
56    @Override
57    public void close() {
58        this.stream = null;
59        this.prefix = null;
60    }
61
62    @Override
63    public String toString() {
64        return "Print to " + (target == STD_OUT ? "System.out" : "System.err");
65    }
66}

可以看到它就是实现了 RichSinkFunction 抽象类,然后实现了 invoke 方法,这里 invoke 方法就是把记录打印出来了就是,没做其他的额外操作。

如何使用?

1SingleOutputStreamOperator.addSink(new PrintSinkFunction<>();

这样就可以了,如果是其他的 Sink Function 的话需要换成对应的。

使用这个 Function 其效果就是打印从 Source 过来的数据,和直接 Source.print() 效果一样。

图片

下篇文章我们将讲解下如何自定义自己的 Sink Function,并使用一个 demo 来教大家,让大家知道这个套路,且能够在自己工作中自定义自己需要的 Sink Function,来完成自己的工作需求。

最后

本文主要讲了下 Flink 的 Data Sink,并介绍了常见的 Data Sink,也看了下源码的 SinkFunction,介绍了一个简单的 Function 使用, 告诉了大家自定义 Sink Function 的套路,下篇文章带大家写个。

关注我

转载请务必注明原创地址为:http://www.54tianzhisheng.cn/2018/10/29/flink-sink/

另外我自己整理了些 Flink 的学习资料,目前已经全部放到微信公众号了。你可以加我的微信:zhisheng_tian,然后回复关键字:Flink 即可无条件获取到。

图片

相关文章

1、《从0到1学习Flink》—— Apache Flink 介绍

2、《从0到1学习Flink》—— Mac 上搭建 Flink 1.6.0 环境并构建运行简单程序入门

3、《从0到1学习Flink》—— Flink 配置文件详解

4、《从0到1学习Flink》—— Data Source 介绍

5、《从0到1学习Flink》—— 如何自定义 Data Source ?

6、《从0到1学习Flink》—— Data Sink 介绍

7、《从0到1学习Flink》—— 如何自定义 Data Sink ?


标签:target,Flink,Source,Sink,Data,public
来源: https://blog.51cto.com/15060469/2682543

本站声明: 1. iCode9 技术分享网(下文简称本站)提供的所有内容,仅供技术学习、探讨和分享;
2. 关于本站的所有留言、评论、转载及引用,纯属内容发起人的个人观点,与本站观点和立场无关;
3. 关于本站的所有言论和文字,纯属内容发起人的个人观点,与本站观点和立场无关;
4. 本站文章均是网友提供,不完全保证技术分享内容的完整性、准确性、时效性、风险性和版权归属;如您发现该文章侵犯了您的权益,可联系我们第一时间进行删除;
5. 本站为非盈利性的个人网站,所有内容不会用来进行牟利,也不会利用任何形式的广告来间接获益,纯粹是为了广大技术爱好者提供技术内容和技术思想的分享性交流网站。

专注分享技术,共同学习,共同进步。侵权联系[81616952@qq.com]

Copyright (C)ICode9.com, All Rights Reserved.

ICode9版权所有