栏目分类:
子分类:
返回
名师互学网用户登录
快速导航关闭
当前搜索
当前分类
子分类
实用工具
热门搜索
名师互学网 > IT > 前沿技术 > 大数据 > 大数据系统

Flink: Function And Rich Function

Flink: Function And Rich Function

序言

       了解了Flink提供的算子,那我们就可以自定义算子了.自定义算子的目的是为了更加灵活的处理我们的业务数据,并将满足条件的结果Sink到目标存储地cuiyaonan2000@163.com

       Function有2中类型即 Function  和 Rich Function .从字面意思我们可以了解 Rich Function 肯定是比Function提供了更多的功能的.

参考版本为: v1.13.2

官网地址:用户自定义 Functions | Apache Flink

Function

以算子Map举例 我们自定义的Function需要实现类(针对其它算子比如Filter,则是实现接口FilterFunction)

MapFunction

举例实现类

class MyMapFunction implements MapFunction {
  public Integer map(String value) { return Integer.parseInt(value); }
};
data.map(new MyMapFunction());

Rich Functions

Rich functions 对比 Function 还提供了这些方法:open、close、getRuntimeContext 和 setRuntimeContext

由上可见一个RichFunctions的调用绝非是一次性买卖,提供了一个生命周期的管理,即从open开始到close结束.同时提供了getRuntimeContext 以获取Flink的上下文环境.

需要注意的是Flink 流数据处理_难得糊涂-CSDN博客_flink数据处理范式 里已经很清楚地表明了,算子是针对窗口的计算执行者,且算子有并行度,则每个算子的生命周期是相互独立的.这个很好理解因为算子的有并行度,且相同算子的可能有多个线程,且分布在不同的服务器上.cuiyaonan2000@163.com.  

  1. open()方法: 是rich function的初始化方法,当一个算子例如map或者filter被调用之前open()会被调用。
  2. close()方法: 是生命周期中的最后一个调用的方法,做一些清理工作。
  3. getRuntimeContext()方法: 提供了函数的RuntimeContext的一些信息,例如函数执行的并行度,任务的名字,以及state状态

同理针对Map算子,创建我们自己的自定义算子需要实现接口RichMapFunction

class MyMapFunction extends RichMapFunction {
  public Integer map(String value) { return Integer.parseInt(value); }
};

转载请注明:文章转载自 www.mshxw.com
本文地址:https://www.mshxw.com/it/350027.html
我们一直用心在做
关于我们 文章归档 网站地图 联系我们

版权所有 (c)2021-2022 MSHXW.COM

ICP备案号:晋ICP备2021003244-6号